- 4688
} - 4689
match events.try_send(AgentEvent::Stream(ev)) { - 4690
Ok(()) => {} - 4691
Err(mpsc::error::TrySendError::Full(sent)) => { - 4692
if let Some(event) = into_stream_event(sent) { - 4693
pending.push_back(event); - 4694
} - 4695
} - 4696
Err(mpsc::error::TrySendError::Closed(_)) => cancel.cancel(), - 4697
} - 4698
} - 4699
// The provider's stream ended: flush whatever backlog - 4700
// is left with an awaited send, so a slow listener - 4701
// never loses the tail of an otherwise-completed step - 4702
// to backpressure. - 4703
if forward { - 4704
for event in pending { - 4705
let _ = events.send(AgentEvent::Stream(event)).await; - 4706
} - 4707
} - 4708
stream - 4709
.result() - 4710
.await - 4711
.map(|message| (message, first_token_ms)) - 4712
}; - 4713
let step = std::panic::AssertUnwindSafe(step).catch_unwind(); - 4714
- 4715
let mut domain_override: Option<FailureDomain> = None; - 4716
let outcome = match self.config.request_timeout { - 4717
Some(t) => match tokio::time::timeout(t, step).await { - 4718
Ok(Ok(r)) => r, - 4719
Ok(Err(_)) => Err(LlmError::Network( - 4720
"provider dispatch panicked and was contained".into(), - 4721
)), - 4722
Err(_) => { - 4723
// The provider owns a spawned producer keyed by - 4724
// this token. A watchdog timeout must revoke - 4725
// the attempt before the retry/fallback path - 4726
// can release capacity and dispatch again -- - 4727
// but only THIS attempt: cancelling the run's - 4728
// own token here used to abort the whole run - 4729
// (the very next retry's backoff wait would - 4730
// see it already cancelled and return Aborted - 4731
// instead of actually retrying). - 4732
attempt_cancel.cancel(); - 4733
domain_override = Some(FailureDomain::Deadline); - 4734
Err(LlmError::Network(format!( - 4735
"model step exceeded deadline of {}s", - 4736
t.as_secs() - 4737
))) - 4738
} - 4739
}, - 4740
None => match step.await { - 4741
Ok(r) => r, - 4742
Err(_) => Err(LlmError::Network( - 4743
"provider dispatch panicked and was contained".into(), - 4744
)), - 4745
}, - 4746
}; - 4747
let elapsed_ms = started.elapsed().as_millis() as u64; - 4748
- 4749
match outcome { - 4750
Ok((r, first_token_ms)) => { - 4751
if let Some(breaker) = &self.config.circuit_breaker { - 4752
breaker.record_success_key(&breaker_key); - 4753
} - 4754
ledger.last_first_token_ms = first_token_ms; - 4755
ledger.receipt.record( - 4756
reason, - 4757
FailureDomain::Unknown, - 4758
Settlement::Ok, - 4759
elapsed_ms, - 4760
Some(r.usage.clone()), - 4761
None, - 4762
); - 4763
return Ok(r); - 4764
} - 4765
Err(e @ LlmError::Aborted { .. }) => { - 4766
ledger.receipt.record( - 4767
reason, - 4768
FailureDomain::Unknown, - 4769
Settlement::Cancelled, - 4770
elapsed_ms, - 4771
None, - 4772
None, - 4773
); - 4774
return Err(e); - 4775
} - 4776
Err(e) => { - 4777
let (domain, settlement) = vak_llm::work::classify_error(&e); - 4778
ledger.receipt.record( - 4779
reason, - 4780
domain_override.unwrap_or(domain), - 4781
settlement, - 4782
elapsed_ms, - 4783
None, - 4784
Some(e.to_string()), - 4785
); - 4786
if e.is_retryable() && attempt < self.config.max_retries { - 4787
if trips_breaker(&e) - 4788
&& let Some(breaker) = &self.config.circuit_breaker - 4789
{ - 4790
breaker.record_failure_key(&breaker_key); - 4791
} - 4792
attempt += 1; - 4793
let delay = backoff_delay( - 4794
attempt, - 4795
e.retry_after_secs(), - 4796
self.config.retry_base_backoff_ms, - 4797
); - 4798
self.record_activity( - 4799
vak_session::ActivityKind::Retry, - 4800
vak_session::ActivityStatus::Running, - 4801
format!("Retry attempt {attempt}"), - 4802
Some(e.to_string()), - 4803
[ - 4804
("attempt".into(), attempt.to_string()), - 4805
("delay_ms".into(), delay.as_millis().to_string()), - 4806
] - 4807
.into(), - 4808
) - 4809
.await; - 4810
let _ = events - 4811
.send(AgentEvent::RetryScheduled { - 4812
attempt, - 4813
delay_ms: delay.as_millis() as u64, - 4814
reason: e.to_string(), - 4815
}) - 4816
.await; - 4817
tokio::select! { - 4818
_ = cancel.cancelled() => { - 4819
return Err(LlmError::Aborted { partial: None }); - 4820
} - 4821
_ = tokio::time::sleep(delay) => {} - 4822
} - 4823
} else { - 4824
// Leg exhausted (retries burned or permanent error): - 4825
// defer to the next frozen candidate. - 4826
if trips_breaker(&e) - 4827
&& let Some(breaker) = &self.config.circuit_breaker - 4828
{ - 4829
breaker.record_failure_key(&breaker_key); - 4830
} - 4831
last_err = Some(e); - 4832
continue 'legs; - 4833
} - 4834
} - 4835
} - 4836
} - 4837
} - 4838
Err(last_err.unwrap_or_else(|| LlmError::Network("route ladder exhausted".into()))) - 4839
} - 4840
- 4841
/// Runs a batch, except that a card the user already sees is not shown - 4842
/// again. A card call is on screen after its first success, so an identical - 4843
/// one (a small model repeats it until a breaker fires) is answered with a - 4844
/// plain ack instead of being run — never an error, which the stop guard - 4845
/// would count as an unresolved failure and answer with another model turn. - 4846
/// Likewise a call identical to one that already delivered a file gets - 4847
/// that call's result without writing a second draft or asking for - 4848
/// approval again, and a card previewing a delivered file is not shown. - 4849
async fn execute_batch( - 4850
&self, - 4851
calls: Vec<PendingToolCall>, - 4852
cancel: &CancellationToken, - 4853
events: &mpsc::Sender<AgentEvent>, - 4854
) -> Vec<(String, ToolRunOutput)> { - 4855
// `recall` is answered from the session here, before dispatch — - 4856
// never sent to a worker (docs/design/68-context-engine.md §3/§7). - 4857
// Intercepting by name mirrors how `emit_*_card` results are - 4858
// rewritten after `execute_batch` below, just earlier: `recall` has - 4859
// no side effects to execute, only session state to read, and - 4860
// `ToolContext` carries no session handle for `RecallTool::execute` - 4861
// to use (AGENTS.md invariant 14). - 4862
let (recall_calls, calls): (Vec<_>, Vec<_>) = calls - 4863
.into_iter() - 4864
.partition(|call| vak_tools::canonical_tool_name(&call.name) == "recall"); - 4865
let mut results = Vec::with_capacity(recall_calls.len()); - 4866
for call in recall_calls { - 4867
let output = match self.resolve_recall(&call.input).await { - 4868
ToolRunOutput::Ok(content) => { - 4869
ToolRunOutput::Ok(windowed_result(&call.id, content, &self.call_yields)) - 4870
} - 4871
failed => failed, - 4872
}; - 4873
results.push((call.id, output)); - 4874
} - 4875
results.extend(self.execute_batch_inner(calls, cancel, events).await); - 4876
results - 4877
} - 4878
- 4879
/// Resolves one `recall` call's arguments against the current session: - 4880
/// `turn` → the resolved turn's full record as text; `presentation` → - 4881
/// the canonical payload; `id` → the evidence content, optionally - 4882
/// sliced by line range. The result is a current-turn tool result and - 4883
/// is verbatim for the rest of that turn like any other result. - 4884
async fn resolve_recall(&self, input: &Value) -> ToolRunOutput { - 4885
let request = match vak_tools::parse_recall_args(input) { - 4886
Ok(request) => request, - 4887
Err(message) => { - 4888
return ToolRunOutput::Err(format!( - 4889
r#"{{"type":"invalid_arguments","message":"{message}"}}"# - 4890
)); - 4891
} - 4892
}; - 4893
let session = self.session.lock().await; - 4894
match request { - 4895
RecallRequest::Turn(n) => { - 4896
let index = TurnIndex::from_log(&session); - 4897
match index.turn_by_number(n as usize) { - 4898
Some(turn) => ToolRunOutput::Ok(render_full_record(&turn.full_record())), - 4899
None => ToolRunOutput::Err(format!( - 4900
r#"{{"type":"invalid_arguments","message":"no turn numbered {n}"}}"# - 4901
)), - 4902
} - 4903
} - 4904
RecallRequest::Presentation(id) => { - 4905
match session - 4906
.presentations() - 4907
.into_iter() - 4908
.find(|(pid, _)| *pid == id) - 4909
{ - 4910
Some((_, record)) => ToolRunOutput::Ok(record.payload.to_string()), - 4911
None => ToolRunOutput::Err(format!( - 4912
r#"{{"type":"invalid_arguments","message":"no presentation {id}"}}"# - 4913
)), - 4914
} - 4915
} - 4916
RecallRequest::Id { id, range } => match session.evidence(&id) { - 4917
Some(evidence) => { - 4918
ToolRunOutput::Ok(vak_tools::apply_range(&evidence.content, range)) - 4919
} - 4920
None => ToolRunOutput::Err(format!( - 4921
r#"{{"type":"invalid_arguments","message":"no evidence {id}"}}"# - 4922
)), - 4923
}, - 4924
} - 4925
} - 4926
- 4927
async fn execute_batch_inner( - 4928
&self, - 4929
calls: Vec<PendingToolCall>, - 4930
cancel: &CancellationToken, - 4931
events: &mpsc::Sender<AgentEvent>, - 4932
) -> Vec<(String, ToolRunOutput)> { - 4933
let key = |call: &PendingToolCall| { - 4934
format!( - 4935
"{}\u{0}{}", - 4936
call.name, - 4937
serde_json::to_string(&call.input).unwrap_or_default() - 4938
) - 4939
}; - 4940
let mut shown = self - 4941
.presented_cards - 4942
.lock() - 4943
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4944
.clone(); - 4945
let mut repeats: Vec<String> = Vec::new(); - 4946
let mut answered: Vec<(String, ToolRunOutput)> = Vec::new(); - 4947
let mut live_cards: HashMap<String, String> = HashMap::new(); - 4948
let mut live_deliveries: HashMap<String, (String, String)> = HashMap::new(); - 4949
let mut live = Vec::with_capacity(calls.len()); - 4950
for call in calls { - 4951
let mut deliveries = self - 4952
.deliveries - 4953
.lock() - 4954
.unwrap_or_else(std::sync::PoisonError::into_inner); - 4955
if let Some(path) = self.delivered_file(&call.name, &call.input) { - 4956
let delivery_key = key(&call); - 4957
if let Some(first) = deliveries.results.get(&delivery_key) { - 4958
answered.push(( - 4959
call.id.clone(), - 4960
ToolRunOutput::Ok(format!("{DRAFT_REPEAT_ACK}\n{first}")), - 4961
)); - 4962
continue; - 4963
} - 4964
live_deliveries.insert(call.id.clone(), (delivery_key, path)); - 4965
} - 4966
if vak_tools::canonical_tool_name(&call.name) == "bash" - 4967
&& let Some(command) = call.input.get("command").and_then(Value::as_str) - 4968
&& let Some(delivered) = copies_delivered_draft(command, &deliveries.paths).cloned() - 4969
{ - 4970
answered.push(( - 4971
call.id.clone(), - 4972
ToolRunOutput::Ok(format!( - 4973
"{DRAFT_COPY_REFUSED} {delivered} is a draft waiting for the person's review, and copying it out of .vak/scratch/ would skip that review. It reaches the workspace when they accept it. Answer with one sentence saying what you changed." - 4974
)), - 4975
)); - 4976
continue; - 4977
} - 4978
if self.tool_presents_cards(&call.name) - 4979
&& let Some(path) = call - 4980
.input - 4981
.pointer("/payload/artifact_path") - 4982
.and_then(Value::as_str) - 4983
&& let Some(delivered) = deliveries - 4984
.paths - 4985
.iter() - 4986
.find(|delivered| same_workspace_path(delivered, path)) - 4987
.cloned() - 4988
{ - 4989
deliveries.withheld_cards.insert(call.id.clone()); - 4990
answered.push(( - 4991
call.id.clone(), - 4992
ToolRunOutput::Ok(format!( - 4993
"{WITHHELD_CARD_ACK} {delivered} is already in front of the person as a draft they review with its change list, and a card cannot show it better. Do not present it again; answer with one sentence saying what you changed." - 4994
)), - 4995
)); - 4996
continue; - 4997
} - 4998
drop(deliveries); - 4999
if self.tool_presents_cards(&call.name) { - 5000
let card_key = key(&call); - 5001
if !shown.insert(card_key.clone()) { - 5002
repeats.push(call.id.clone()); - 5003
continue; - 5004
} - 5005
live_cards.insert(call.id.clone(), card_key); - 5006
} - 5007
live.push(call); - 5008
} - 5009
let mut results = self.execute_batch_calls(live, cancel, events).await; - 5010
{ - 5011
let mut presented = self - 5012
.presented_cards - 5013
.lock() - 5014
.unwrap_or_else(std::sync::PoisonError::into_inner); - 5015
for (id, output) in &results { - 5016
if let (Some(card_key), ToolRunOutput::Ok(_)) = (live_cards.get(id), output) { - 5017
presented.insert(card_key.clone()); - 5018
} - 5019
} - 5020
} - 5021
{ - 5022
let mut deliveries = self - 5023
.deliveries - 5024
.lock() - 5025
.unwrap_or_else(std::sync::PoisonError::into_inner); - 5026
for (id, output) in &results { - 5027
if let (Some((delivery_key, path)), ToolRunOutput::Ok(text)) = - 5028
(live_deliveries.get(id), output) - 5029
{ - 5030
deliveries - 5031
.results - 5032
.insert(delivery_key.clone(), text.clone()); - 5033
deliveries.paths.insert(path.clone()); - 5034
} - 5035
} - 5036
} - 5037
results.extend(answered); - 5038
results.extend( - 5039
repeats - 5040
.into_iter() - 5041
.map(|id| (id, ToolRunOutput::Ok(CARD_REPEAT_ACK.into()))), - 5042
); - 5043
results - 5044
} - 5045
- 5046
/// The turn's answer when a current value was asked for and nothing - 5047
/// was retrieved after the one repair (docs/design/68 §7): an honest - 5048
/// statement naming the last figure this conversation recorded and - 5049
/// when, never that figure presented as current. - 5050
async fn stale_data_outcome(&self) -> TurnOutcome { - 5051
let last_known = { - 5052
let session = self.session.lock().await; - 5053
let turn_id = session.latest_directive_entry_id(); - 5054
let entries = session.chain_to_root(); - 5055
session - 5056
.presentations() - 5057
.into_iter() - 5058
.rev() - 5059
.find(|(_, record)| Some(record.turn_id.as_str()) != turn_id.as_deref()) - 5060
.map(|(id, record)| { - 5061
let when = entries - 5062
.iter() - 5063
.find(|entry| entry.id == id) - 5064
.map(|entry| entry.ts.format("%Y-%m-%d %H:%M UTC").to_string()) - 5065
.unwrap_or_else(|| "an earlier turn".to_string()); - 5066
format!( - 5067
" The most recent figure in this conversation was recorded at {when}: {}.", - 5068
record.identity_digest - 5069
) - 5070
}) - 5071
.unwrap_or_default() - 5072
}; - 5073
self.record_activity( - 5074
vak_session::ActivityKind::Diagnostic, - 5075
vak_session::ActivityStatus::Failed, - 5076
"stale-data-refused".into(), - 5077
Some( - 5078
"a current value was asked for, nothing was retrieved this turn after one repair, \ - 5079
and no carried-over figure was presented as current" - 5080
.into(), - 5081
), - 5082
std::collections::BTreeMap::new(), - 5083
) - 5084
.await; - 5085
TurnOutcome::Completed { - 5086
response: AssistantMessage { - 5087
content: vec![ContentBlock::text(format!( - 5088
"I could not retrieve a current value on this turn, so I am not presenting a \ - 5089
carried-over figure as current.{last_known} Ask again when a retrieval tool \ - 5090
is available, or ask for the last known figure explicitly." - 5091
))], - 5092
stop_reason: StopReason::EndTurn, - 5093
usage: Usage::default(), - 5094
model: self.config.model.clone(), - 5095
response_id: None, - 5096
}, - 5097
} - 5098
} - 5099
- 5100
/// The turn's answer when a card repeatedly refuses to be about what was - 5101
/// asked (docs/design/68 §7): honest text instead of a wrong card. When - 5102
/// this turn's own retrieval actually succeeded, its raw result is - 5103
/// quoted rather than left out — the evidence exists, only the card - 5104
/// built from it did not. - 5105
async fn topic_mismatch_outcome(&self, evidence: Option<String>) -> TurnOutcome { - 5106
self.record_activity( - 5107
vak_session::ActivityKind::Diagnostic, - 5108
vak_session::ActivityStatus::Failed, - 5109
"topic-mismatch-refused".into(), - 5110
Some( - 5111
"a card was offered twice whose content did not match the directive; the turn \ - 5112
closed on an honest statement instead of showing a wrong card" - 5113
.into(), - 5114
), - 5115
std::collections::BTreeMap::new(), - 5116
) - 5117
.await; - 5118
let text = match evidence { - 5119
Some(found) => format!( - 5120
"I could not turn what I found into a card that actually answers this, so here is \ - 5121
the raw result instead:\n\n{found}" - 5122
), - 5123
None => "I could not produce a card that matches what was asked, and had nothing else \ - 5124
to fall back on this turn. Could you rephrase the question?" - 5125
.to_string(), - 5126
}; - 5127
TurnOutcome::Completed { - 5128
response: AssistantMessage { - 5129
content: vec![ContentBlock::text(text)], - 5130
stop_reason: StopReason::EndTurn, - 5131
usage: Usage::default(), - 5132
model: self.config.model.clone(), - 5133
response_id: None, - 5134
}, - 5135
} - 5136
} - 5137
- 5138
/// The turn's answer when the model kept re-emitting a card it had - 5139
/// already shown: the card stands, the loop stops paying for acks. - 5140
async fn card_repeat_outcome(&self) -> TurnOutcome { - 5141
self.record_activity( - 5142
vak_session::ActivityKind::Diagnostic, - 5143
vak_session::ActivityStatus::Succeeded, - 5144
"card-repeat-exhausted".into(), - 5145
Some(format!( - 5146
"{CARD_REPEAT_EXHAUSTION_THRESHOLD} consecutive steps re-emitted an already-shown card; \ - 5147
the turn closed on that card as its answer" - 5148
)), - 5149
std::collections::BTreeMap::new(), - 5150
) - 5151
.await; - 5152
TurnOutcome::Completed { - 5153
response: AssistantMessage { - 5154
content: Vec::new(), - 5155
stop_reason: StopReason::EndTurn, - 5156
usage: Usage::default(), - 5157
model: self.config.model.clone(), - 5158
response_id: None, - 5159
}, - 5160
} - 5161
} - 5162
- 5163
async fn execute_batch_calls( - 5164
&self, - 5165
calls: Vec<PendingToolCall>, - 5166
cancel: &CancellationToken, - 5167
events: &mpsc::Sender<AgentEvent>, - 5168
) -> Vec<(String, ToolRunOutput)> { - 5169
let index = self - 5170
.config - 5171
.mcp_tool_index - 5172
.lock() - 5173
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5174
.clone(); - 5175
let calls = calls - 5176
.into_iter() - 5177
.map(normalize_tool_call) - 5178
.map(|call| normalize_mcp_call(call, &index)) - 5179
.map(|call| normalize_schema_wrapper(call, &self.config.tools)) - 5180
.collect::<Vec<_>>(); - 5181
let n = calls.len(); - 5182
let cwd = self - 5183
.session - 5184
.lock() - 5185
.await - 5186
.header() - 5187
.map(|h| h.contract_cwd()) - 5188
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| ".".into())); - 5189
let sandbox = self.config.sandbox.clone(); - 5190
let hooks = self.config.hooks.clone(); - 5191
let hook_recorder = self.config.hook_recorder.clone(); - 5192
let tool_activity_recorder = self.config.tool_activity_recorder.clone(); - 5193
let session_id = self - 5194
.session - 5195
.lock() - 5196
.await - 5197
.header() - 5198
.map(|h| h.session_id.clone()) - 5199
.unwrap_or_default(); - 5200
let agent_id = self - 5201
.session - 5202
.lock() - 5203
.await - 5204
.header() - 5205
.and_then(|h| h.agent.as_ref().map(|a| a.id.clone())); - 5206
let skill_names = self - 5207
.session - 5208
.lock() - 5209
.await - 5210
.header() - 5211
.map(|h| { - 5212
h.contract - 5213
.capabilities - 5214
.iter() - 5215
.filter(|capability| { - 5216
capability.kind == vak_session::types::CapabilityKind::Skill - 5217
}) - 5218
.map(|capability| capability.name.clone()) - 5219
.collect::<Vec<_>>() - 5220
}) - 5221
.unwrap_or_default(); - 5222
- 5223
let mut authz: Vec<Result<(), String>> = Vec::with_capacity(n); - 5224
for call in &calls { - 5225
authz.push( - 5226
authorize( - 5227
&self.config, - 5228
call, - 5229
&cwd, - 5230
&self.run_call_counts, - 5231
&self.config.tools, - 5232
) - 5233
.await, - 5234
); - 5235
} - 5236
let ids: Vec<String> = calls.iter().map(|c| c.id.clone()).collect(); - 5237
- 5238
if !self.config.parallel_tools - 5239
|| n == 1 - 5240
|| self.config.work_mode == WorkMode::Managed - 5241
|| calls.iter().any(|call| call.name == "work") - 5242
|| batch_has_file_dependency(&calls) - 5243
{ - 5244
let mut out = Vec::with_capacity(n); - 5245
for (call, verdict) in calls.into_iter().zip(authz) { - 5246
match verdict { - 5247
Err(reason) => out.push((call.id, ToolRunOutput::Err(reason))), - 5248
Ok(()) => { - 5249
if cancel.is_cancelled() { - 5250
out.push((call.id, ToolRunOutput::Err("cancelled".into()))); - 5251
continue; - 5252
} - 5253
let managed_item = (self.config.work_mode == WorkMode::Managed - 5254
&& call.name != "work") - 5255
.then(|| call.name.clone()); - 5256
let managed_item = match managed_item { - 5257
Some(tool) => { - 5258
self.begin_managed_tool_item( - 5259
&tool, - 5260
(call.name == "task" || call.name == "flow") - 5261
.then(|| { - 5262
Some(( - 5263
call.input.get("contract_id")?.as_str()?, - 5264
call.input.get("work_item_id")?.as_str()?, - 5265
)) - 5266
}) - 5267
.flatten(), - 5268
call.input.get("flow").and_then(|value| value.as_str()), - 5269
) - 5270
.await - 5271
} - 5272
None => None, - 5273
}; - 5274
if self.config.work_mode == WorkMode::Managed - 5275
&& call.name != "work" - 5276
&& managed_item.is_none() - 5277
{ - 5278
out.push(( - 5279
call.id, - 5280
ToolRunOutput::Err( - 5281
"managed execution requires a compatible Ready or Running work item" - 5282
.into(), - 5283
), - 5284
)); - 5285
continue; - 5286
} - 5287
out.push(if call.name == "work" { - 5288
let result = self.execute_work_call(&call.input).await; - 5289
self.emit_work_state(events).await; - 5290
(call.id.clone(), result) - 5291
} else { - 5292
let owns_lifecycle = call.name == "task" || call.name == "flow"; - 5293
let (returned_id, result) = execute_one( - 5294
call, - 5295
&self.config.tools, - 5296
&cwd, - 5297
&session_id, - 5298
agent_id.as_deref(), - 5299
&skill_names, - 5300
hooks.as_ref(), - 5301
hook_recorder.as_ref(), - 5302
tool_activity_recorder.as_ref(), - 5303
sandbox.as_ref(), - 5304
&self.call_yields, - 5305
cancel, - 5306
events, - 5307
) - 5308
.await; - 5309
if !owns_lifecycle - 5310
&& let Some((contract_id, item_id, attempt)) = managed_item - 5311
{ - 5312
self.finish_managed_tool_item( - 5313
&contract_id, - 5314
&item_id, - 5315
attempt, - 5316
&returned_id, - 5317
&result, - 5318
) - 5319
.await; - 5320
self.emit_work_state(events).await; - 5321
} - 5322
(returned_id, result) - 5323
}); - 5324
} - 5325
} - 5326
} - 5327
return out; - 5328
} - 5329
- 5330
// Greedy wave scheduling: claimed calls that conflict are placed in - 5331
// separate waves; unclaimed and read-only calls share wave 0. - 5332
let claims_for = |call: &PendingToolCall| -> vak_tools::ResourceClaims { - 5333
self.config - 5334
.tools - 5335
.iter() - 5336
.find(|t| t.name() == call.name) - 5337
.map(|t| t.claims(&call.input)) - 5338
.unwrap_or_default() - 5339
}; - 5340
let mut waves: Vec<Vec<usize>> = vec![Vec::new()]; - 5341
for (idx, call) in calls.iter().enumerate() { - 5342
let claims = claims_for(call); - 5343
if claims.is_unclaimed() { - 5344
waves[0].push(idx); - 5345
continue; - 5346
} - 5347
let mut placed = false; - 5348
for wave in waves.iter_mut() { - 5349
let compatible = wave.iter().all(|j| { - 5350
let other = claims_for(&calls[*j]); - 5351
!other.conflicts(&claims) && !claims.conflicts(&other) - 5352
}); - 5353
if compatible { - 5354
wave.push(idx); - 5355
placed = true; - 5356
break; - 5357
} - 5358
} - 5359
if !placed { - 5360
waves.push(vec![idx]); - 5361
} - 5362
} - 5363
- 5364
let mut ordered: Vec<Option<(String, ToolRunOutput)>> = (0..n).map(|_| None).collect(); - 5365
for wave in waves { - 5366
let mut join = tokio::task::JoinSet::new(); - 5367
for &idx in &wave { - 5368
if authz[idx].is_err() { - 5369
continue; - 5370
} - 5371
let call = match calls.get(idx) { - 5372
Some(c) => c.clone(), - 5373
None => continue, - 5374
}; - 5375
let tools = self.config.tools.clone(); - 5376
let cancel = cancel.clone(); - 5377
let events = events.clone(); - 5378
let cwd = cwd.clone(); - 5379
let sandbox = sandbox.clone(); - 5380
let session_id = session_id.clone(); - 5381
let skill_names = skill_names.clone(); - 5382
let hooks = hooks.clone(); - 5383
let hook_recorder = hook_recorder.clone(); - 5384
let tool_activity_recorder = tool_activity_recorder.clone(); - 5385
let agent_id = agent_id.clone(); - 5386
let yields = self.call_yields.clone(); - 5387
join.spawn(async move { - 5388
let r = execute_one( - 5389
call, - 5390
&tools, - 5391
&cwd, - 5392
&session_id, - 5393
agent_id.as_deref(), - 5394
&skill_names, - 5395
hooks.as_ref(), - 5396
hook_recorder.as_ref(), - 5397
tool_activity_recorder.as_ref(), - 5398
sandbox.as_ref(), - 5399
&yields, - 5400
&cancel, - 5401
&events, - 5402
) - 5403
.await; - 5404
(idx, r) - 5405
}); - 5406
} - 5407
while let Some(res) = join.join_next().await { - 5408
if let Ok((idx, pair)) = res { - 5409
ordered[idx] = Some(pair); - 5410
} - 5411
} - 5412
} - 5413
for (idx, verdict) in authz.into_iter().enumerate() { - 5414
if let Err(reason) = verdict { - 5415
ordered[idx] = Some((ids[idx].clone(), ToolRunOutput::Err(reason))); - 5416
} - 5417
} - 5418
ordered.into_iter().flatten().collect() - 5419
} - 5420
- 5421
async fn execute_work_call(&self, args: &serde_json::Value) -> ToolRunOutput { - 5422
let operation = args.get("operation").and_then(|value| value.as_str()); - 5423
let mut session = self.session.lock().await; - 5424
let projection = match session.work_projection() { - 5425
Ok(Some(projection)) => projection, - 5426
Ok(None) => return ToolRunOutput::Err("no active managed work contract".into()), - 5427
Err(error) => return ToolRunOutput::Err(format!("invalid work ledger: {error}")), - 5428
}; - 5429
match operation { - 5430
Some("get") => ToolRunOutput::Ok( - 5431
serde_json::to_string(&projection).unwrap_or_else(|_| "{}".into()), - 5432
), - 5433
Some("transition") => { - 5434
let Some(item_id) = args.get("item_id").and_then(|value| value.as_str()) else { - 5435
return ToolRunOutput::Err("work transition requires item_id".into()); - 5436
}; - 5437
let Some(to) = args.get("to").and_then(|value| value.as_str()) else { - 5438
return ToolRunOutput::Err("work transition requires to".into()); - 5439
}; - 5440
let Some(state) = projection.items.get(item_id) else { - 5441
return ToolRunOutput::Err(format!("unknown work item '{item_id}'")); - 5442
}; - 5443
let Some(target) = parse_model_item_status(to) else { - 5444
return ToolRunOutput::Err(format!("unsupported model transition '{to}'")); - 5445
}; - 5446
let event = vak_session::types::WorkEvent { - 5447
contract_id: projection.contract.contract_id.clone(), - 5448
revision: projection.contract.revision, - 5449
kind: vak_session::types::WorkEventKind::ItemStatusChanged { - 5450
item_id: item_id.into(), - 5451
from: state.status.clone(), - 5452
to: target, - 5453
attempt: state.attempt, - 5454
reason: args - 5455
.get("reason") - 5456
.and_then(|value| value.as_str()) - 5457
.unwrap_or_default() - 5458
.into(), - 5459
}, - 5460
}; - 5461
match session.append_work(event) { - 5462
Ok(_) => { - 5463
ToolRunOutput::Ok(format!("work item '{item_id}' transitioned to {to}")) - 5464
} - 5465
Err(error) => ToolRunOutput::Err(error.to_string()), - 5466
} - 5467
} - 5468
Some("attach_evidence") => { - 5469
ToolRunOutput::Err( - 5470
"model evidence attachment is disabled; successful tool results are attached automatically" - 5471
.into(), - 5472
) - 5473
} - 5474
_ => ToolRunOutput::Err( - 5475
"work operation must be get, transition, or attach_evidence".into(), - 5476
), - 5477
} - 5478
} - 5479
- 5480
async fn begin_managed_tool_item( - 5481
&self, - 5482
tool: &str, - 5483
requested: Option<(&str, &str)>, - 5484
flow_name: Option<&str>, - 5485
) -> Option<(String, String, u32)> { - 5486
let mut session = self.session.lock().await; - 5487
let projection = session.work_projection().ok().flatten()?; - 5488
let item = projection.items.values().find(|state| { - 5489
(state.status == vak_session::types::WorkItemStatus::Ready - 5490
|| state.status == vak_session::types::WorkItemStatus::Running) - 5491
&& projection - 5492
.contract - 5493
.items - 5494
.iter() - 5495
.find(|definition| definition.item_id == state.item_id) - 5496
.is_some_and(|definition| { - 5497
let dependencies_ready = definition.dependencies.iter().all(|dependency| { - 5498
matches!( - 5499
projection.items.get(dependency).map(|item| &item.status), - 5500
Some(vak_session::types::WorkItemStatus::Succeeded) - 5501
| Some(vak_session::types::WorkItemStatus::Skipped) - 5502
) - 5503
}); - 5504
dependencies_ready - 5505
&& match requested { - 5506
Some((contract_id, item_id)) => { - 5507
projection.contract.contract_id == contract_id - 5508
&& state.item_id == item_id - 5509
&& match (&definition.owner, tool, flow_name) { - 5510
(vak_session::types::WorkOwner::Worker, "task", _) => true, - 5511
( - 5512
vak_session::types::WorkOwner::Flow { name }, - 5513
"flow", - 5514
Some(requested_flow), - 5515
) => name == requested_flow, - 5516
_ => false, - 5517
} - 5518
} - 5519
None => { - 5520
matches!(definition.owner, vak_session::types::WorkOwner::ParentAgent) - 5521
|| matches!(&definition.owner, vak_session::types::WorkOwner::Tool { name } if name == tool) - 5522
} - 5523
} - 5524
}) - 5525
})?; - 5526
let contract_id = projection.contract.contract_id.clone(); - 5527
let item_id = item.item_id.clone(); - 5528
let attempt = item.attempt.saturating_add(1); - 5529
if item.status == vak_session::types::WorkItemStatus::Ready { - 5530
if requested.is_some() { - 5531
let owner = if tool == "flow" { - 5532
vak_session::types::WorkOwner::Flow { - 5533
name: flow_name.unwrap_or_default().into(), - 5534
} - 5535
} else { - 5536
vak_session::types::WorkOwner::Worker - 5537
}; - 5538
session - 5539
.append_work(vak_session::types::WorkEvent { - 5540
contract_id: contract_id.clone(), - 5541
revision: projection.contract.revision, - 5542
kind: vak_session::types::WorkEventKind::ItemAssigned { - 5543
item_id: item_id.clone(), - 5544
owner, - 5545
child_session_id: None, - 5546
}, - 5547
}) - 5548
.ok()?; - 5549
} - 5550
session - 5551
.append_work(vak_session::types::WorkEvent { - 5552
contract_id: contract_id.clone(), - 5553
revision: projection.contract.revision, - 5554
kind: vak_session::types::WorkEventKind::ItemStatusChanged { - 5555
item_id: item_id.clone(), - 5556
from: vak_session::types::WorkItemStatus::Ready, - 5557
to: vak_session::types::WorkItemStatus::Running, - 5558
attempt, - 5559
reason: format!("executing {tool}"), - 5560
}, - 5561
}) - 5562
.ok()?; - 5563
} - 5564
Some((contract_id, item_id, attempt)) - 5565
} - 5566
- 5567
async fn finish_managed_tool_item( - 5568
&self, - 5569
contract_id: &str, - 5570
item_id: &str, - 5571
attempt: u32, - 5572
tool_use_id: &str, - 5573
result: &ToolRunOutput, - 5574
) { - 5575
let mut session = self.session.lock().await; - 5576
let Ok(Some(projection)) = session.work_projection() else { - 5577
return; - 5578
}; - 5579
if projection.contract.contract_id != contract_id - 5580
|| projection - 5581
.items - 5582
.get(item_id) - 5583
.is_none_or(|state| state.status != vak_session::types::WorkItemStatus::Running) - 5584
{ - 5585
return; - 5586
} - 5587
let next = match result { - 5588
ToolRunOutput::Ok(_) => vak_session::types::WorkItemStatus::ReadyForVerification, - 5589
ToolRunOutput::Err(_) => vak_session::types::WorkItemStatus::Failed, - 5590
}; - 5591
let session_id = session - 5592
.header() - 5593
.map(|header| header.session_id.clone()) - 5594
.unwrap_or_default(); - 5595
if matches!(result, ToolRunOutput::Ok(_)) - 5596
&& session - 5597
.append_work(vak_session::types::WorkEvent { - 5598
contract_id: contract_id.into(), - 5599
revision: projection.contract.revision, - 5600
kind: vak_session::types::WorkEventKind::EvidenceAttached { - 5601
item_id: item_id.into(), - 5602
evidence: vak_session::types::EvidenceRef::ToolResult { - 5603
session_id, - 5604
tool_use_id: tool_use_id.into(), - 5605
}, - 5606
}, - 5607
}) - 5608
.is_err() - 5609
{ - 5610
return; - 5611
} - 5612
let _ = session.append_work(vak_session::types::WorkEvent { - 5613
contract_id: contract_id.into(), - 5614
revision: projection.contract.revision, - 5615
kind: vak_session::types::WorkEventKind::ItemStatusChanged { - 5616
item_id: item_id.into(), - 5617
from: vak_session::types::WorkItemStatus::Running, - 5618
to: next, - 5619
attempt, - 5620
reason: "parent tool execution returned".into(), - 5621
}, - 5622
}); - 5623
} - 5624
- 5625
async fn record_worker_work( - 5626
&self, - 5627
assignments: &[(String, String, String)], - 5628
results: &[(String, ToolRunOutput)], - 5629
) { - 5630
if self.config.work_mode != WorkMode::Managed { - 5631
return; - 5632
} - 5633
let mut session = self.session.lock().await; - 5634
for (call_id, contract_id, item_id) in assignments { - 5635
let Ok(Some(projection)) = session.work_projection() else { - 5636
continue; - 5637
}; - 5638
if projection.contract.contract_id != *contract_id { - 5639
continue; - 5640
} - 5641
let Some(state) = projection.items.get(item_id) else { - 5642
continue; - 5643
}; - 5644
if state.status != vak_session::types::WorkItemStatus::Running { - 5645
continue; - 5646
} - 5647
let Some((_, output)) = results.iter().find(|(id, _)| id == call_id) else { - 5648
continue; - 5649
}; - 5650
let child_id = match output { - 5651
ToolRunOutput::Ok(text) | ToolRunOutput::Err(text) => extract_worker_id(text), - 5652
}; - 5653
let revision = projection.contract.revision; - 5654
let assigned = vak_session::types::WorkEvent { - 5655
contract_id: contract_id.clone(), - 5656
revision, - 5657
kind: vak_session::types::WorkEventKind::ItemAssigned { - 5658
item_id: item_id.clone(), - 5659
owner: vak_session::types::WorkOwner::Worker, - 5660
child_session_id: child_id.clone(), - 5661
}, - 5662
}; - 5663
if session.append_work(assigned).is_err() { - 5664
continue; - 5665
} - 5666
let outcome = if matches!(output, ToolRunOutput::Ok(_)) { - 5667
vak_session::types::WorkItemStatus::ReadyForVerification - 5668
} else { - 5669
vak_session::types::WorkItemStatus::Failed - 5670
}; - 5671
if let Some(child_id) = child_id - 5672
&& matches!(output, ToolRunOutput::Ok(_)) - 5673
&& session - 5674
.append_work(vak_session::types::WorkEvent { - 5675
contract_id: contract_id.clone(), - 5676
revision, - 5677
kind: vak_session::types::WorkEventKind::EvidenceAttached { - 5678
item_id: item_id.clone(), - 5679
evidence: vak_session::types::EvidenceRef::ChildSession { - 5680
session_id: child_id, - 5681
}, - 5682
}, - 5683
}) - 5684
.is_err() - 5685
{ - 5686
continue; - 5687
}
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.