- 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") - 5001
{ - 5002
return ( - 5003
StatusCode::BAD_REQUEST, - 5004
Json(serde_json::json!({"error": format!("unknown work_mode '{mode}'")})), - 5005
) - 5006
.into_response(); - 5007
} - 5008
let managed = matches!(body.work_mode.as_deref(), Some("managed")); - 5009
let automatic = matches!(body.work_mode.as_deref(), Some("auto")); - 5010
if (managed || automatic) && !(body.attachments.is_empty() && body.files.is_empty()) { - 5011
return ( - 5012
StatusCode::BAD_REQUEST, - 5013
Json(serde_json::json!({"error": "managed work currently requires text-only input"})), - 5014
) - 5015
.into_response(); - 5016
} - 5017
let mut attached = Vec::new(); - 5018
for file in &body.files { - 5019
match inbox::attached(handle.core.cwd(), file) { - 5020
Ok(file) => attached.push(file), - 5021
Err(error) => { - 5022
return ( - 5023
StatusCode::BAD_REQUEST, - 5024
Json(serde_json::json!({ "error": error })), - 5025
) - 5026
.into_response(); - 5027
} - 5028
} - 5029
} - 5030
if let Err(e) = handle.core.provider() { - 5031
return provider_unavailable(e); - 5032
} - 5033
- 5034
let expanded_prompt = body.prompt.clone(); - 5035
let start = if let Some(objective) = body.goal.clone() { - 5036
TurnStart::Goal { - 5037
prompt: expanded_prompt.clone(), - 5038
objective, - 5039
criteria: body.criteria.clone(), - 5040
} - 5041
} else if managed { - 5042
TurnStart::Managed(expanded_prompt.clone()) - 5043
} else if automatic { - 5044
TurnStart::Auto(expanded_prompt.clone()) - 5045
} else if body.attachments.is_empty() && attached.is_empty() { - 5046
TurnStart::message(vak_llm::Message::user_text(expanded_prompt.clone())) - 5047
} else { - 5048
let mut blocks = vec![vak_llm::ContentBlock::text(expanded_prompt.clone())]; - 5049
for a in &body.attachments { - 5050
if a.data.trim().is_empty() { - 5051
continue; - 5052
} - 5053
blocks.push(vak_llm::ContentBlock::image_base64( - 5054
a.mime.clone(), - 5055
a.data.trim().to_string(), - 5056
)); - 5057
} - 5058
let mut attachments = Vec::with_capacity(attached.len()); - 5059
for (note, mut file) in attached { - 5060
file.block = blocks.len(); - 5061
blocks.push(vak_llm::ContentBlock::text(note)); - 5062
attachments.push(file); - 5063
} - 5064
TurnStart::Message(vak_session::MessageRecord { - 5065
message: vak_llm::Message { - 5066
role: vak_llm::Role::User, - 5067
content: blocks, - 5068
}, - 5069
meta: (!attachments.is_empty()).then(|| vak_session::MessageMeta { - 5070
attachments, - 5071
..Default::default() - 5072
}), - 5073
}) - 5074
}; - 5075
let restricted = matches!( - 5076
start, - 5077
TurnStart::Goal { .. } | TurnStart::Managed(_) | TurnStart::Auto(_) - 5078
); - 5079
let queue_message = start.preview_message(); - 5080
- 5081
// ---- Admit at the busy boundary (docs/design/64, "Request durability - 5082
// and delivery"): idle starts a chain now; busy queues durably or, for - 5083
// a goal/managed/auto request that cannot be queued, is rejected - 5084
// explicitly. Never a bare 202 that silently discards the input - 5085
// (finding 1). ---- - 5086
let mut taken = match admit_or_queue(&handle, restricted, queue_message.clone()) { - 5087
Admission::RejectedBusy => { - 5088
return ( - 5089
StatusCode::CONFLICT, - 5090
Json(serde_json::json!({ - 5091
"error": "a goal or managed/auto request cannot be queued while the session is busy; wait for the current run to finish" - 5092
})), - 5093
) - 5094
.into_response(); - 5095
} - 5096
Admission::Queued => { - 5097
let request_id = request_id.unwrap_or_else(|| format!("run-{}", uuid::Uuid::now_v7())); - 5098
let mut data = std::collections::BTreeMap::new(); - 5099
data.insert("request_id".into(), request_id.clone()); - 5100
if let Some(routing) = body.routing.as_ref() { - 5101
record_routing_data(&mut data, routing); - 5102
} - 5103
record_activity_or_buffer( - 5104
&handle, - 5105
vak_session::ActivityRecord { - 5106
activity_id: format!("admission-{request_id}"), - 5107
turn: None, - 5108
kind: vak_session::ActivityKind::Run, - 5109
status: vak_session::ActivityStatus::Pending, - 5110
label: "Request queued".into(), - 5111
detail: None, - 5112
data, - 5113
}, - 5114
); - 5115
return ( - 5116
StatusCode::ACCEPTED, - 5117
Json(serde_json::json!({"request_id": request_id, "state": "queued"})), - 5118
) - 5119
.into_response(); - 5120
} - 5121
Admission::Started(taken) => taken, - 5122
}; - 5123
- 5124
if taken.is_read_only() { - 5125
match vak_session::SessionLog::open(taken.path().to_path_buf()) { - 5126
Ok(writable) => { - 5127
taken = writable; - 5128
} - 5129
Err(vak_session::SessionError::Locked(_)) => { - 5130
*handle - 5131
.session - 5132
.lock() - 5133
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 5134
return ( - 5135
StatusCode::CONFLICT, - 5136
Json(serde_json::json!({ - 5137
"error": "This conversation is currently active in Vakyartha Desktop. Close or finish the task in Desktop before continuing here." - 5138
})), - 5139
) - 5140
.into_response(); - 5141
} - 5142
Err(e) => { - 5143
*handle - 5144
.session - 5145
.lock() - 5146
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 5147
return ( - 5148
StatusCode::CONFLICT, - 5149
Json(serde_json::json!({ "error": e.to_string() })), - 5150
) - 5151
.into_response(); - 5152
} - 5153
} - 5154
} - 5155
if let Some(request_id) = request_id.as_deref() - 5156
&& taken.has_request_admission(request_id) - 5157
{ - 5158
*handle - 5159
.session - 5160
.lock() - 5161
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 5162
return ( - 5163
StatusCode::ACCEPTED, - 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 (
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.