- 5164
Json(serde_json::json!({"request_id": request_id, "state": "duplicate"})), - 5165
) - 5166
.into_response(); - 5167
} - 5168
if let Some(request_id) = request_id.as_deref() { - 5169
let mut data = std::collections::BTreeMap::new(); - 5170
data.insert("request_id".into(), request_id.to_owned()); - 5171
if let Some(routing) = body.routing.as_ref() { - 5172
record_routing_data(&mut data, routing); - 5173
} - 5174
if taken - 5175
.append_activity(vak_session::ActivityRecord { - 5176
activity_id: format!("admission-{request_id}"), - 5177
turn: None, - 5178
kind: vak_session::ActivityKind::Run, - 5179
status: vak_session::ActivityStatus::Running, - 5180
label: "Request accepted".into(), - 5181
detail: None, - 5182
data, - 5183
}) - 5184
.is_err() - 5185
{ - 5186
*handle - 5187
.session - 5188
.lock() - 5189
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 5190
return ( - 5191
StatusCode::INTERNAL_SERVER_ERROR, - 5192
axum::Json( - 5193
serde_json::json!({"error": "could not durably record request admission"}), - 5194
), - 5195
) - 5196
.into_response(); - 5197
} - 5198
handle - 5199
.admissions - 5200
.lock() - 5201
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5202
.insert(request_id.to_owned()); - 5203
} - 5204
- 5205
// Preview intent so `/control-state` has something to report even - 5206
// before the ledger reflects a real IntentRecord entry for this leg. - 5207
let core = handle.core.clone(); - 5208
let preview_intent = core.resolve_turn_intent(&taken, &queue_message); - 5209
let mut preview_outcome = - 5210
vak_intent::OutcomeSpec::from_intent(&expanded_prompt, &preview_intent); - 5211
preview_outcome.evidence_max_age_secs = Some(core.effective_evidence_max_age_secs()); - 5212
*handle - 5213
.intent - 5214
.lock() - 5215
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 5216
Some(vak_session::types::IntentRecord { - 5217
reading: preview_intent.reading, - 5218
strands: preview_intent.strands, - 5219
engagement: preview_intent.engagement, - 5220
provenance: preview_intent.provenance, - 5221
outcome: Some(preview_outcome), - 5222
model_visible: None, - 5223
commitment_id: None, - 5224
strand_commitments: Default::default(), - 5225
}); - 5226
- 5227
// Give SSE consumers a moment to attach so terminal events are seen — - 5228
// but only when nobody is watching yet (finding 2). - 5229
wait_for_external_subscriber(&handle).await; - 5230
- 5231
let response_request_id = request_id.clone(); - 5232
spawn_http_turn_chain(&state, handle, core, taken, start, request_id); - 5233
- 5234
( - 5235
StatusCode::ACCEPTED, - 5236
Json(serde_json::json!({"request_id": response_request_id, "state": "started"})), - 5237
) - 5238
.into_response() - 5239
} - 5240
- 5241
#[derive(serde::Deserialize)] - 5242
struct SteeringBody { - 5243
text: String, - 5244
/// Caller-owned id used to recover a retry without enqueuing duplicate - 5245
/// steering input. Older callers may omit it; the server then generates - 5246
/// one for the single attempt. - 5247
#[serde(default)] - 5248
request_id: Option<String>, - 5249
#[serde(default)] - 5250
routing: Option<RoutingEnvelope>, - 5251
/// Origin is metadata for the audit trail, never an authority grant. - 5252
#[serde(default = "default_intervention_source")] - 5253
source: String, - 5254
/// Optional base64 images appended to the steered prompt, mirroring - 5255
/// /run so queued input is never degraded to bare text. - 5256
#[serde(default)] - 5257
attachments: Vec<RunAttachment>, - 5258
} - 5259
- 5260
fn default_intervention_source() -> String { - 5261
"human".into() - 5262
} - 5263
- 5264
/// Who is behind an HTTP intervention. - 5265
/// - 5266
/// The authenticated loopback client is the operator's own surface, so the - 5267
/// default is a person on that surface. A caller may declare itself *lower* - 5268
/// — `agent` (one of ours, steering a child) or `system` (an integration) — - 5269
/// and the control plane then treats it accordingly. Free text never grants - 5270
/// authority; the source and the explicit command do - 5271
/// (docs/design/47-commitment-kernel.md, control plane). - 5272
fn control_source_for(declared: &str, surface: &str, target: &str) -> vak_intent::ControlSource { - 5273
match declared.trim().to_ascii_lowercase().as_str() { - 5274
// An agent reaching a session over HTTP is, as far as the control - 5275
// plane can tell, acting on the session it names — not on a child - 5276
// it dispatched. Worker control goes through the worker endpoints, - 5277
// which check parentage. - 5278
"agent" | "worker" => vak_intent::ControlSource::Agent { - 5279
session_id: target.to_string(), - 5280
parent_session_id: None, - 5281
}, - 5282
"system" | "webhook" | "cron" | "integration" => vak_intent::ControlSource::System { - 5283
origin: declared.trim().to_string(), - 5284
}, - 5285
_ => vak_intent::ControlSource::Human { - 5286
surface: surface.to_string(), - 5287
principal: None, - 5288
}, - 5289
} - 5290
} - 5291
- 5292
/// Kind and text for a message on a control endpoint: an explicit command, - 5293
/// or steering text. - 5294
fn intervention_of(text: &str) -> (vak_intent::InterventionKind, String) { - 5295
match vak_intent::parse_command(text) { - 5296
Some(command) => { - 5297
let kind = command.intervention_kind(); - 5298
let text = command - 5299
.text() - 5300
.map(str::to_string) - 5301
.unwrap_or_else(|| text.to_string()); - 5302
(kind, text) - 5303
} - 5304
None => (vak_intent::InterventionKind::Steer, text.to_string()), - 5305
} - 5306
} - 5307
- 5308
async fn send_steering( - 5309
State(state): State<AppState>, - 5310
Path(id): Path<String>, - 5311
Json(body): Json<SteeringBody>, - 5312
) -> axum::response::Response { - 5313
let Some(handle) = state.get(&id) else { - 5314
return StatusCode::NOT_FOUND.into_response(); - 5315
}; - 5316
if let Some(routing) = body.routing.as_ref() - 5317
&& let Some(expected) = routing.outcome_revision - 5318
&& !routing_revision_is_current(&handle, expected) - 5319
{ - 5320
return ( - 5321
StatusCode::CONFLICT, - 5322
Json(serde_json::json!({ - 5323
"error": "target result is stale; refresh before continuing", - 5324
"target_revision": expected, - 5325
})), - 5326
) - 5327
.into_response(); - 5328
} - 5329
let request_id = body - 5330
.request_id - 5331
.as_deref() - 5332
.map(str::trim) - 5333
.filter(|value| !value.is_empty()) - 5334
.map(ToOwned::to_owned) - 5335
.unwrap_or_else(|| format!("intervention-{}", uuid::Uuid::now_v7())); - 5336
let already_admitted = handle - 5337
.admissions - 5338
.lock() - 5339
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5340
.contains(&request_id) - 5341
|| handle - 5342
.session - 5343
.lock() - 5344
.ok() - 5345
.and_then(|guard| { - 5346
guard - 5347
.as_ref() - 5348
.map(|log| log.has_request_admission(&request_id)) - 5349
}) - 5350
.unwrap_or(false); - 5351
if already_admitted { - 5352
return ( - 5353
StatusCode::ACCEPTED, - 5354
Json(serde_json::json!({ - 5355
"request_id": request_id, - 5356
"decision": "duplicate", - 5357
"state": "already_admitted", - 5358
})), - 5359
) - 5360
.into_response(); - 5361
} - 5362
handle - 5363
.admissions - 5364
.lock() - 5365
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5366
.insert(request_id.clone()); - 5367
let (kind, command_text) = intervention_of(&body.text); - 5368
let evaluation = vak_intent::evaluate_intervention(vak_intent::InterventionRequest { - 5369
request_id: request_id.clone(), - 5370
kind: kind.clone(), - 5371
text: command_text, - 5372
source: control_source_for(&body.source, state.core.surface().slug(), &id), - 5373
target_revision: None, - 5374
target_session_id: Some(id.clone()), - 5375
target_parent_session_id: target_parent_of(&handle), - 5376
}); - 5377
record_activity_or_buffer( - 5378
&handle, - 5379
vak_session::ActivityRecord { - 5380
activity_id: evaluation.request.request_id.clone(), - 5381
turn: None, - 5382
kind: vak_session::ActivityKind::Diagnostic, - 5383
status: if evaluation.decision == vak_intent::InterventionDecision::Queued { - 5384
vak_session::ActivityStatus::Pending - 5385
} else { - 5386
vak_session::ActivityStatus::Succeeded - 5387
}, - 5388
label: if evaluation.decision == vak_intent::InterventionDecision::Queued { - 5389
"Intervention queued" - 5390
} else { - 5391
"Intervention accepted" - 5392
} - 5393
.into(), - 5394
detail: Some(body.text.clone()), - 5395
data: std::collections::BTreeMap::from([ - 5396
("kind".into(), kind.as_str().into()), - 5397
("decision".into(), evaluation.decision.as_str().into()), - 5398
("source".into(), body.source.clone()), - 5399
("reason".into(), evaluation.reason.clone()), - 5400
("routing".into(), routing_summary(body.routing.as_ref())), - 5401
]), - 5402
}, - 5403
); - 5404
if matches!( - 5405
evaluation.decision, - 5406
vak_intent::InterventionDecision::RequiresHuman - 5407
| vak_intent::InterventionDecision::Rejected - 5408
) { - 5409
handle - 5410
.admissions - 5411
.lock() - 5412
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5413
.remove(&request_id); - 5414
let status = if evaluation.decision == vak_intent::InterventionDecision::Rejected { - 5415
StatusCode::FORBIDDEN - 5416
} else { - 5417
StatusCode::CONFLICT - 5418
}; - 5419
return ( - 5420
status, - 5421
Json(serde_json::json!({ - 5422
"request_id": evaluation.request.request_id, - 5423
"decision": evaluation.decision.as_str(), - 5424
"reason": evaluation.reason, - 5425
})), - 5426
) - 5427
.into_response(); - 5428
} - 5429
match kind { - 5430
vak_intent::InterventionKind::Replan - 5431
| vak_intent::InterventionKind::Reprioritize - 5432
| vak_intent::InterventionKind::AddRequirement - 5433
| vak_intent::InterventionKind::RemoveRequirement => { - 5434
handle - 5435
.admissions - 5436
.lock() - 5437
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5438
.remove(&request_id); - 5439
return plan_change( - 5440
State(state.clone()), - 5441
Path(id), - 5442
Json(PlanChangeBody { - 5443
text: body.text.clone(), - 5444
source: body.source.clone(), - 5445
target_revision: None, - 5446
}), - 5447
) - 5448
.await; - 5449
} - 5450
vak_intent::InterventionKind::Status => { - 5451
handle - 5452
.admissions - 5453
.lock() - 5454
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5455
.remove(&request_id); - 5456
return ( - 5457
StatusCode::ACCEPTED, - 5458
Json(serde_json::json!({ - 5459
"request_id": evaluation.request.request_id, - 5460
"decision": evaluation.decision.as_str(), - 5461
"state": "status_requested", - 5462
})), - 5463
) - 5464
.into_response(); - 5465
} - 5466
vak_intent::InterventionKind::Pause => { - 5467
handle - 5468
.admissions - 5469
.lock() - 5470
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5471
.remove(&request_id); - 5472
handle.steering.pause(); - 5473
record_control_activity(&handle, "Run paused", "pause"); - 5474
return ( - 5475
StatusCode::ACCEPTED, - 5476
Json(serde_json::json!({ - 5477
"request_id": evaluation.request.request_id, - 5478
"decision": evaluation.decision.as_str(), - 5479
"state": "paused", - 5480
"reason": evaluation.reason, - 5481
})), - 5482
) - 5483
.into_response(); - 5484
} - 5485
vak_intent::InterventionKind::Resume => { - 5486
handle - 5487
.admissions - 5488
.lock() - 5489
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5490
.remove(&request_id); - 5491
handle.steering.resume(); - 5492
record_control_activity(&handle, "Run resumed", "resume"); - 5493
return ( - 5494
StatusCode::ACCEPTED, - 5495
Json(serde_json::json!({ - 5496
"request_id": evaluation.request.request_id, - 5497
"decision": evaluation.decision.as_str(), - 5498
"state": "resumed", - 5499
"reason": evaluation.reason, - 5500
})), - 5501
) - 5502
.into_response(); - 5503
} - 5504
vak_intent::InterventionKind::Cancel => { - 5505
handle - 5506
.admissions - 5507
.lock() - 5508
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5509
.remove(&request_id); - 5510
handle - 5511
.cancel - 5512
.lock() - 5513
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5514
.cancel(); - 5515
record_control_activity(&handle, "Run cancelled", "cancel"); - 5516
return ( - 5517
StatusCode::ACCEPTED, - 5518
Json(serde_json::json!({ - 5519
"request_id": evaluation.request.request_id, - 5520
"decision": evaluation.decision.as_str(), - 5521
"state": "cancelled", - 5522
"reason": evaluation.reason, - 5523
})), - 5524
) - 5525
.into_response(); - 5526
} - 5527
_ => {} - 5528
} - 5529
let usable: Vec<&RunAttachment> = body - 5530
.attachments - 5531
.iter() - 5532
.filter(|a| !a.data.trim().is_empty()) - 5533
.collect(); - 5534
let message = if usable.is_empty() { - 5535
vak_llm::Message::user_text(body.text.clone()) - 5536
} else { - 5537
let mut blocks = vec![vak_llm::ContentBlock::text(body.text.clone())]; - 5538
for a in usable { - 5539
blocks.push(vak_llm::ContentBlock::image_base64( - 5540
a.mime.clone(), - 5541
a.data.trim().to_string(), - 5542
)); - 5543
} - 5544
vak_llm::Message { - 5545
role: vak_llm::Role::User, - 5546
content: blocks, - 5547
} - 5548
}; - 5549
- 5550
// Admit at the same busy boundary `/run` uses (docs/design/64, "Request - 5551
// durability and delivery"): a steer that lands on an IDLE session is a - 5552
// fresh admission and starts its own chain, rather than sitting in a - 5553
// queue nothing is left to drain (finding 1b). A steer is never a - 5554
// restricted (goal/managed/auto) request, so this never rejects. - 5555
match admit_or_queue(&handle, false, message.clone()) { - 5556
Admission::RejectedBusy => { - 5557
unreachable!("send_steering never admits a restricted request kind") - 5558
} - 5559
Admission::Queued => { - 5560
// The ledger was busy: the durable "Intervention queued" / - 5561
// "accepted" activity recorded above already covers this - 5562
// request. If it landed in the buffer rather than the ledger - 5563
// (session was busy at that check too), the busy leg's own - 5564
// `http_settle` flush removes this admission guard when it - 5565
// settles; nothing here needs to. - 5566
( - 5567
StatusCode::ACCEPTED, - 5568
Json(serde_json::json!({ - 5569
"request_id": evaluation.request.request_id, - 5570
"decision": evaluation.decision.as_str(), - 5571
"state": "steering_queued", - 5572
})), - 5573
) - 5574
.into_response() - 5575
} - 5576
Admission::Started(taken) => { - 5577
// The activity recorded above is durable now (whether it was - 5578
// written straight into the ledger or is about to be, via this - 5579
// very chain) — the in-memory admission guard can be released. - 5580
handle - 5581
.admissions - 5582
.lock() - 5583
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5584
.remove(&request_id); - 5585
let core = handle.core.clone(); - 5586
let preview_intent = core.resolve_turn_intent(&taken, &message); - 5587
let mut preview_outcome = - 5588
vak_intent::OutcomeSpec::from_intent(&body.text, &preview_intent); - 5589
preview_outcome.evidence_max_age_secs = Some(core.effective_evidence_max_age_secs()); - 5590
*handle - 5591
.intent - 5592
.lock() - 5593
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 5594
Some(vak_session::types::IntentRecord { - 5595
reading: preview_intent.reading, - 5596
strands: preview_intent.strands, - 5597
engagement: preview_intent.engagement, - 5598
provenance: preview_intent.provenance, - 5599
outcome: Some(preview_outcome), - 5600
model_visible: None, - 5601
commitment_id: None, - 5602
strand_commitments: Default::default(), - 5603
}); - 5604
wait_for_external_subscriber(&handle).await; - 5605
spawn_http_turn_chain( - 5606
&state, - 5607
handle, - 5608
core, - 5609
taken, - 5610
TurnStart::message(message), - 5611
None, - 5612
); - 5613
( - 5614
StatusCode::ACCEPTED, - 5615
Json(serde_json::json!({ - 5616
"request_id": evaluation.request.request_id, - 5617
"decision": evaluation.decision.as_str(), - 5618
"state": "started", - 5619
})), - 5620
) - 5621
.into_response() - 5622
} - 5623
} - 5624
} - 5625
- 5626
async fn cancel_run(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 5627
let Some(handle) = state.get(&id) else { - 5628
return StatusCode::NOT_FOUND; - 5629
}; - 5630
handle - 5631
.cancel - 5632
.lock() - 5633
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5634
.cancel(); - 5635
record_control_activity(&handle, "Run cancelled", "cancel"); - 5636
deny_pending_approvals(&handle); - 5637
// A stop means stop: whatever was queued for a continuation leg is - 5638
// discarded rather than silently running as the "next" turn once the - 5639
// cancelled run unwinds (finding 1c). Input that arrives AFTER this - 5640
// drain but before the run actually unwinds is a fresh push into the - 5641
// same queue and is unaffected — it becomes the next leg of the chain, - 5642
// same as any other steering. - 5643
let discarded = handle.steering.drain(vak_agent::DrainMode::All); - 5644
if !discarded.is_empty() { - 5645
record_activity_or_buffer( - 5646
&handle, - 5647
vak_session::ActivityRecord { - 5648
activity_id: format!("cancel-discard-{}", uuid::Uuid::now_v7()), - 5649
turn: None, - 5650
kind: vak_session::ActivityKind::Diagnostic, - 5651
status: vak_session::ActivityStatus::Cancelled, - 5652
label: "Queued input discarded by stop".into(), - 5653
detail: Some(format!( - 5654
"{} queued message{} discarded", - 5655
discarded.len(), - 5656
if discarded.len() == 1 { "" } else { "s" } - 5657
)), - 5658
data: std::collections::BTreeMap::new(), - 5659
}, - 5660
); - 5661
} - 5662
// No synthesized `RunFinished` here: the run's own settle path - 5663
// (`run_turn_chain`) emits the terminal event once it actually stops. - 5664
// Broadcasting one from here raced the real one — cancel_run's - 5665
// synthetic event could arrive, mark the chat finished on the client, - 5666
// and then the run's genuine `RunFinished` landed afterward and - 5667
// un-terminated it. - 5668
StatusCode::ACCEPTED - 5669
} - 5670
- 5671
async fn pause_run(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 5672
let Some(handle) = state.get(&id) else { - 5673
return StatusCode::NOT_FOUND; - 5674
}; - 5675
handle.steering.pause(); - 5676
record_control_activity(&handle, "Run paused", "pause"); - 5677
StatusCode::ACCEPTED - 5678
} - 5679
- 5680
async fn resume_run(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 5681
let Some(handle) = state.get(&id) else { - 5682
return StatusCode::NOT_FOUND; - 5683
}; - 5684
handle.steering.resume(); - 5685
record_control_activity(&handle, "Run resumed", "resume"); - 5686
StatusCode::ACCEPTED - 5687
} - 5688
- 5689
async fn control_state( - 5690
State(state): State<AppState>, - 5691
Path(id): Path<String>, - 5692
) -> axum::response::Response { - 5693
let Some(handle) = state.get(&id) else { - 5694
return StatusCode::NOT_FOUND.into_response(); - 5695
}; - 5696
let running = handle - 5697
.session - 5698
.lock() - 5699
.map(|session| session.is_none()) - 5700
.unwrap_or(false); - 5701
let revision = handle - 5702
.session - 5703
.lock() - 5704
.ok() - 5705
.and_then(|guard| { - 5706
guard.as_ref().and_then(|session| { - 5707
session.chain_to_root().iter().rev().find_map(|entry| { - 5708
if let vak_session::EntryPayload::Intent(record) = &entry.payload { - 5709
record.outcome.as_ref().map(|outcome| outcome.revision) - 5710
} else { - 5711
None - 5712
} - 5713
}) - 5714
}) - 5715
}) - 5716
.or_else(|| { - 5717
handle.intent.lock().ok().and_then(|guard| { - 5718
guard - 5719
.as_ref()? - 5720
.outcome - 5721
.as_ref() - 5722
.map(|outcome| outcome.revision) - 5723
}) - 5724
}) - 5725
.unwrap_or(0); - 5726
Json(serde_json::json!({ - 5727
"running": running, - 5728
"paused": handle.steering.is_paused(), - 5729
"revision": revision, - 5730
})) - 5731
.into_response() - 5732
} - 5733
- 5734
#[derive(serde::Deserialize)] - 5735
struct PlanChangeBody { - 5736
text: String, - 5737
#[serde(default = "default_intervention_source")] - 5738
source: String, - 5739
#[serde(default)] - 5740
target_revision: Option<u64>, - 5741
} - 5742
- 5743
fn apply_plan_change( - 5744
mut outcome: vak_intent::OutcomeSpec, - 5745
request: &vak_intent::InterventionRequest, - 5746
) -> (vak_intent::OutcomeSpec, (String, String)) { - 5747
let before = outcome - 5748
.requirements - 5749
.iter() - 5750
.map(|requirement| requirement.id.clone()) - 5751
.collect::<Vec<_>>(); - 5752
match request.kind { - 5753
vak_intent::InterventionKind::AddRequirement => { - 5754
outcome.requirements.push(vak_intent::OutcomeRequirement { - 5755
id: format!("intervention-{}", request.request_id), - 5756
kind: vak_intent::RequirementKind::Deliverable, - 5757
description: request.text.clone(), - 5758
origin: vak_intent::RequirementOrigin::Explicit, - 5759
importance: vak_intent::RequirementImportance::Must, - 5760
target: None, - 5761
}) - 5762
} - 5763
vak_intent::InterventionKind::RemoveRequirement => { - 5764
if let Some(target) = request.text.split_whitespace().last() { - 5765
outcome - 5766
.requirements - 5767
.retain(|requirement| requirement.id != target); - 5768
} - 5769
} - 5770
vak_intent::InterventionKind::Replan | vak_intent::InterventionKind::Reprioritize => { - 5771
outcome.assumptions.push(request.text.clone()); - 5772
} - 5773
_ => {} - 5774
} - 5775
outcome.revision = outcome.revision.saturating_add(1); - 5776
let after = outcome - 5777
.requirements - 5778
.iter() - 5779
.map(|requirement| requirement.id.clone()) - 5780
.collect::<Vec<_>>(); - 5781
(outcome, (before.join(","), after.join(","))) - 5782
} - 5783
- 5784
async fn plan_change( - 5785
State(state): State<AppState>, - 5786
Path(id): Path<String>, - 5787
Json(body): Json<PlanChangeBody>, - 5788
) -> axum::response::Response { - 5789
let Some(handle) = state.get(&id) else { - 5790
return StatusCode::NOT_FOUND.into_response(); - 5791
}; - 5792
let (kind, command_text) = intervention_of(&body.text); - 5793
if !matches!( - 5794
kind, - 5795
vak_intent::InterventionKind::Replan - 5796
| vak_intent::InterventionKind::Reprioritize - 5797
| vak_intent::InterventionKind::AddRequirement - 5798
| vak_intent::InterventionKind::RemoveRequirement - 5799
) { - 5800
return (StatusCode::BAD_REQUEST, "request is not a plan change").into_response(); - 5801
} - 5802
let revision = handle - 5803
.session - 5804
.lock() - 5805
.ok() - 5806
.and_then(|guard| { - 5807
guard.as_ref().and_then(|session| { - 5808
session.chain_to_root().iter().rev().find_map(|entry| { - 5809
if let vak_session::EntryPayload::Intent(record) = &entry.payload { - 5810
record.outcome.as_ref().map(|outcome| outcome.revision) - 5811
} else { - 5812
None - 5813
} - 5814
}) - 5815
}) - 5816
}) - 5817
.or_else(|| { - 5818
handle.intent.lock().ok().and_then(|guard| { - 5819
guard - 5820
.as_ref() - 5821
.and_then(|record| record.outcome.as_ref().map(|outcome| outcome.revision)) - 5822
}) - 5823
}) - 5824
.unwrap_or(0); - 5825
if body - 5826
.target_revision - 5827
.is_some_and(|target| target != revision) - 5828
{ - 5829
return ( - 5830
StatusCode::CONFLICT, - 5831
Json(serde_json::json!({ - 5832
"decision": "rejected", - 5833
"reason": "plan revision is stale", - 5834
"revision": revision, - 5835
})), - 5836
) - 5837
.into_response(); - 5838
} - 5839
let evaluation = vak_intent::evaluate_intervention(vak_intent::InterventionRequest { - 5840
request_id: format!("plan-change-{}", uuid::Uuid::now_v7()), - 5841
kind, - 5842
text: command_text, - 5843
source: control_source_for(&body.source, state.core.surface().slug(), &id), - 5844
target_revision: Some(revision), - 5845
target_session_id: Some(id.clone()), - 5846
target_parent_session_id: target_parent_of(&handle), - 5847
}); - 5848
let mut admitted_revision = None; - 5849
let mut requirement_diff = None; - 5850
if evaluation.decision == vak_intent::InterventionDecision::Queued - 5851
&& let Ok(mut guard) = handle.session.lock() - 5852
&& let Some(session) = guard.as_mut() - 5853
&& let Some(record) = session.chain_to_root().iter().rev().find_map(|entry| { - 5854
if let vak_session::EntryPayload::Intent(record) = &entry.payload { - 5855
Some(record.as_ref()) - 5856
} else { - 5857
None - 5858
} - 5859
}) - 5860
&& let Some(outcome) = record.outcome.clone() - 5861
{ - 5862
let (outcome, diff) = apply_plan_change(outcome, &evaluation.request); - 5863
requirement_diff = Some(diff); - 5864
let update = vak_session::types::IntentRecord { - 5865
reading: record.reading.clone(), - 5866
strands: record.strands.clone(), - 5867
engagement: record.engagement.clone(), - 5868
provenance: record.provenance.clone(), - 5869
outcome: Some(outcome), - 5870
model_visible: None, - 5871
commitment_id: record.commitment_id.clone(), - 5872
strand_commitments: record.strand_commitments.clone(), - 5873
}; - 5874
if session.append_intent(update.clone()).is_ok() { - 5875
*handle - 5876
.intent - 5877
.lock() - 5878
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(update); - 5879
admitted_revision = Some(revision + 1); - 5880
} - 5881
} - 5882
if evaluation.decision == vak_intent::InterventionDecision::Queued - 5883
&& admitted_revision.is_none() - 5884
&& let Some(mut record) = handle.intent.lock().ok().and_then(|guard| guard.clone()) - 5885
&& let Some(outcome) = record.outcome.take() - 5886
{ - 5887
let (outcome, diff) = apply_plan_change(outcome, &evaluation.request); - 5888
requirement_diff = Some(diff); - 5889
let update = vak_session::types::IntentRecord { - 5890
reading: record.reading, - 5891
strands: record.strands, - 5892
engagement: record.engagement, - 5893
provenance: record.provenance, - 5894
outcome: Some(outcome.clone()), - 5895
model_visible: None, - 5896
commitment_id: record.commitment_id, - 5897
strand_commitments: record.strand_commitments, - 5898
}; - 5899
*handle - 5900
.intent - 5901
.lock() - 5902
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(update.clone()); - 5903
handle.steering.push_outcome_update(update); - 5904
admitted_revision = Some(outcome.revision); - 5905
} - 5906
record_activity_or_buffer( - 5907
&handle, - 5908
vak_session::ActivityRecord { - 5909
activity_id: evaluation.request.request_id.clone(), - 5910
turn: None, - 5911
kind: vak_session::ActivityKind::Diagnostic, - 5912
status: if evaluation.decision == vak_intent::InterventionDecision::RequiresHuman { - 5913
vak_session::ActivityStatus::Denied - 5914
} else { - 5915
vak_session::ActivityStatus::Pending - 5916
}, - 5917
label: "Plan change evaluated".into(), - 5918
detail: Some(evaluation.reason.clone()), - 5919
data: std::collections::BTreeMap::from([ - 5920
("decision".into(), evaluation.decision.as_str().into()), - 5921
("source".into(), body.source), - 5922
( - 5923
"revision".into(), - 5924
admitted_revision.unwrap_or(revision).to_string(), - 5925
), - 5926
( - 5927
"requirements_before".into(), - 5928
requirement_diff - 5929
.as_ref() - 5930
.map(|diff| diff.0.clone()) - 5931
.unwrap_or_else(String::new), - 5932
), - 5933
( - 5934
"requirements_after".into(), - 5935
requirement_diff - 5936
.as_ref() - 5937
.map(|diff| diff.1.clone()) - 5938
.unwrap_or_else(String::new), - 5939
), - 5940
("request".into(), body.text), - 5941
]), - 5942
}, - 5943
); - 5944
let status = if evaluation.decision == vak_intent::InterventionDecision::RequiresHuman { - 5945
StatusCode::CONFLICT - 5946
} else { - 5947
handle - 5948
.steering - 5949
.push_steering(evaluation.request.text.clone()); - 5950
StatusCode::ACCEPTED - 5951
}; - 5952
( - 5953
status, - 5954
Json(serde_json::json!({ - 5955
"request_id": evaluation.request.request_id, - 5956
"decision": evaluation.decision.as_str(), - 5957
"revision": admitted_revision.unwrap_or(revision), - 5958
"reason": evaluation.reason, - 5959
})), - 5960
) - 5961
.into_response() - 5962
} - 5963
- 5964
pub(crate) fn record_control_activity(handle: &SessionHandle, label: &str, control: &str) { - 5965
record_activity_or_buffer( - 5966
handle, - 5967
vak_session::ActivityRecord { - 5968
activity_id: format!("control-{}", uuid::Uuid::now_v7()), - 5969
turn: None, - 5970
kind: vak_session::ActivityKind::Diagnostic, - 5971
status: vak_session::ActivityStatus::Succeeded, - 5972
label: label.into(), - 5973
detail: Some("operator control-plane request".into()), - 5974
data: std::collections::BTreeMap::from([ - 5975
("control".into(), control.into()), - 5976
("source".into(), "human".into()), - 5977
]), - 5978
}, - 5979
); - 5980
} - 5981
- 5982
fn record_activity_or_buffer(handle: &SessionHandle, activity: vak_session::ActivityRecord) { - 5983
if let Ok(mut session) = handle.session.lock() - 5984
&& let Some(session) = session.as_mut() - 5985
{ - 5986
let _ = session.append_activity(activity); - 5987
} else if let Ok(mut activities) = handle.activity_buffer.lock() { - 5988
activities.push(activity); - 5989
} - 5990
} - 5991
- 5992
fn record_routing_data( - 5993
data: &mut std::collections::BTreeMap<String, String>, - 5994
routing: &RoutingEnvelope, - 5995
) { - 5996
if let Some(value) = routing.message_id.as_deref() { - 5997
data.insert("message_id".into(), value.into()); - 5998
} - 5999
if let Some(value) = routing.conversation_id.as_deref() { - 6000
data.insert("conversation_id".into(), value.into()); - 6001
} - 6002
if let Some(value) = routing.target_work_id.as_deref() { - 6003
data.insert("target_work_id".into(), value.into()); - 6004
} - 6005
if let Some(value) = routing.target_result_id.as_deref() { - 6006
data.insert("target_result_id".into(), value.into()); - 6007
} - 6008
if let Some(value) = routing.relation.as_deref() { - 6009
data.insert("relation".into(), value.into()); - 6010
} - 6011
if let Some(value) = routing.outcome_revision { - 6012
data.insert("outcome_revision".into(), value.to_string()); - 6013
} - 6014
if let Some(value) = routing.provenance.as_deref() { - 6015
data.insert("routing_provenance".into(), value.into()); - 6016
} - 6017
} - 6018
- 6019
/// The session that dispatched this one, from its own header: what "an agent - 6020
/// may control only its own children" is checked against. Unknown while the - 6021
/// ledger is out on a running turn, and unknown refuses an agent's control. - 6022
fn target_parent_of(handle: &SessionHandle) -> Option<String> { - 6023
handle.session.lock().ok().and_then(|guard| { - 6024
guard.as_ref().and_then(|log| { - 6025
log.header() - 6026
.and_then(|header| header.parent_session_id.clone()) - 6027
}) - 6028
}) - 6029
} - 6030
- 6031
fn routing_revision_is_current(handle: &SessionHandle, expected: u64) -> bool { - 6032
handle.intent.lock().ok().is_some_and(|record| { - 6033
record - 6034
.as_ref() - 6035
.and_then(|intent| intent.outcome.as_ref()) - 6036
.is_some_and(|outcome| outcome.revision == expected) - 6037
}) - 6038
} - 6039
- 6040
fn routing_summary(routing: Option<&RoutingEnvelope>) -> String { - 6041
let Some(routing) = routing else { - 6042
return String::new(); - 6043
}; - 6044
let mut data = std::collections::BTreeMap::new(); - 6045
record_routing_data(&mut data, routing); - 6046
serde_json::to_string(&data).unwrap_or_default() - 6047
} - 6048
- 6049
fn deny_pending_approvals(handle: &SessionHandle) { - 6050
let requests: Vec<ApprovalRequest> = handle - 6051
.pending - 6052
.lock() - 6053
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6054
.drain() - 6055
.map(|(_, request)| request) - 6056
.collect(); - 6057
for request in requests { - 6058
request.respond(false); - 6059
} - 6060
} - 6061
- 6062
// ---- Worker control plane ------------------------------------------------- - 6063
// - 6064
// Children already stream lifecycle/tool events into the parent session's - 6065
// SSE channel; these endpoints add the missing half: listing, steering, and - 6066
// stopping from a remote surface. Scope-checked against the parent so one - 6067
// session can never touch another's child. - 6068
- 6069
async fn list_workers( - 6070
State(state): State<AppState>, - 6071
Path(id): Path<String>, - 6072
) -> Json<serde_json::Value> { - 6073
let children = state.core.workers().active_for(&id); - 6074
Json(serde_json::json!({ "workers": children })) - 6075
} - 6076
- 6077
#[derive(serde::Deserialize)] - 6078
struct WorkerSteerBody { - 6079
text: String, - 6080
} - 6081
- 6082
async fn steer_worker( - 6083
State(state): State<AppState>, - 6084
Path((id, child)): Path<(String, String)>, - 6085
Json(body): Json<WorkerSteerBody>, - 6086
) -> StatusCode { - 6087
if state.core.workers().parent_of(&child).as_deref() != Some(id.as_str()) { - 6088
return StatusCode::NOT_FOUND; - 6089
} - 6090
if state.core.workers().steer(&child, &body.text) { - 6091
StatusCode::ACCEPTED - 6092
} else { - 6093
StatusCode::CONFLICT - 6094
} - 6095
} - 6096
- 6097
async fn stop_worker( - 6098
State(state): State<AppState>, - 6099
Path((id, child)): Path<(String, String)>, - 6100
) -> StatusCode { - 6101
if state.core.workers().parent_of(&child).as_deref() != Some(id.as_str()) { - 6102
return StatusCode::NOT_FOUND; - 6103
} - 6104
if state.core.workers().stop(&child) { - 6105
StatusCode::ACCEPTED - 6106
} else { - 6107
StatusCode::CONFLICT - 6108
} - 6109
} - 6110
- 6111
#[derive(serde::Deserialize)] - 6112
struct ApprovalBody { - 6113
approve: bool, - 6114
/// "…and don't ask again for calls like this one". Derives the narrowest - 6115
/// rule that covers this call and persists it to - 6116
/// `.vak/permissions.local.toml`. Only meaningful alongside - 6117
/// `approve: true` — remembering a refusal would be a deny rule, which - 6118
/// is a different and much heavier decision than answering one gate. - 6119
#[serde(default)] - 6120
remember: bool, - 6121
} - 6122
- 6123
#[derive(serde::Deserialize)] - 6124
struct OutcomeReviewBody { - 6125
verdict: String, - 6126
#[serde(default)] - 6127
turn: Option<usize>, - 6128
#[serde(default)] - 6129
note: Option<String>, - 6130
} - 6131
- 6132
/// Record an operator's review of a computed outcome. This is deliberately - 6133
/// separate from approval: reviewing a result never authorizes a tool call. - 6134
async fn record_outcome_review( - 6135
State(state): State<AppState>, - 6136
Path(id): Path<String>, - 6137
Json(body): Json<OutcomeReviewBody>, - 6138
) -> axum::response::Response { - 6139
use axum::response::IntoResponse; - 6140
if !matches!( - 6141
body.verdict.as_str(), - 6142
"accepted" | "needs_work" | "rejected" - 6143
) { - 6144
return (StatusCode::BAD_REQUEST, "invalid outcome review verdict").into_response(); - 6145
} - 6146
let Some(handle) = state.get(&id) else { - 6147
return StatusCode::NOT_FOUND.into_response(); - 6148
}; - 6149
let Ok(mut guard) = handle.session.lock() else { - 6150
return StatusCode::CONFLICT.into_response(); - 6151
}; - 6152
let Some(session) = guard.as_mut() else { - 6153
return StatusCode::NOT_FOUND.into_response(); - 6154
}; - 6155
let reviewed_turn = body.turn.or_else(|| { - 6156
session - 6157
.chain_to_root() - 6158
.iter() - 6159
.filter_map(|entry| match &entry.payload { - 6160
vak_session::EntryPayload::Activity(activity) - 6161
if activity.label == "Outcome evaluation" => - 6162
{ - 6163
activity.turn
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.