- 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 ( - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.