- 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
} - 5688
if let Ok(Some(after_running)) = session.work_projection() { - 5689
let _ = session.append_work(vak_session::types::WorkEvent { - 5690
contract_id: contract_id.clone(), - 5691
revision, - 5692
kind: vak_session::types::WorkEventKind::ItemStatusChanged { - 5693
item_id: item_id.clone(), - 5694
from: vak_session::types::WorkItemStatus::Running, - 5695
to: outcome, - 5696
attempt: after_running.items[item_id].attempt, - 5697
reason: "worker returned".into(), - 5698
}, - 5699
}); - 5700
} - 5701
} - 5702
} - 5703
} - 5704
- 5705
/// Decides `verification_stale` from the model's actual call-issue order, - 5706
/// not from whatever order the tool results happened to come back in. - 5707
/// `code_mutation_ids` codepaths turn the flag on (an edit to a code path - 5708
/// just landed and hasn't been re-verified); `bash_ids` turn it off (a - 5709
/// bash run just re-verified, or is at least the most recent evidence). - 5710
/// Only the LAST succeeded call, in issue order, that matches either list - 5711
/// decides the outcome — everything else in the batch is irrelevant to it. - 5712
/// A call that never succeeded doesn't count as either kind of evidence. - 5713
fn resolve_verification_stale( - 5714
call_issue_order: &[String], - 5715
bash_ids: &[&str], - 5716
code_mutation_ids: &[&str], - 5717
succeeded: &std::collections::HashSet<&str>, - 5718
current: bool, - 5719
) -> bool { - 5720
let mut stale = current; - 5721
for id in call_issue_order { - 5722
if !succeeded.contains(id.as_str()) { - 5723
continue; - 5724
} - 5725
if bash_ids.contains(&id.as_str()) { - 5726
stale = false; - 5727
} else if code_mutation_ids.contains(&id.as_str()) { - 5728
stale = true; - 5729
} - 5730
} - 5731
stale - 5732
} - 5733
- 5734
#[cfg(test)] - 5735
mod verification_stale_tests { - 5736
use super::resolve_verification_stale; - 5737
use std::collections::HashSet; - 5738
- 5739
#[test] - 5740
fn edit_issued_after_bash_stays_stale_regardless_of_result_order() { - 5741
// Model issues bash first, then edits a code file — the bash - 5742
// verification is now stale, no matter which result comes back - 5743
// first from a concurrent batch. - 5744
let order = vec!["b1".to_string(), "e1".to_string()]; - 5745
let succeeded: HashSet<&str> = ["b1", "e1"].into_iter().collect(); - 5746
assert!(resolve_verification_stale( - 5747
&order, - 5748
&["b1"], - 5749
&["e1"], - 5750
&succeeded, - 5751
false - 5752
)); - 5753
} - 5754
- 5755
#[test] - 5756
fn bash_issued_after_edit_clears_stale_regardless_of_result_order() { - 5757
// Model edits a code file, then runs bash to verify it — no - 5758
// longer stale, no matter which result comes back first. - 5759
let order = vec!["e1".to_string(), "b1".to_string()]; - 5760
let succeeded: HashSet<&str> = ["b1", "e1"].into_iter().collect(); - 5761
assert!(!resolve_verification_stale( - 5762
&order, - 5763
&["b1"], - 5764
&["e1"], - 5765
&succeeded, - 5766
false - 5767
)); - 5768
} - 5769
- 5770
#[test] - 5771
fn a_failed_call_is_not_evidence_either_way() { - 5772
// Bash issued after the edit, but the bash call FAILED — the edit - 5773
// is still unverified, so staleness must not clear. - 5774
let order = vec!["e1".to_string(), "b1".to_string()]; - 5775
let succeeded: HashSet<&str> = ["e1"].into_iter().collect(); // b1 not in succeeded - 5776
assert!(resolve_verification_stale( - 5777
&order, - 5778
&["b1"], - 5779
&["e1"], - 5780
&succeeded, - 5781
false - 5782
)); - 5783
} - 5784
- 5785
#[test] - 5786
fn irrelevant_calls_in_the_batch_do_not_affect_the_flag() { - 5787
let order = vec!["e1".to_string(), "r1".to_string()]; - 5788
let succeeded: HashSet<&str> = ["e1", "r1"].into_iter().collect(); - 5789
assert!(resolve_verification_stale( - 5790
&order, - 5791
&["b1"], - 5792
&["e1"], - 5793
&succeeded, - 5794
false - 5795
)); - 5796
} - 5797
} - 5798
- 5799
/// A model may issue `write(path)` and `read(path)` in one tool batch. Those - 5800
/// calls have an order dependency even when parallel tools are enabled: the - 5801
/// read must observe the saved bytes, not race the write in another worker. - 5802
fn batch_has_file_dependency(calls: &[PendingToolCall]) -> bool { - 5803
calls.iter().enumerate().any(|(i, first)| { - 5804
let Some(path) = first.input.get("path").and_then(Value::as_str) else { - 5805
return false; - 5806
}; - 5807
matches!(first.name.as_str(), "write" | "edit" | "read") - 5808
&& calls.iter().skip(i + 1).any(|second| { - 5809
second.input.get("path").and_then(Value::as_str) == Some(path) - 5810
&& matches!(second.name.as_str(), "write" | "edit" | "read") - 5811
&& (first.name != "read" || second.name != "read") - 5812
}) - 5813
}) - 5814
} - 5815
- 5816
#[cfg(test)] - 5817
mod file_batch_dependency_tests { - 5818
use super::{PendingToolCall, batch_has_file_dependency}; - 5819
- 5820
fn call(name: &str, path: &str) -> PendingToolCall { - 5821
PendingToolCall { - 5822
id: format!("{name}-{path}"), - 5823
name: name.into(), - 5824
input: serde_json::json!({"path": path}), - 5825
} - 5826
} - 5827
- 5828
#[test] - 5829
fn same_file_write_and_read_run_in_model_order() { - 5830
assert!(batch_has_file_dependency(&[ - 5831
call("write", "report.csv"), - 5832
call("read", "report.csv") - 5833
])); - 5834
assert!(batch_has_file_dependency(&[ - 5835
call("read", "report.csv"), - 5836
call("edit", "report.csv") - 5837
])); - 5838
assert!(!batch_has_file_dependency(&[ - 5839
call("write", "a.csv"), - 5840
call("read", "b.csv") - 5841
])); - 5842
assert!(!batch_has_file_dependency(&[ - 5843
call("read", "a.csv"), - 5844
call("read", "a.csv") - 5845
])); - 5846
} - 5847
} - 5848
- 5849
fn normalize_tool_call(mut call: PendingToolCall) -> PendingToolCall { - 5850
let canonical = vak_tools::canonical_tool_name(&call.name); - 5851
if canonical != call.name { - 5852
call.name = canonical.to_string(); - 5853
} - 5854
if (call.name == "read" || call.name == "write" || call.name == "edit") - 5855
&& call.input.is_object() - 5856
&& let Some(obj) = call.input.as_object_mut() - 5857
&& !obj.contains_key("path") - 5858
&& let Some(file_path) = obj.get("file_path").or_else(|| obj.get("file")).cloned() - 5859
{ - 5860
obj.insert("path".into(), file_path); - 5861
} - 5862
if call.name == "bash" - 5863
&& call.input.is_object() - 5864
&& let Some(obj) = call.input.as_object_mut() - 5865
&& !obj.contains_key("command") - 5866
&& let Some(cmd) = obj - 5867
.get("cmd") - 5868
.or_else(|| obj.get("script")) - 5869
.or_else(|| obj.get("code")) - 5870
.or_else(|| obj.get("input")) - 5871
.cloned() - 5872
{ - 5873
obj.insert("command".into(), cmd); - 5874
} - 5875
call - 5876
} - 5877
- 5878
/// Name only retrieval routes actually offered on this step. A generic - 5879
/// "use a retrieval tool" repair left small models claiming they had no live - 5880
/// access even when an MCP broker was admitted and had worked in this same - 5881
/// session. Server and tool instances remain discovered at demand time. - 5882
fn freshness_retrieval_hint(definitions: &[vak_llm::ToolDefinition]) -> String { - 5883
let mut routes = Vec::new(); - 5884
for definition in definitions { - 5885
match definition.name.as_str() { - 5886
"mcp" => routes - 5887
.push("mcp (list a configured server's tools, then call an exact discovered tool)"), - 5888
"browse" => routes.push("browse"), - 5889
"webfetch" => routes.push("webfetch"), - 5890
"find_tools" => routes.push("find_tools (discover another admitted retrieval tool)"), - 5891
_ => {} - 5892
} - 5893
} - 5894
if routes.is_empty() { - 5895
"Use an admitted retrieval tool if one is available.".into() - 5896
} else { - 5897
format!( - 5898
"Admitted retrieval routes on this step: {}.", - 5899
routes.join(", ") - 5900
) - 5901
} - 5902
} - 5903
- 5904
fn is_search_results_fetch(input: &Value) -> bool { - 5905
let Some(url) = input.get("url").and_then(Value::as_str) else { - 5906
return false; - 5907
}; - 5908
let Some((_, authority_and_path)) = url.split_once("://") else { - 5909
return false; - 5910
}; - 5911
let Some((_, path_and_query)) = authority_and_path.split_once('/') else { - 5912
return false; - 5913
}; - 5914
let Some((path, query)) = path_and_query.split_once('?') else { - 5915
return false; - 5916
}; - 5917
let path = path.trim_end_matches('/'); - 5918
(path == "search" || path.ends_with("/search")) - 5919
&& query.split('&').any(|part| { - 5920
part.starts_with("q=") || part.starts_with("query=") || part.starts_with("search=") - 5921
}) - 5922
} - 5923
- 5924
/// Whether an answer honestly says it has no data. The one list the - 5925
/// freshness and grounding checks share (a judgement, listed in - 5926
/// docs/design/30-output-engineering.md). - 5927
fn admits_no_data(text: &str) -> bool { - 5928
let lower = text.to_lowercase(); - 5929
[ - 5930
"don't have", - 5931
"do not have", - 5932
"no access to", - 5933
"couldn't find", - 5934
"could not find", - 5935
"unable to find", - 5936
"no live data", - 5937
] - 5938
.iter() - 5939
.any(|phrase| lower.contains(phrase)) - 5940
} - 5941
- 5942
/// Whether `answer` visibly draws on `evidence`: it repeats a host or a - 5943
/// figure the retrieval returned. Evidence with nothing so checkable says - 5944
/// nothing either way, so it counts as used — the check exists to catch an - 5945
/// answer written from memory beside a retrieval, not to demand a format. - 5946
fn uses_evidence(answer: &str, evidence: &str) -> bool { - 5947
let answer = answer.to_lowercase(); - 5948
let mut markers = evidence - 5949
.split(|c: char| c.is_whitespace() || "()[]{}<>\"'`,;|".contains(c)) - 5950
.map(|token| { - 5951
token - 5952
.trim_matches(|c: char| ".:!?".contains(c)) - 5953
.to_lowercase() - 5954
}) - 5955
.filter_map(|token| { - 5956
let host = token - 5957
.split_once("://") - 5958
.map(|(_, rest)| rest) - 5959
.unwrap_or(&token) - 5960
.split('/') - 5961
.next() - 5962
.unwrap_or_default() - 5963
.trim_start_matches("www.") - 5964
.to_string(); - 5965
let is_host = host.contains('.') - 5966
&& host.rsplit('.').next().is_some_and(|tld| { - 5967
tld.len() >= 2 && tld.chars().all(|c| c.is_ascii_alphabetic()) - 5968
}); - 5969
let is_figure = token.chars().filter(char::is_ascii_digit).count() >= 2; - 5970
if is_host { - 5971
Some(host) - 5972
} else if is_figure { - 5973
Some(token) - 5974
} else { - 5975
None - 5976
} - 5977
}) - 5978
.peekable(); - 5979
if markers.peek().is_none() { - 5980
return true; - 5981
} - 5982
markers.any(|marker| answer.contains(&marker)) - 5983
} - 5984
- 5985
fn normalize_mcp_call( - 5986
mut call: PendingToolCall, - 5987
index: &std::collections::HashMap<String, String>, - 5988
) -> PendingToolCall { - 5989
if let Some(server) = index.get(&call.name) { - 5990
let tool = std::mem::replace(&mut call.name, "mcp".into()); - 5991
call.input = serde_json::json!({ - 5992
"action": "call", - 5993
"server": server, - 5994
"tool": tool, - 5995
"arguments": call.input, - 5996
}); - 5997
} else if call.name == "mcp" { - 5998
// Some providers ignore the `oneOf` discriminator and omit `action`. - 5999
// Recover the unambiguous shapes at the broker boundary so a malformed - 6000
// call does not strand an otherwise valid turn: server+tool means call; - 6001
// anything else is the harmless catalog request. - 6002
if let Some(obj) = call.input.as_object_mut() - 6003
&& !obj.contains_key("action") - 6004
{ - 6005
let has_server = obj - 6006
.get("server") - 6007
.and_then(|v| v.as_str()) - 6008
.is_some_and(|s| !s.is_empty()); - 6009
let has_tool = obj - 6010
.get("tool") - 6011
.and_then(|v| v.as_str()) - 6012
.is_some_and(|s| !s.is_empty()); - 6013
obj.insert( - 6014
"action".into(), - 6015
serde_json::Value::String( - 6016
if has_server && has_tool { - 6017
"call" - 6018
} else { - 6019
"list" - 6020
} - 6021
.into(), - 6022
), - 6023
); - 6024
} - 6025
// Dynamic broker auto-resolution: if the model called `mcp` with `action: "call"` - 6026
// and specified `tool`, but omitted or left `server` empty, resolve `server` - 6027
// dynamically if the tool name uniquely maps to an admitted server in `index`. - 6028
if let Some(obj) = call.input.as_object_mut() { - 6029
let is_call = obj.get("action").and_then(|a| a.as_str()) == Some("call"); - 6030
let server_missing = obj - 6031
.get("server") - 6032
.is_none_or(|s| s.is_null() || s.as_str() == Some("")); - 6033
if is_call - 6034
&& server_missing - 6035
&& let Some(tool_name) = obj.get("tool").and_then(|t| t.as_str()) - 6036
&& let Some(server) = index.get(tool_name) - 6037
{ - 6038
obj.insert("server".into(), serde_json::Value::String(server.clone())); - 6039
} - 6040
if is_call - 6041
&& obj - 6042
.get("server") - 6043
.and_then(|s| s.as_str()) - 6044
.is_none_or(str::is_empty) - 6045
{ - 6046
obj.insert("action".into(), serde_json::Value::String("list".into())); - 6047
obj.remove("server"); - 6048
obj.remove("tool"); - 6049
obj.remove("arguments"); - 6050
} - 6051
} - 6052
} - 6053
call - 6054
} - 6055
- 6056
/// Some provider dialects wrap a tool's arguments once beneath a generated - 6057
/// label such as `metric_card` or `arguments`. Recover only the unambiguous - 6058
/// case: the outer call fails the declared schema, contains exactly one - 6059
/// unknown object value, and that inner object satisfies the schema in full. - 6060
/// - 6061
/// This is deliberately schema-driven rather than a list of card/tool names. - 6062
/// It cannot make an invalid payload valid, discard sibling fields, or weaken - 6063
/// the authorization boundary: the canonical value is still validated again - 6064
/// by `authorize` before dispatch. - 6065
fn normalize_schema_wrapper(mut call: PendingToolCall, tools: &[Arc<dyn Tool>]) -> PendingToolCall { - 6066
let Some(tool) = tools.iter().find(|tool| tool.name() == call.name) else { - 6067
return call; - 6068
}; - 6069
let schema = tool.schema(); - 6070
if let Some(input) = unwrapped_schema_input(&schema, &call.input) { - 6071
call.input = input; - 6072
} - 6073
call - 6074
} - 6075
- 6076
fn unwrapped_schema_input(schema: &Value, input: &Value) -> Option<Value> { - 6077
if vak_tools::validate_input(schema, input).is_ok() { - 6078
return None; - 6079
} - 6080
let object = input.as_object()?; - 6081
let (outer_key, inner) = object.iter().next().filter(|_| object.len() == 1)?; - 6082
if !inner.is_object() - 6083
|| schema - 6084
.get("properties") - 6085
.and_then(Value::as_object) - 6086
.is_some_and(|properties| properties.contains_key(outer_key)) - 6087
|| vak_tools::validate_input(schema, inner).is_err() - 6088
{ - 6089
return None; - 6090
} - 6091
Some(inner.clone()) - 6092
} - 6093
- 6094
fn extract_worker_id(text: &str) -> Option<String> { - 6095
let prefix = "worker '"; - 6096
let start = text.find(prefix)? + prefix.len(); - 6097
let rest = &text[start..]; - 6098
Some(rest.split('\'').next()?.to_string()) - 6099
} - 6100
- 6101
#[derive(serde::Deserialize)] - 6102
struct AuthoredContract { - 6103
objective: String, - 6104
#[serde(default)] - 6105
constraints: Vec<vak_session::types::WorkConstraint>, - 6106
#[serde(default)] - 6107
assumptions: Vec<vak_session::types::WorkAssumption>, - 6108
#[serde(default)] - 6109
criteria: Vec<vak_session::types::WorkCriterion>, - 6110
items: Vec<vak_session::types::WorkItemDefinition>, - 6111
} - 6112
- 6113
pub fn validate_work_paths(contract: &vak_session::types::WorkContract) -> Result<(), String> { - 6114
for item in &contract.items { - 6115
for claim in &item.path_claims { - 6116
let path = std::path::Path::new(claim); - 6117
if path.is_absolute() - 6118
|| path - 6119
.components() - 6120
.any(|component| component == std::path::Component::ParentDir) - 6121
{ - 6122
return Err(format!( - 6123
"work item '{}' claims a path outside the workspace", - 6124
item.item_id - 6125
)); - 6126
} - 6127
} - 6128
} - 6129
for criterion in &contract.criteria { - 6130
let path = match &criterion.kind { - 6131
vak_session::types::CriterionKind::FileExists { path } - 6132
| vak_session::types::CriterionKind::FileContains { path, .. } => Some(path), - 6133
_ => None, - 6134
}; - 6135
if let Some(path) = path - 6136
&& (path.is_absolute() - 6137
|| path - 6138
.components() - 6139
.any(|component| component == std::path::Component::ParentDir)) - 6140
{ - 6141
return Err(format!( - 6142
"criterion '{}' claims a path outside the workspace", - 6143
criterion.criterion_id - 6144
)); - 6145
} - 6146
} - 6147
Ok(()) - 6148
} - 6149
- 6150
fn parse_model_item_status(value: &str) -> Option<vak_session::types::WorkItemStatus> { - 6151
match value { - 6152
"running" => Some(vak_session::types::WorkItemStatus::Running), - 6153
"blocked" => Some(vak_session::types::WorkItemStatus::Blocked), - 6154
"ready_for_verification" => Some(vak_session::types::WorkItemStatus::ReadyForVerification), - 6155
_ => None, - 6156
} - 6157
} - 6158
- 6159
fn workspace_criterion_path( - 6160
cwd: &std::path::Path, - 6161
path: &std::path::Path, - 6162
) -> Option<std::path::PathBuf> { - 6163
let workspace = std::fs::canonicalize(cwd).ok()?; - 6164
let candidate = if path.is_absolute() { - 6165
path.to_path_buf() - 6166
} else { - 6167
cwd.join(path) - 6168
}; - 6169
let resolved = std::fs::canonicalize(candidate).ok()?; - 6170
resolved.starts_with(workspace).then_some(resolved) - 6171
} - 6172
#[allow(clippy::too_many_arguments)] - 6173
async fn execute_one( - 6174
call: PendingToolCall, - 6175
tools: &[Arc<dyn Tool>], - 6176
cwd: &std::path::Path, - 6177
session_id: &str, - 6178
agent_id: Option<&str>, - 6179
skill_names: &[String], - 6180
hooks: Option<&Arc<Vec<vak_hooks::HookDef>>>, - 6181
hook_recorder: Option<&HookRecorder>, - 6182
tool_activity_recorder: Option<&ToolActivityRecorder>, - 6183
sandbox: Option<&Arc<dyn vak_tools::sandbox::Sandbox>>, - 6184
yields: &CallYields, - 6185
cancel: &CancellationToken, - 6186
events: &mpsc::Sender<AgentEvent>, - 6187
) -> (String, ToolRunOutput) { - 6188
let started = std::time::Instant::now(); - 6189
let activity_input = call.input.clone(); - 6190
let _ = events - 6191
.send(AgentEvent::ToolCallStart { - 6192
id: call.id.clone(), - 6193
name: call.name.clone(), - 6194
args_json: args_preview(&call.input), - 6195
}) - 6196
.await; - 6197
- 6198
if let Some(hooks) = hooks { - 6199
let pre = vak_hooks::run_hooks_with_recorder( - 6200
hooks.clone(), - 6201
vak_hooks::HookEvent::PreToolUse, - 6202
session_id, - 6203
cwd, - 6204
Some((&call.name, &call.input)), - 6205
None, - 6206
cancel, - 6207
hook_recorder.map(|recorder| &**recorder), - 6208
) - 6209
.await; - 6210
if pre.blocked { - 6211
let reason = pre.reason.unwrap_or_else(|| "blocked by hook".into()); - 6212
let content = format!("blocked by hook: {reason}"); - 6213
let _ = events - 6214
.send(AgentEvent::ToolCallEnd { - 6215
id: call.id.clone(), - 6216
name: call.name.clone(), - 6217
is_error: true, - 6218
result_preview: result_preview(&content), - 6219
}) - 6220
.await; - 6221
return (call.id, ToolRunOutput::Err(content)); - 6222
} - 6223
} - 6224
- 6225
let tool = tools.iter().find(|t| t.name() == call.name); - 6226
let hook_name = call.name.clone(); - 6227
let hook_input = call.input.clone(); - 6228
- 6229
let matched_skill = skill_names.iter().find(|name| { - 6230
name.as_str() == call.name - 6231
|| name.replace('-', "_") == call.name - 6232
|| name.as_str() == call.name.replace('_', "-") - 6233
}); - 6234
- 6235
let mut output = match tool { - 6236
None => { - 6237
if let Some(actual_skill) = matched_skill { - 6238
ToolRunOutput::Err(format!( - 6239
r#"{{"type":"capability_kind_mismatch","name":{},"actual_kind":"skill","invocation":{{"tool":"skill","arguments":{{"name":{}}}}}}}"#, - 6240
serde_json::to_string(&call.name).unwrap_or_else(|_| "\"invalid\"".into()), - 6241
serde_json::to_string(actual_skill).unwrap_or_else(|_| "\"invalid\"".into()) - 6242
)) - 6243
} else if tools.iter().any(|t| t.name() == "mcp") { - 6244
ToolRunOutput::Err(format!( - 6245
r#"{{"type":"unknown_capability","requested_kind":"tool","name":{},"available_tools":{},"recovery_advice":"Tool '{}' is an external MCP capability. Invoke it via the 'mcp' tool: mcp(action: \"call\", server: \"<server_name>\", tool: \"{}\", arguments: {{ ... }})"}}"#, - 6246
serde_json::to_string(&call.name).unwrap_or_else(|_| "\"invalid\"".into()), - 6247
serde_json::to_string(&tools.iter().map(|t| t.name()).collect::<Vec<_>>()) - 6248
.unwrap_or_else(|_| "[]".into()), - 6249
call.name, - 6250
call.name - 6251
)) - 6252
} else { - 6253
ToolRunOutput::Err(format!( - 6254
r#"{{"type":"unknown_capability","requested_kind":"tool","name":{},"available_tools":{}}}"#, - 6255
serde_json::to_string(&call.name).unwrap_or_else(|_| "\"invalid\"".into()), - 6256
serde_json::to_string(&tools.iter().map(|t| t.name()).collect::<Vec<_>>()) - 6257
.unwrap_or_else(|_| "[]".into()) - 6258
)) - 6259
} - 6260
} - 6261
Some(tool) => { - 6262
let (sandbox_sink, mut sandbox_rx) = - 6263
vak_tools::SandboxEventSink::new_with_id(call.id.clone()); - 6264
let sandbox_sink = sandbox_sink.with_owner_session(session_id.to_string()); - 6265
let events_tx = events.clone(); - 6266
let forwarder = tokio::spawn(async move { - 6267
while let Some(sb_ev) = sandbox_rx.recv().await { - 6268
let _ = events_tx.send(AgentEvent::Sandbox(sb_ev)).await; - 6269
} - 6270
}); - 6271
- 6272
let ctx = vak_tools::ToolContext { - 6273
cwd: cwd.to_path_buf(), - 6274
cancel: cancel.child_token(), - 6275
sandbox: sandbox.cloned(), - 6276
sandbox_sink: Some(sandbox_sink), - 6277
agent_id: agent_id.map(|s| s.to_string()), - 6278
new_documents: Vec::new(), - 6279
}; - 6280
let tool = tool.clone(); - 6281
// The context (and its event sender) moves into the task and is - 6282
// dropped when it ends, so joining the forwarder below cannot - 6283
// wait on a sender this function still holds. - 6284
let res = tokio::spawn(async move { tool.execute(&call.input, &ctx).await }).await; - 6285
let output = match res { - 6286
Ok(out) => { - 6287
if let Some(delegated) = out.delegated { - 6288
yields - 6289
.lock() - 6290
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6291
.entry(call.id.clone()) - 6292
.or_default() - 6293
.delegated = Some(delegated); - 6294
} - 6295
let content = windowed_result(&call.id, out.content, yields); - 6296
if out.is_error { - 6297
ToolRunOutput::Err(content) - 6298
} else { - 6299
ToolRunOutput::Ok(content) - 6300
} - 6301
}
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.