- 4001
Json(serde_json::json!({ "error": e.to_string() })), - 4002
) - 4003
.into_response(), - 4004
} - 4005
} - 4006
- 4007
/// Resolve a durable conversation into the live handle map at an admission - 4008
/// boundary. A browser may keep its page and EventSources across a server - 4009
/// restart; neither a follow-up run nor a reconnected stream can assume an - 4010
/// earlier explicit `/attach` call still exists in this process. - 4011
async fn ensure_session_handle( - 4012
state: &AppState, - 4013
session_id: &str, - 4014
) -> Result<(String, Arc<SessionHandle>), vak_core::CoreError> { - 4015
// Already attached? Return before touching the file. - 4016
// - 4017
// The live handle owns an exclusive lock on the session JSONL for its - 4018
// whole lifetime, and the lock is per open-file-description: opening the - 4019
// same path again from THIS process conflicts with our own handle just - 4020
// as it would with a stranger's. Re-attaching is routine — the desktop - 4021
// calls it on every task switch, and mid-run the handle's session is - 4022
// temporarily owned by the agent — so this must be a no-op, not a - 4023
// second open. - 4024
if vak_core::trash::is_trashed(&state.core.shared_data_home(), session_id) { - 4025
return Err(vak_core::CoreError::Session(vak_session::SessionError::Io( - 4026
std::io::Error::new( - 4027
std::io::ErrorKind::NotFound, - 4028
format!("session is in the trash: {session_id}"), - 4029
), - 4030
))); - 4031
} - 4032
if let Some(handle) = state.get(session_id) { - 4033
return Ok((session_id.to_owned(), handle)); - 4034
} - 4035
let session = if let Ok(s) = state.active_core().open_session(session_id).await { - 4036
Ok(s) - 4037
} else if let Ok(s) = state.core.open_session(session_id).await { - 4038
Ok(s) - 4039
} else if let Ok(s) = state.active_core().open_session_read_only(session_id).await { - 4040
Ok(s) - 4041
} else if let Ok(s) = state.core.open_session_read_only(session_id).await { - 4042
Ok(s) - 4043
} else { - 4044
find_session_on_disk(&state.core, session_id).ok_or_else(|| { - 4045
vak_core::CoreError::Session(vak_session::SessionError::Io(std::io::Error::new( - 4046
std::io::ErrorKind::NotFound, - 4047
format!("session not found: {session_id}"), - 4048
))) - 4049
}) - 4050
}; - 4051
session.map(|session| { - 4052
let header = session.header(); - 4053
let id = header - 4054
.map(|h| h.session_id.clone()) - 4055
.unwrap_or_else(|| session_id.to_owned()); - 4056
let session_cwd = header - 4057
.map(|h| h.cwd.clone()) - 4058
.unwrap_or_else(|| state.core.cwd().clone()); - 4059
let handle_core = match header { - 4060
Some(header) => resolve_core_for_header(state, header), - 4061
None => resolve_process_core_for_cwd(state, &state.active_core(), &session_cwd), - 4062
}; - 4063
// The header id can differ from the requested one; if that handle - 4064
// is already live, keep it rather than replacing it. - 4065
let handle = state.get(&id).unwrap_or_else(|| { - 4066
register_handle(state, id.clone(), session, session_cwd, handle_core) - 4067
}); - 4068
(id, handle) - 4069
}) - 4070
} - 4071
- 4072
/// The one derivation of "which `Core` should this session's next turn run - 4073
/// under," from its own recorded header (finding 6). A custom Agent's - 4074
/// session gets the exact pins `agent_chats::resolve_agent_core` applies — - 4075
/// permission-mode cap, sandbox backend override, provider instance - 4076
/// override, and shared sessions_home — via `agent_chats:: - 4077
/// pinned_core_for_workspace`, keyed by the session's OWN recorded - 4078
/// `header.cwd` rather than a workspace path re-derived from whatever - 4079
/// happens to be the CURRENT active workspace (which can disagree once the - 4080
/// active workspace has moved on since the session was created — see that - 4081
/// function's doc comment). Before this, `/attach`, `/run` and the SSE - 4082
/// endpoints resolved a plain pooled `Core` for the cwd with none of those - 4083
/// pins, so the same session's security ceiling depended on which endpoint - 4084
/// happened to touch it first. The built-in `vak` identity (and a - 4085
/// pre-Agent ledger with no `agent` on its header at all) keeps the plain - 4086
/// active/process core resolution. - 4087
fn resolve_core_for_header(state: &AppState, header: &vak_session::types::SessionHeader) -> Core { - 4088
let active = state.active_core(); - 4089
match header.agent.as_ref() { - 4090
Some(identity) if identity.id != "vak" => { - 4091
agent_chats::pinned_core_for_workspace(state, &active, identity, &header.cwd) - 4092
.unwrap_or_else(|_| resolve_process_core_for_cwd(state, &active, &header.cwd)) - 4093
} - 4094
_ => resolve_process_core_for_cwd(state, &active, &header.cwd), - 4095
} - 4096
} - 4097
- 4098
/// Plain cwd-keyed pooled `Core` resolution with no Agent-specific pins — - 4099
/// the built-in `vak` identity's own path, and `resolve_core_for_header`'s - 4100
/// fallback when a custom Agent's pinned resolution itself fails. - 4101
fn resolve_process_core_for_cwd(state: &AppState, active: &Core, cwd: &std::path::Path) -> Core { - 4102
if cwd == active.cwd().as_path() { - 4103
return active.clone(); - 4104
} - 4105
if cwd == state.core.cwd().as_path() { - 4106
return state.core.clone(); - 4107
} - 4108
if let Ok(c) = state - 4109
.gateway - 4110
.core_pool - 4111
.resolve_at(cwd, None, std::time::Instant::now()) - 4112
{ - 4113
return c; - 4114
} - 4115
// `resolve_at` failing here is not a trust decision — an unconditional - 4116
// `true` would let a workspace whose trust prompt an operator declined - 4117
// have its hooks/MCP servers/secret scope applied anyway. Recompute - 4118
// trust the same way `resolve_at` does rather than assuming it. - 4119
vak_core::Core::new_with_trust(cwd.to_path_buf(), vak_core::trust::is_trusted(cwd)) - 4120
.unwrap_or_else(|_| state.core.clone()) - 4121
} - 4122
- 4123
#[derive(serde::Deserialize, Default)] - 4124
struct ListSessionsQuery { - 4125
/// Lists the trash instead: only the sessions moved there, so a person - 4126
/// can restore one. Nothing else reads a trashed session. - 4127
#[serde(default)] - 4128
trash: bool, - 4129
} - 4130
- 4131
/// Sidebar projection over the persisted store: one summary per JSONL file. - 4132
async fn list_sessions( - 4133
State(state): State<AppState>, - 4134
axum::extract::Query(query): axum::extract::Query<ListSessionsQuery>, - 4135
) -> Json<serde_json::Value> { - 4136
// Sessions are stored per workspace, so this follows the workspace the - 4137
// client has open rather than the one the process started in. - 4138
let active = state.active_core(); - 4139
let dir = vak_session::SessionPath::sessions_dir(&state.core.sessions_home(), active.cwd()); - 4140
let active_cwd = active.cwd().to_string_lossy().into_owned(); - 4141
let archive_map = read_archive(&state.core); - 4142
let trashed = vak_core::trash::trashed(&state.core.shared_data_home()); - 4143
let mut sessions = Vec::new(); - 4144
let mut entries: Vec<std::fs::DirEntry> = std::fs::read_dir(&dir) - 4145
.map(|read| read.flatten().collect()) - 4146
.unwrap_or_default(); - 4147
// A workspace can be renamed or canonicalized between runs (notably - 4148
// `/var` vs `/private/var` on macOS). Recover sessions by their durable - 4149
// header cwd when the hashed directory no longer matches, while still - 4150
// filtering strictly to the active workspace. - 4151
if let Ok(projects) = std::fs::read_dir(state.core.sessions_home().join("sessions")) { - 4152
for project in projects.flatten() { - 4153
if let Ok(files) = std::fs::read_dir(project.path()) { - 4154
for file in files.flatten() { - 4155
let duplicate = entries - 4156
.iter() - 4157
.any(|existing| existing.path() == file.path()); - 4158
if !duplicate { - 4159
entries.push(file); - 4160
} - 4161
} - 4162
} - 4163
} - 4164
} - 4165
let shared_home = state.core.shared_data_home(); - 4166
if let Ok(agents) = std::fs::read_dir(shared_home.join("agents")) { - 4167
for agent in agents.flatten() { - 4168
if let Ok(projects) = std::fs::read_dir(agent.path().join("sessions")) { - 4169
for project in projects.flatten() { - 4170
if let Ok(files) = std::fs::read_dir(project.path()) { - 4171
for file in files.flatten() { - 4172
let duplicate = entries - 4173
.iter() - 4174
.any(|existing| existing.path() == file.path()); - 4175
if !duplicate { - 4176
entries.push(file); - 4177
} - 4178
} - 4179
} - 4180
} - 4181
} - 4182
} - 4183
} - 4184
for entry in entries { - 4185
let path = entry.path(); - 4186
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { - 4187
continue; - 4188
} - 4189
let Some(session_id) = path.file_stem().and_then(|s| s.to_str()).map(String::from) else { - 4190
continue; - 4191
}; - 4192
if trashed.contains(&session_id) != query.trash { - 4193
continue; - 4194
} - 4195
let updated_at = std::fs::metadata(&path) - 4196
.ok() - 4197
.and_then(|m| m.modified().ok()) - 4198
.map(|t| chrono::DateTime::<chrono::Utc>::from(t).to_rfc3339()); - 4199
let (created_at, title, entry_count, cwd, agent) = summarize_jsonl(&path); - 4200
// A user-created Agent lives in its own isolated workspace (see - 4201
// agent_chats::open / agent_workspace), independent of whichever - 4202
// default workspace the browser client currently has open — that - 4203
// switch only ever applied to the built-in "vak" identity, which - 4204
// still shares the process's default workspace. So only the "vak" - 4205
// (or header-less/legacy) sessions are filtered by `active_cwd`; - 4206
// every other Agent's sessions are always its own to show. - 4207
let is_default_agent = agent.as_ref().is_none_or(|a| a.id == "vak"); - 4208
if is_default_agent && cwd.as_deref() != Some(active_cwd.as_str()) { - 4209
continue; - 4210
} - 4211
// Header-only sessions are abandoned drafts (for example, creating a - 4212
// task and immediately switching away). Keep the ledger append-only, - 4213
// but do not let empty drafts accumulate in the task switcher. This - 4214
// applies equally to built-in and user-created Agents, but actively - 4215
// registered sessions must remain discoverable. - 4216
let is_active = state.get(&session_id).is_some(); - 4217
if entry_count <= 1 && !is_active { - 4218
continue; - 4219
} - 4220
let running = state.get(&session_id).is_some_and(|handle| { - 4221
handle - 4222
.session - 4223
.lock() - 4224
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4225
.is_none() - 4226
}); - 4227
let archived = archive_map.get(&session_id).copied().unwrap_or(false); - 4228
sessions.push(serde_json::json!({ - 4229
"session_id": session_id, - 4230
"cwd": cwd.unwrap_or_else(|| state.core.cwd().to_string_lossy().into_owned()), - 4231
"created_at": created_at, - 4232
"updated_at": updated_at, - 4233
"entries": entry_count, - 4234
"agent": agent, - 4235
"title": title, - 4236
"running": running, - 4237
"archived": archived, - 4238
})); - 4239
} - 4240
sessions.sort_by_key(|s| s["updated_at"].as_str().unwrap_or("").to_string()); - 4241
sessions.reverse(); - 4242
Json(serde_json::json!({ "sessions": sessions })) - 4243
} - 4244
- 4245
/// Bounded scan: header line for created_at + first user message as title. - 4246
fn summarize_jsonl( - 4247
path: &std::path::Path, - 4248
) -> ( - 4249
Option<String>, - 4250
Option<String>, - 4251
u64, - 4252
Option<String>, - 4253
Option<vak_session::types::AgentIdentity>, - 4254
) { - 4255
use std::io::BufRead; - 4256
let Ok(file) = std::fs::File::open(path) else { - 4257
return (None, None, 0, None, None); - 4258
}; - 4259
let mut reader = std::io::BufReader::new(file); - 4260
let mut created_at = None; - 4261
let mut title = None; - 4262
let mut cwd = None; - 4263
let mut agent = None; - 4264
let mut entries = 0u64; - 4265
let mut line = String::new(); - 4266
loop { - 4267
line.clear(); - 4268
match reader.read_line(&mut line) { - 4269
Ok(0) => break, - 4270
Ok(_) => { - 4271
entries += 1; - 4272
if let Ok(entry) = serde_json::from_str::<vak_session::Entry>(line.trim()) { - 4273
match entry.payload { - 4274
vak_session::EntryPayload::Header(h) => { - 4275
created_at = Some(h.created_at.to_rfc3339()); - 4276
cwd = Some(h.cwd.to_string_lossy().into_owned()); - 4277
agent = h.agent; - 4278
} - 4279
vak_session::EntryPayload::Message(rec) => { - 4280
if title.is_none() - 4281
&& rec.message.role == vak_llm::Role::User - 4282
&& rec.control_kind().is_none() - 4283
{ - 4284
let text = rec.message.text_content(); - 4285
let text = text.trim(); - 4286
if !text.is_empty() { - 4287
let first_line = text.lines().next().unwrap_or(text).trim(); - 4288
let mut snippet: String = first_line.chars().take(72).collect(); - 4289
if first_line.chars().count() > 72 { - 4290
snippet.push('…'); - 4291
} - 4292
title = Some(snippet); - 4293
} - 4294
} - 4295
} - 4296
vak_session::EntryPayload::Compaction(_) => {} - 4297
vak_session::EntryPayload::Receipt(_) => {} - 4298
vak_session::EntryPayload::Goal(_) => {} - 4299
vak_session::EntryPayload::GoalUpdate(_) => {} - 4300
vak_session::EntryPayload::Activity(_) => {} - 4301
vak_session::EntryPayload::Work(_) => {} - 4302
vak_session::EntryPayload::Intent(_) => {} - 4303
vak_session::EntryPayload::TurnCapabilitiesBound(_) - 4304
| vak_session::EntryPayload::TurnCapabilitiesRef(_) => {} - 4305
vak_session::EntryPayload::ChildRun { .. } => {} - 4306
vak_session::EntryPayload::Presentation(_) => {} - 4307
vak_session::EntryPayload::TurnCard(_) => {} - 4308
vak_session::EntryPayload::EvidenceBody(_) => {} - 4309
} - 4310
} - 4311
if title.is_some() && entries > 400 { - 4312
break; - 4313
} - 4314
} - 4315
Err(_) => break, - 4316
} - 4317
} - 4318
(created_at, title, entries, cwd, agent) - 4319
} - 4320
- 4321
#[derive(serde::Deserialize)] - 4322
struct RoutingEnvelope { - 4323
/// The user message that caused this admission. This is metadata, not a - 4324
/// capability or an instruction to the model. - 4325
#[serde(default)] - 4326
message_id: Option<String>, - 4327
#[serde(default)] - 4328
conversation_id: Option<String>, - 4329
#[serde(default)] - 4330
target_work_id: Option<String>, - 4331
#[serde(default)] - 4332
target_result_id: Option<String>, - 4333
/// `independent`, `follow_up`, `correction`, `status`, `cancel`, or - 4334
/// `schedule`; unknown values are retained as provenance but never used - 4335
/// to authorize work. - 4336
#[serde(default)] - 4337
relation: Option<String>, - 4338
#[serde(default)] - 4339
outcome_revision: Option<u64>, - 4340
#[serde(default)] - 4341
provenance: Option<String>, - 4342
} - 4343
- 4344
#[derive(serde::Deserialize)] - 4345
struct RunBody { - 4346
prompt: String, - 4347
/// Stable client identity used to make network retries idempotent. - 4348
#[serde(default)] - 4349
request_id: Option<String>, - 4350
#[serde(default)] - 4351
routing: Option<RoutingEnvelope>, - 4352
/// Optional run-scoped work profile. `managed` creates and persists a - 4353
/// work contract before the agent can execute tools. - 4354
#[serde(default)] - 4355
work_mode: Option<String>, - 4356
/// Optional base64 images appended to the prompt as vision content - 4357
/// (docs/design/22-gateway.md media passthrough). - 4358
#[serde(default)] - 4359
attachments: Vec<RunAttachment>, - 4360
/// Files already saved in the workspace inbox (`POST /fs/inbox`), named - 4361
/// by the paths that route returned. Each reaches the model as a note - 4362
/// saying where it is, never as its bytes (docs/design/72, F1). - 4363
#[serde(default)] - 4364
files: Vec<String>, - 4365
/// Goal mode (docs/design/42-managed-work-contracts.md): durable objective; completion - 4366
/// is audited against `criteria`, never self-reported. - 4367
#[serde(default)] - 4368
goal: Option<String>, - 4369
/// Acceptance criteria for goal mode (`verify:` prefixed criteria run - 4370
/// as brokered shell commands; others are judged from evidence). - 4371
#[serde(default)] - 4372
criteria: Vec<String>, - 4373
} - 4374
- 4375
#[derive(serde::Deserialize)] - 4376
struct RunAttachment { - 4377
#[serde(default = "default_image_mime")] - 4378
mime: String, - 4379
data: String, - 4380
} - 4381
- 4382
fn default_image_mime() -> String { - 4383
"image/png".into() - 4384
} - 4385
- 4386
pub(crate) fn mpsc_to_broadcast(tx: events::EventBus) -> mpsc::Sender<AgentEvent> { - 4387
let (tx_in, mut rx) = mpsc::channel::<AgentEvent>(512); - 4388
tokio::spawn(async move { - 4389
// Forward into the BROADCAST channel (sync send). Forwarding into - 4390
// tx_in would feed the channel back into itself. - 4391
// - 4392
// Headless consumers (gateway turns, cron routines) legitimately run - 4393
// with zero broadcast subscribers; send errors must NEVER tear the - 4394
// pump down — the agent treats a dropped mpsc receiver as a lost - 4395
// consumer and cancels the run mid-flight. - 4396
while let Some(ev) = rx.recv().await { - 4397
let _ = tx.send(ev); - 4398
} - 4399
}); - 4400
tx_in - 4401
} - 4402
- 4403
/// A run cannot start without a working provider credential. - 4404
/// - 4405
/// Every one of these three call sites used to return a bare 503 with no - 4406
/// body and nothing logged, which made "the agent never replied" a - 4407
/// symptom with no server-side trail: a client saw an empty response, an - 4408
/// operator reading gateway.log saw nothing at all, and diagnosing it - 4409
/// meant reading this file. `Core::provider()` already carries a precise - 4410
/// `CoreError::MissingAuth { env, provider }` — this puts it where an - 4411
/// operator and a client can both actually see it. - 4412
/// - 4413
/// The body is `{"error": <message>}`, matching every other handler in - 4414
/// this file. A `{"error": <code>, "detail": <message>}` shape was tried - 4415
/// first and reverted: the desktop frontend's error handling already - 4416
/// reads `.error` as the human-readable string every other endpoint puts - 4417
/// there, so a two-field body would have shown the user the machine code - 4418
/// ("provider_unavailable") instead of the message that says what to fix. - 4419
/// - 4420
/// A missing credential also carries `"kind": "no_ai_service"` (see - 4421
/// `provider_error_body`). - 4422
fn provider_unavailable(err: vak_core::CoreError) -> axum::response::Response { - 4423
use axum::response::IntoResponse; - 4424
eprintln!("[run] refused: {err}"); - 4425
( - 4426
StatusCode::SERVICE_UNAVAILABLE, - 4427
axum::Json(provider_error_body(&err)), - 4428
) - 4429
.into_response() - 4430
} - 4431
- 4432
/// `{"error": <message>}` for a provider failure, plus `"kind": - 4433
/// "no_ai_service"` when the cause is a missing credential, so a client can - 4434
/// say so in plain words without matching on the message text; `error` - 4435
/// stays the precise message an operator or the CLI needs. The one place - 4436
/// that decides the kind, for a refused turn and a model catalogue alike. - 4437
fn provider_error_body(err: &vak_core::CoreError) -> serde_json::Value { - 4438
match err { - 4439
vak_core::CoreError::MissingAuth { .. } => { - 4440
serde_json::json!({ "error": err.to_string(), "kind": "no_ai_service" }) - 4441
} - 4442
_ => serde_json::json!({ "error": err.to_string() }), - 4443
} - 4444
} - 4445
- 4446
/// `run_prompt`, `side_chat`, and `start_bestofn` all fall back to - 4447
/// `provider_unavailable` when `Core::provider()` fails; this pins the - 4448
/// response it produces so a regression — an empty body, or the - 4449
/// `{"error": <code>, "detail": <message>}` shape tried and reverted - 4450
/// above — fails a fast unit test instead of surfacing as "the agent - 4451
/// never replied" with nothing in gateway.log to explain why. - 4452
#[cfg(test)] - 4453
#[allow(clippy::unwrap_used, clippy::expect_used)] - 4454
mod provider_unavailable_tests { - 4455
use super::provider_unavailable; - 4456
use axum::response::IntoResponse as _; - 4457
use http_body_util::BodyExt as _; - 4458
- 4459
#[tokio::test] - 4460
async fn reports_status_and_a_body_naming_the_missing_credential() { - 4461
let err = vak_core::CoreError::MissingAuth { - 4462
env: "ANTHROPIC_API_KEY".into(), - 4463
provider: "anthropic".into(), - 4464
}; - 4465
let response = provider_unavailable(err).into_response(); - 4466
assert_eq!( - 4467
response.status(), - 4468
axum::http::StatusCode::SERVICE_UNAVAILABLE - 4469
); - 4470
- 4471
let bytes = response - 4472
.into_body() - 4473
.collect() - 4474
.await - 4475
.expect("body readable") - 4476
.to_bytes(); - 4477
assert!( - 4478
!bytes.is_empty(), - 4479
"body must not be empty — that was the original bug" - 4480
); - 4481
- 4482
let body: serde_json::Value = serde_json::from_slice(&bytes).expect("body is JSON"); - 4483
// `error` carries the human-readable message, same as every other - 4484
// handler in this file — not a machine code with the message hidden - 4485
// in a `detail` the frontend never reads — and `kind` types it. - 4486
let fields: Vec<&String> = body.as_object().expect("object body").keys().collect(); - 4487
assert_eq!( - 4488
fields, - 4489
vec!["error", "kind"], - 4490
"body must have the `error` message and its `kind`" - 4491
); - 4492
assert_eq!(body["kind"], "no_ai_service"); - 4493
let message = body["error"].as_str().expect("error is a string"); - 4494
assert!( - 4495
message.contains("ANTHROPIC_API_KEY"), - 4496
"message must name the env var to set, got: {message}" - 4497
); - 4498
assert!( - 4499
message.contains("anthropic"), - 4500
"message must name the provider, got: {message}" - 4501
); - 4502
} - 4503
- 4504
#[tokio::test] - 4505
async fn only_a_missing_credential_is_typed_as_no_ai_service() { - 4506
let err = vak_core::CoreError::InvalidConfig("bad route".into()); - 4507
let bytes = provider_unavailable(err) - 4508
.into_response() - 4509
.into_body() - 4510
.collect() - 4511
.await - 4512
.expect("body readable") - 4513
.to_bytes(); - 4514
let body: serde_json::Value = serde_json::from_slice(&bytes).expect("body is JSON"); - 4515
let fields: Vec<&String> = body.as_object().expect("object body").keys().collect(); - 4516
assert_eq!(fields, vec!["error"], "other refusals carry no kind"); - 4517
} - 4518
} - 4519
- 4520
// ---- Shared turn-chain executor (invariant 30; docs/design/ - 4521
// 64-agent-owned-platform.md, "Request durability and delivery") ---------- - 4522
// - 4523
// `run_prompt`, `send_steering`, and `gateway::execute_turn_chain` all - 4524
// admit a prompt, run it, and — if more input arrived while the run was - 4525
// settling — keep going rather than silently stranding it. Before this, - 4526
// each surface implemented that loop separately: the HTTP path did not - 4527
// implement it at all (steering queued after a run's last internal drain - 4528
// was never picked back up), and the gateway's own version restored - 4529
// `handle.session` before draining, leaving a race window where a - 4530
// concurrent admission could steal the ledger. `admit_or_queue` and - 4531
// `continue_or_release` are the one busy/idle decision, in both - 4532
// directions; `run_turn_chain` is the one loop that runs a leg and decides - 4533
// whether to continue, parameterized by approver and run kind so each - 4534
// surface keeps its own settle bookkeeping (durable activity records vs. - 4535
// reply channel + rendered text) without duplicating the loop mechanics. - 4536
- 4537
/// What a turn chain's FIRST leg runs. Every leg after the first is always - 4538
/// a plain message turn: draining `handle.steering` only ever produces a - 4539
/// `vak_llm::Message` via `SteeringQueues::merge_prompt`, never a fresh - 4540
/// goal/managed/auto request — that is `/run`'s own admission, which a - 4541
/// queued steering message never claims to be. - 4542
enum TurnStart { - 4543
/// The person's message as recorded, with its metadata (attached files). - 4544
Message(vak_session::MessageRecord), - 4545
Managed(String), - 4546
Auto(String), - 4547
Goal { - 4548
prompt: String, - 4549
objective: String, - 4550
criteria: Vec<String>, - 4551
}, - 4552
} - 4553
- 4554
impl TurnStart { - 4555
fn message(message: vak_llm::Message) -> Self { - 4556
TurnStart::Message(vak_session::MessageRecord { - 4557
message, - 4558
meta: None, - 4559
}) - 4560
} - 4561
- 4562
/// The message this leg would present — used both to seed the preview - 4563
/// intent before a run starts and, on the busy path, as the queued - 4564
/// steering entry (attachments and all; invariant 1, model-visible - 4565
/// input is never degraded to bare text). - 4566
fn preview_message(&self) -> vak_llm::Message { - 4567
match self { - 4568
TurnStart::Message(m) => m.message.clone(), - 4569
TurnStart::Managed(p) | TurnStart::Auto(p) => vak_llm::Message::user_text(p), - 4570
TurnStart::Goal { prompt, .. } => vak_llm::Message::user_text(prompt), - 4571
} - 4572
} - 4573
- 4574
#[allow(clippy::too_many_arguments)] - 4575
async fn run( - 4576
self, - 4577
core: &Core, - 4578
session: SessionLog, - 4579
cancel: CancellationToken, - 4580
approver: Arc<dyn Approver>, - 4581
steering: Arc<SteeringQueues>, - 4582
events: mpsc::Sender<AgentEvent>, - 4583
) -> Result<(vak_agent::TurnOutcome, SessionLog), vak_core::CoreError> { - 4584
match self { - 4585
TurnStart::Message(m) => { - 4586
core.run_turn_with_message( - 4587
session, - 4588
m, - 4589
cancel, - 4590
Some(approver), - 4591
None, - 4592
Some(steering), - 4593
events, - 4594
) - 4595
.await - 4596
} - 4597
TurnStart::Managed(prompt) => { - 4598
core.run_managed_turn_with( - 4599
session, - 4600
&prompt, - 4601
cancel, - 4602
Some(approver), - 4603
None, - 4604
Some(steering), - 4605
events, - 4606
) - 4607
.await - 4608
} - 4609
TurnStart::Auto(prompt) => { - 4610
core.run_auto_turn_with( - 4611
session, - 4612
&prompt, - 4613
cancel, - 4614
Some(approver), - 4615
None, - 4616
Some(steering), - 4617
events, - 4618
) - 4619
.await - 4620
} - 4621
TurnStart::Goal { - 4622
prompt, - 4623
objective, - 4624
criteria, - 4625
} => { - 4626
core.run_goal_turn_with( - 4627
session, - 4628
&prompt, - 4629
&objective, - 4630
criteria, - 4631
cancel, - 4632
Some(approver), - 4633
None, - 4634
Some(steering), - 4635
events, - 4636
) - 4637
.await - 4638
} - 4639
} - 4640
} - 4641
} - 4642
- 4643
/// Result of admitting input at the busy boundary. - 4644
enum Admission { - 4645
/// The ledger was idle; the caller now owns it and must run a chain. - 4646
Started(SessionLog), - 4647
/// Busy: `message` was pushed onto `handle.steering` durably. - 4648
Queued, - 4649
/// Busy, and this input kind (goal/managed/auto) cannot be queued. - 4650
RejectedBusy, - 4651
} - 4652
- 4653
/// Admit input at the busy boundary shared by `/run` and `/steering`: busy - 4654
/// input is queued durably or explicitly rejected, and a new request is - 4655
/// never acknowledged and discarded (finding 1). Locking `handle.session` - 4656
/// for the whole decision is what makes it safe against a chain settling in - 4657
/// [`continue_or_release`] at the same instant — the two can never observe - 4658
/// a window where the ledger looks idle to one caller and busy to the - 4659
/// other. - 4660
fn admit_or_queue( - 4661
handle: &SessionHandle, - 4662
restricted: bool, - 4663
message: vak_llm::Message, - 4664
) -> Admission { - 4665
let mut slot = handle - 4666
.session - 4667
.lock() - 4668
.unwrap_or_else(std::sync::PoisonError::into_inner); - 4669
if let Some(taken) = slot.take() { - 4670
return Admission::Started(taken); - 4671
} - 4672
drop(slot); - 4673
if restricted { - 4674
return Admission::RejectedBusy; - 4675
} - 4676
handle.steering.push_steering_message(message); - 4677
Admission::Queued - 4678
} - 4679
- 4680
/// After a leg settles, atomically decide whether the chain continues. - 4681
/// Mirrors [`admit_or_queue`]'s locking discipline from the other - 4682
/// direction: the ledger is written back to `handle.session` (marking the - 4683
/// session idle again) only when nothing is queued, so steering that - 4684
/// arrives while a leg is settling either lands in THIS drain or is queued - 4685
/// against a session that is genuinely idle once this returns — never a - 4686
/// session that looks idle while the drain that would have picked it up - 4687
/// already happened and is gone. - 4688
fn continue_or_release( - 4689
handle: &SessionHandle, - 4690
ledger: SessionLog, - 4691
) -> Option<(SessionLog, vak_llm::Message)> { - 4692
let mut slot = handle - 4693
.session - 4694
.lock() - 4695
.unwrap_or_else(std::sync::PoisonError::into_inner); - 4696
let queued = handle.steering.drain(vak_agent::DrainMode::All); - 4697
match SteeringQueues::merge_prompt(queued) { - 4698
None => { - 4699
*slot = Some(ledger); - 4700
None - 4701
} - 4702
Some(merged) => Some((ledger, merged)), - 4703
} - 4704
} - 4705
- 4706
/// Give an SSE consumer a moment to attach before a freshly admitted chain - 4707
/// starts, so its terminal event is seen — but only when nobody is - 4708
/// watching yet. `handle.subscribed` is a single-permit `Notify`: the first - 4709
/// SSE connection's `notify_one()` satisfies exactly one `.notified()` - 4710
/// call, so unconditionally waiting here made every leg after the first - 4711
/// block for the full 2s even with a client already attached (finding 2). - 4712
async fn wait_for_external_subscriber(handle: &SessionHandle) { - 4713
if handle.events_tx.external_subscribers() == 0 { - 4714
let _ = tokio::time::timeout(Duration::from_secs(2), handle.subscribed.notified()).await; - 4715
} - 4716
} - 4717
- 4718
/// What a leg's surface-specific settle step hands back to - 4719
/// [`run_turn_chain`]: the ledger to keep running with (`None` when a - 4720
/// `CoreError` left nothing recoverable), and the `(summary, is_error)` - 4721
/// pair the executor broadcasts as this leg's `RunFinished`. - 4722
type SettleResult = (Option<SessionLog>, String, bool); - 4723
- 4724
/// The one turn-chain executor shared by the HTTP path (`run_prompt`, - 4725
/// `send_steering`) and the gateway (`gateway::execute_turn_chain`). - 4726
/// `start` runs as the first leg; `settle` performs the surface-specific - 4727
/// bookkeeping for each leg's outcome — durable activity + presentation - 4728
/// snapshot + hub summary for HTTP, reply channel + rendered text + - 4729
/// reflection logging for the gateway — and returns the ledger to continue - 4730
/// with (reflection itself stays inside `settle`, since the two surfaces - 4731
/// react to its outcome differently; see `http_settle` and - 4732
/// `gateway::execute_turn_chain`). After settling, steering queued in the - 4733
/// meantime is drained and merged into a continuation leg via - 4734
/// [`continue_or_release`] — under the same lock that decides whether the - 4735
/// ledger goes back to `handle.session` — rather than being silently - 4736
/// stranded once the chain looks idle again (finding 1). - 4737
async fn run_turn_chain<F, Fut>( - 4738
core: Core, - 4739
handle: Arc<SessionHandle>, - 4740
mut taken: SessionLog, - 4741
mut start: TurnStart, - 4742
approver_factory: impl Fn(&str) -> Arc<dyn Approver>, - 4743
mut settle: F, - 4744
) where - 4745
F: FnMut(&str, Result<(vak_agent::TurnOutcome, SessionLog), vak_core::CoreError>) -> Fut, - 4746
Fut: std::future::Future<Output = SettleResult>, - 4747
{ - 4748
loop { - 4749
let session_id = taken - 4750
.header() - 4751
.map(|h| h.session_id.clone()) - 4752
.unwrap_or_else(|| handle.id.clone()); - 4753
let approver = approver_factory(&session_id); - 4754
let events = mpsc_to_broadcast(handle.events_tx.clone()); - 4755
let steering = handle.steering.clone(); - 4756
let cancel = handle - 4757
.cancel - 4758
.lock() - 4759
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4760
.clone(); - 4761
let outcome = start - 4762
.run(&core, taken, cancel, approver, steering, events) - 4763
.await; - 4764
// Reset the token so the next leg is not born already-cancelled. - 4765
*handle - 4766
.cancel - 4767
.lock() - 4768
.unwrap_or_else(std::sync::PoisonError::into_inner) = CancellationToken::new(); - 4769
- 4770
let (ledger, summary, is_error) = settle(&session_id, outcome).await; - 4771
let _ = handle - 4772
.events_tx - 4773
.send(AgentEvent::RunFinished { summary, is_error }); - 4774
- 4775
let Some(ledger) = ledger else { - 4776
return; - 4777
}; - 4778
- 4779
match continue_or_release(&handle, ledger) { - 4780
None => return, - 4781
Some((ledger, merged)) => { - 4782
taken = ledger; - 4783
start = TurnStart::message(merged); - 4784
} - 4785
} - 4786
} - 4787
} - 4788
- 4789
/// `run_prompt`'s per-leg settle: durable "Run finished" activity, buffered - 4790
/// activity flush (with the same request_id admissions cleanup on BOTH the - 4791
/// success and error path — the error path used to skip it, leaking an - 4792
/// admissions entry for any steering that had been buffered before the - 4793
/// failure), presentation snapshot, hub summary, and FTS indexing. - 4794
async fn http_settle( - 4795
handle: Arc<SessionHandle>, - 4796
hub: events::EventHub, - 4797
admin_store: Option<vak_store::Store>, - 4798
sessions_home: std::path::PathBuf, - 4799
run_id: String, - 4800
outcome: Result<(vak_agent::TurnOutcome, SessionLog), vak_core::CoreError>, - 4801
) -> SettleResult { - 4802
fn flush_buffered(handle: &SessionHandle, log: &mut SessionLog) { - 4803
let buffered = std::mem::take( - 4804
&mut *handle - 4805
.activity_buffer - 4806
.lock() - 4807
.unwrap_or_else(std::sync::PoisonError::into_inner), - 4808
); - 4809
for activity in buffered { - 4810
let request_id = activity.data.get("request_id").cloned(); - 4811
let _ = log.append_activity(activity); - 4812
if let Some(request_id) = request_id { - 4813
handle - 4814
.admissions - 4815
.lock() - 4816
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4817
.remove(&request_id); - 4818
} - 4819
} - 4820
} - 4821
- 4822
match outcome { - 4823
Ok((o, mut session_log)) => { - 4824
let (summary, is_error) = match &o { - 4825
vak_agent::TurnOutcome::Completed { .. } => ("completed".to_string(), false), - 4826
vak_agent::TurnOutcome::Aborted { .. } => ("aborted".to_string(), false), - 4827
vak_agent::TurnOutcome::Failed { error } => (format!("failed: {error}"), true), - 4828
vak_agent::TurnOutcome::MaxTurnsReached => ("max_turns".to_string(), true), - 4829
}; - 4830
let activity_status = match &o { - 4831
vak_agent::TurnOutcome::Completed { .. } => vak_session::ActivityStatus::Succeeded, - 4832
vak_agent::TurnOutcome::Aborted { .. } => vak_session::ActivityStatus::Cancelled, - 4833
vak_agent::TurnOutcome::Failed { .. } => vak_session::ActivityStatus::Failed, - 4834
vak_agent::TurnOutcome::MaxTurnsReached => vak_session::ActivityStatus::Partial, - 4835
}; - 4836
let _ = session_log.append_activity(vak_session::ActivityRecord { - 4837
activity_id: format!("run-{run_id}-{}", chrono::Utc::now().timestamp_micros()), - 4838
turn: None, - 4839
kind: vak_session::ActivityKind::Run, - 4840
status: activity_status, - 4841
label: "Run finished".into(), - 4842
detail: Some(summary.clone()), - 4843
data: std::collections::BTreeMap::new(), - 4844
}); - 4845
flush_buffered(&handle, &mut session_log); - 4846
*handle - 4847
.presentation - 4848
.lock() - 4849
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 4850
live_presentation_snapshot(&handle.core, &run_id, &session_log); - 4851
hub.emit_agent_summary(&summary, Some(run_id.clone())); - 4852
index_session_later(admin_store, sessions_home, run_id.clone()); - 4853
// Background reflection seam (docs/design/29 P1): after the - 4854
// summary is recorded and while this leg still owns the ledger - 4855
// (a second in-process handle cannot take the file lock). - 4856
// Bounded; the result is deliberately ignored — a completed run - 4857
// never fails on reflection. - 4858
if !is_error && handle.core.config().memory.reflection { - 4859
let _ = tokio::time::timeout( - 4860
REFLECTION_CALL_TIMEOUT, - 4861
handle.core.reflect_after_turn(&session_log, ""), - 4862
) - 4863
.await; - 4864
} - 4865
(Some(session_log), summary, is_error) - 4866
} - 4867
Err(e) => { - 4868
// Same leak class: restore from the durable ledger so the - 4869
// handle does not stay wedged on "run in progress". - 4870
let restored = reopen_ledger(&handle.core, &run_id).map(|mut restored| { - 4871
flush_buffered(&handle, &mut restored); - 4872
*handle - 4873
.presentation - 4874
.lock() - 4875
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 4876
live_presentation_snapshot(&handle.core, &run_id, &restored); - 4877
restored - 4878
}); - 4879
(restored, format!("error: {e}"), true) - 4880
} - 4881
} - 4882
} - 4883
- 4884
/// Spawn an HTTP-surfaced turn chain: an `HttpApprover` built fresh per leg, - 4885
/// `http_settle` bookkeeping, and the `request_id` admissions cleanup once - 4886
/// the WHOLE chain (every leg, not just the first) has settled. Shared by - 4887
/// `run_prompt` and `send_steering`'s own idle-admission path (finding 1b) - 4888
/// — a steer that lands on an idle session IS a fresh admission, not inert - 4889
/// queued input nothing will ever look at again. - 4890
fn spawn_http_turn_chain( - 4891
state: &AppState, - 4892
handle: Arc<SessionHandle>, - 4893
core: Core, - 4894
taken: SessionLog, - 4895
start: TurnStart, - 4896
request_id: Option<String>, - 4897
) { - 4898
let hub = state.hub.clone(); - 4899
let admin_store = state.store.clone(); - 4900
let sessions_home = state.core.sessions_home(); - 4901
let chain_handle = handle; - 4902
tokio::spawn(async move { - 4903
let approver_handle = chain_handle.clone(); - 4904
let settle_handle = chain_handle.clone(); - 4905
run_turn_chain( - 4906
core, - 4907
chain_handle.clone(), - 4908
taken, - 4909
start, - 4910
move |_leg_session_id: &str| -> Arc<dyn Approver> { - 4911
// Driven by a client that is holding the SSE stream open, - 4912
// so a gate raised here reaches a person. - 4913
Arc::new(HttpApprover { - 4914
events_tx: approver_handle.events_tx.clone(), - 4915
pending: approver_handle.pending.clone(), - 4916
session_id: approver_handle.id.clone(), - 4917
activity_buffer: approver_handle.activity_buffer.clone(), - 4918
answerable: true, - 4919
}) - 4920
}, - 4921
move |leg_session_id: &str, outcome| { - 4922
let handle = settle_handle.clone(); - 4923
let hub = hub.clone(); - 4924
let admin_store = admin_store.clone(); - 4925
let sessions_home = sessions_home.clone(); - 4926
let run_id = leg_session_id.to_string(); - 4927
async move { - 4928
http_settle(handle, hub, admin_store, sessions_home, run_id, outcome).await - 4929
} - 4930
}, - 4931
) - 4932
.await; - 4933
- 4934
if let Some(request_id) = request_id { - 4935
chain_handle - 4936
.admissions - 4937
.lock() - 4938
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4939
.remove(&request_id); - 4940
} - 4941
}); - 4942
} - 4943
- 4944
async fn run_prompt( - 4945
State(state): State<AppState>, - 4946
Path(id): Path<String>, - 4947
Json(body): Json<RunBody>, - 4948
) -> axum::response::Response { - 4949
refresh_control_plane(&state); - 4950
use axum::response::IntoResponse; - 4951
let handle = match ensure_session_handle(&state, &id).await { - 4952
Ok((_, handle)) => handle, - 4953
Err(e) => { - 4954
return ( - 4955
StatusCode::NOT_FOUND, - 4956
Json(serde_json::json!({ "error": e.to_string() })), - 4957
) - 4958
.into_response(); - 4959
} - 4960
}; - 4961
if let Some(routing) = body.routing.as_ref() - 4962
&& let Some(expected) = routing.outcome_revision - 4963
&& !routing_revision_is_current(&handle, expected) - 4964
{ - 4965
return ( - 4966
StatusCode::CONFLICT, - 4967
Json(serde_json::json!({ - 4968
"error": "target result is stale; refresh before continuing", - 4969
"target_revision": expected, - 4970
})), - 4971
) - 4972
.into_response(); - 4973
} - 4974
let request_id = body.request_id.clone(); - 4975
if let Some(request_id) = request_id.as_deref() - 4976
&& handle - 4977
.admissions - 4978
.lock() - 4979
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4980
.contains(request_id) - 4981
{ - 4982
return ( - 4983
StatusCode::ACCEPTED, - 4984
Json(serde_json::json!({"request_id": request_id, "state": "duplicate"})), - 4985
) - 4986
.into_response(); - 4987
} - 4988
- 4989
// ---- Validate the request shape before any side effect (finding 4): - 4990
// no durable admission activity, no admissions-set insertion, and no - 4991
// synthesized `RunFinished` for input that never starts a run. ---- - 4992
if body.goal.is_some() && !(body.attachments.is_empty() && body.files.is_empty()) { - 4993
return ( - 4994
StatusCode::BAD_REQUEST, - 4995
Json(serde_json::json!({"error": "goal runs do not support attachments"})), - 4996
) - 4997
.into_response(); - 4998
} - 4999
if let Some(mode) = body.work_mode.as_deref() - 5000
&& !matches!(mode, "direct" | "managed" | "auto")
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.