- 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
} - 6302
Err(join_err) => ToolRunOutput::Err(format!("tool task failed: {join_err}")), - 6303
}; - 6304
let _ = forwarder.await; - 6305
output - 6306
} - 6307
}; - 6308
- 6309
if let Some(hooks) = hooks { - 6310
let post_reason = match &output { - 6311
ToolRunOutput::Ok(content) => content.as_str(), - 6312
ToolRunOutput::Err(content) => content.as_str(), - 6313
}; - 6314
let post = vak_hooks::run_hooks_with_recorder( - 6315
hooks.clone(), - 6316
vak_hooks::HookEvent::PostToolUse, - 6317
session_id, - 6318
cwd, - 6319
Some((&hook_name, &hook_input)), - 6320
Some(post_reason), - 6321
cancel, - 6322
hook_recorder.map(|recorder| &**recorder), - 6323
) - 6324
.await; - 6325
if post.blocked && matches!(output, ToolRunOutput::Ok(_)) { - 6326
let reason = post.reason.unwrap_or_else(|| "flagged by hook".into()); - 6327
output = ToolRunOutput::Err(format!("{}\n[post-tool-use hook]: {reason}", post_reason)); - 6328
} else if post.blocked { - 6329
let reason = post.reason.unwrap_or_else(|| "flagged by hook".into()); - 6330
output = ToolRunOutput::Err(format!("{post_reason}\n[post-tool-use hook]: {reason}")); - 6331
} - 6332
} - 6333
- 6334
// Tool failures are values, but a terse error alone makes weaker models - 6335
// stop instead of repairing the call. Keep the original error intact and - 6336
// attach a bounded, non-authorizing recovery contract. The next model - 6337
// turn is the retry loop; permission, cancellation, and policy failures - 6338
// deliberately do not receive a retry suggestion. - 6339
if let ToolRunOutput::Err(content) = &mut output - 6340
&& let Some(hint) = tool_recovery_hint(content) - 6341
{ - 6342
content.push_str(hint); - 6343
} - 6344
- 6345
let result_content = match &output { - 6346
ToolRunOutput::Ok(content) | ToolRunOutput::Err(content) => content.as_str(), - 6347
}; - 6348
let _ = events - 6349
.send(AgentEvent::ToolCallEnd { - 6350
id: call.id.clone(), - 6351
name: call.name.clone(), - 6352
is_error: matches!(output, ToolRunOutput::Err(_)), - 6353
result_preview: result_preview(result_content), - 6354
}) - 6355
.await; - 6356
if let Some(recorder) = tool_activity_recorder { - 6357
recorder( - 6358
&call.name, - 6359
&activity_input, - 6360
matches!(output, ToolRunOutput::Ok(_)), - 6361
started.elapsed().as_millis() as u64, - 6362
); - 6363
} - 6364
- 6365
(call.id, output) - 6366
} - 6367
- 6368
/// Returns the model recovery contract for a correctable tool error. The - 6369
/// classification that decides this is owned by `vak_tools::ToolErrorKind` — - 6370
/// this function is intentionally a thin shim so the hint and the repair - 6371
/// budget in the turn loop can never drift onto a different keyword list. - 6372
fn tool_recovery_hint(error: &str) -> Option<&'static str> { - 6373
if ToolErrorKind::classify(error).is_correctable() { - 6374
Some( - 6375
"\n[recovery] Treat this as a failed attempt. Inspect the error and the admitted tool/schema inventory, then make at most one corrected or alternative call. Do not repeat identical arguments. If the failure is environmental or the corrected call is unsafe, explain the blocker instead.", - 6376
) - 6377
} else { - 6378
None - 6379
} - 6380
} - 6381
- 6382
/// Reconcile the run repair budget against this turn's correctable tool - 6383
/// failures. Returns `Some(TurnOutcome)` only when the budget is exhausted - 6384
/// and the loop must stop repairing; otherwise `None` (continue to the next - 6385
/// model dispatch). This is the "loop where these don't happen" for tool - 6386
/// failures of any class: the recovery hint is the first nudge; a - 6387
/// system-authored, schema-resurfacing directive is the second; a bounded - 6388
/// degraded stop is the third. - 6389
async fn reconcile_repair_budget( - 6390
agent: &mut Agent, - 6391
failed_correctable: &[(String, ToolErrorKind)], - 6392
) -> Option<TurnOutcome> { - 6393
if failed_correctable.is_empty() { - 6394
agent.repair.consecutive_failed_turns = 0; - 6395
return None; - 6396
} - 6397
if agent.repair.exhausted { - 6398
return None; - 6399
} - 6400
agent.repair.consecutive_failed_turns += 1; - 6401
- 6402
if agent.repair.consecutive_failed_turns == REPAIR_DIRECTIVE_TURN { - 6403
// The model has now failed to self-repair the same fault across two - 6404
// turns: a text hint is no longer enough. The loop takes over and - 6405
// resurfaces the exact admitted schema for the rejected tools so the - 6406
// repair is no longer a guess. - 6407
inject_repair_directive(agent, failed_correctable).await; - 6408
} - 6409
- 6410
if agent.repair.consecutive_failed_turns > MAX_REPAIR_TURNS { - 6411
agent.repair.exhausted = true; - 6412
let outcome = degraded_outcome(agent, failed_correctable).await; - 6413
return Some(outcome); - 6414
} - 6415
None - 6416
} - 6417
- 6418
/// Append an authoritative repair directive that resurfaces the admitted - 6419
/// schema for each tool the model could not get right. Unlike the per-call - 6420
/// `[recovery]` hint, this is issued by the loop itself (not the model) - 6421
/// when the model has demonstrated it cannot repair the failure unprompted. - 6422
async fn inject_repair_directive(agent: &Agent, failed: &[(String, ToolErrorKind)]) { - 6423
let remaining = MAX_REPAIR_TURNS.saturating_sub(agent.repair.consecutive_failed_turns - 1); - 6424
let mut parts: Vec<String> = vec![format!( - 6425
"{} The run is stuck on correctable tool failures that \ - 6426
were not repaired across turns. Do not repeat the failing call shape; \ - 6427
re-issue with the exact arguments this tool requires. The run will \ - 6428
stop retrying after {} more failed repair turn(s).", - 6429
vak_intent::control::ControlKind::RepairDirective.marker(), - 6430
remaining.max(1) - 6431
)]; - 6432
let mut seen = std::collections::HashSet::new(); - 6433
for (name, _kind) in failed { - 6434
if !seen.insert(name.clone()) { - 6435
continue; - 6436
} - 6437
let mut found = false; - 6438
for tool in &agent.config.tools { - 6439
if tool.name() != *name { - 6440
continue; - 6441
} - 6442
found = true; - 6443
let schema_str = - 6444
serde_json::to_string(&tool.schema()).unwrap_or_else(|_| "{}".to_string()); - 6445
parts.push(format!( - 6446
"\nAdmitted tool `{}` (re-surfaced verbatim):\ndescription: {}\ninput_schema: {}", - 6447
tool.name(), - 6448
tool.description(), - 6449
schema_str - 6450
)); - 6451
break; - 6452
} - 6453
if !found { - 6454
// Unknown tool name: list what IS admitted so the model can map - 6455
// the rejected call onto an admitted one (e.g. use the `mcp` - 6456
// broker instead of a raw capability name). - 6457
let admitted: Vec<String> = agent - 6458
.config - 6459
.tools - 6460
.iter() - 6461
.map(|t| t.name().to_string()) - 6462
.collect(); - 6463
parts.push(format!( - 6464
"\nTool `{}` is not in the admitted set for this turn. Admitted \ - 6465
tools: {}. Re-issue using an admitted tool.", - 6466
name, - 6467
admitted.join(", ") - 6468
)); - 6469
} - 6470
} - 6471
// Runtime-authored, so tagged: it must never read as the person's words, - 6472
// to the model or to any client (AGENTS.md, "typed, never sniffed"). - 6473
let _ = agent - 6474
.session - 6475
.lock() - 6476
.await - 6477
.append_message(MessageRecord::control( - 6478
vak_intent::control::ControlKind::RepairDirective, - 6479
parts.join("\n\n"), - 6480
)); - 6481
} - 6482
- 6483
/// Build the degraded, honest completion returned when the repair budget is - 6484
/// exhausted: record a diagnostic (append-only, model-visible) and return a - 6485
/// system-authored answer that states the failure instead of inventing one. - 6486
async fn degraded_outcome(agent: &Agent, failed: &[(String, ToolErrorKind)]) -> TurnOutcome { - 6487
let summary = failed - 6488
.iter() - 6489
.map(|(name, kind)| format!("- `{name}`: correctable fault ({kind:?})")) - 6490
.collect::<Vec<_>>() - 6491
.join("\n"); - 6492
agent - 6493
.record_activity( - 6494
vak_session::ActivityKind::Diagnostic, - 6495
vak_session::ActivityStatus::Failed, - 6496
"Tool repair exhausted".into(), - 6497
Some(format!( - 6498
"correctable tool failures were not repaired within the run repair \ - 6499
budget ({} repair turns); run stopped rather than signing a false \ - 6500
complete", - 6501
MAX_REPAIR_TURNS - 6502
)), - 6503
std::collections::BTreeMap::from([ - 6504
( - 6505
"repair_turns".into(), - 6506
agent.repair.consecutive_failed_turns.to_string(), - 6507
), - 6508
( - 6509
"failed_tools".into(), - 6510
failed - 6511
.iter() - 6512
.map(|(n, _)| n.as_str()) - 6513
.collect::<Vec<_>>() - 6514
.join(","), - 6515
), - 6516
]), - 6517
) - 6518
.await; - 6519
let response = AssistantMessage { - 6520
content: vec![ContentBlock::text(format!( - 6521
"I attempted the requested work, but the supporting tool calls failed \ - 6522
and could not be repaired within the run's recovery budget. I will \ - 6523
not sign off a fabricated answer. What failed:\n{summary}\n\nTo \ - 6524
continue, either correct the inputs above and re-run, or widen the \ - 6525
workspace capabilities / permissions if the failure is an admission \ - 6526
gate." - 6527
))], - 6528
stop_reason: StopReason::EndTurn, - 6529
usage: Usage::default(), - 6530
model: agent.config.model.clone(), - 6531
response_id: None, - 6532
}; - 6533
TurnOutcome::Completed { response } - 6534
} - 6535
- 6536
/// Effectful repetitions reach a human at this count. An identical `read` - 6537
/// remains inside its normal permission check for two more attempts, then - 6538
/// fails as a tool value without asking a human to approve rereading a file. - 6539
const DOOM_LOOP_THRESHOLD: u32 = 3; - 6540
const READ_REPEAT_LIMIT: u32 = 5; - 6541
- 6542
fn repeat_guard_decision(name: &str, count: u32) -> Option<Decision> { - 6543
if name == "read" && count >= READ_REPEAT_LIMIT { - 6544
return Some(Decision::Deny { - 6545
reason: format!( - 6546
"identical read call repeated ×{count} this run; use the successful result already returned or explain why a different inspection is needed" - 6547
), - 6548
}); - 6549
} - 6550
if name != "read" && count >= DOOM_LOOP_THRESHOLD { - 6551
return Some(Decision::Ask { - 6552
reason: format!("identical {name} call repeated ×{count} this run"), - 6553
source: AskSource::CircuitBreaker, - 6554
}); - 6555
} - 6556
None - 6557
} - 6558
- 6559
#[cfg(test)] - 6560
mod repeat_guard_tests { - 6561
use super::{AskSource, Decision, repeat_guard_decision}; - 6562
- 6563
#[test] - 6564
fn repeated_reads_stay_under_normal_permission_before_bounded_denial() { - 6565
assert!(repeat_guard_decision("read", 3).is_none()); - 6566
assert!(repeat_guard_decision("read", 4).is_none()); - 6567
assert!(matches!( - 6568
repeat_guard_decision("read", 5), - 6569
Some(Decision::Deny { .. }) - 6570
)); - 6571
assert!(matches!( - 6572
repeat_guard_decision("write", 3), - 6573
Some(Decision::Ask { - 6574
source: AskSource::CircuitBreaker, - 6575
.. - 6576
}) - 6577
)); - 6578
} - 6579
} - 6580
const MAX_TOOL_INPUT_CHARS: usize = 32_000; - 6581
- 6582
/// After this many **consecutive** turns that end with an unresolved - 6583
/// correctable tool failure (`RepairState::consecutive_failed_turns` exceeds - 6584
/// it), the run stops re-dispatching and degrades the outcome instead of - 6585
/// spinning on failing tool calls. This bounds model-guided repair so a weak - 6586
/// model that ignores the recovery hint cannot burn the whole turn budget - 6587
/// on the same fault class. - 6588
/// Three consecutive model-drift events (docs/design/68-context-engine.md - 6589
/// §7) end the turn with the degraded outcome, mirroring `MAX_REPAIR_TURNS` - 6590
/// for tool repair: enough room for one bad step to self-correct after a - 6591
/// steering nudge, not enough to spend the whole turn serving the wrong - 6592
/// directive. - 6593
const MODEL_DRIFT_EXHAUSTION_THRESHOLD: u32 = 3; - 6594
- 6595
const MAX_REPAIR_TURNS: u32 = 2; - 6596
/// On the Nth consecutive correctable-failure turn the system stops relying - 6597
/// on a text hint alone: it injects an authoritative, schema-resurfacing - 6598
/// directive so the model is no longer guessing what shape was rejected. - 6599
const REPAIR_DIRECTIVE_TURN: u32 = 2; - 6600
- 6601
/// Per-run account of model-guided tool recovery. The loop is: hint on the - 6602
/// first failure, a system-authored directive on the second, and a bounded - 6603
/// degraded stop on the third — instead of unlimited spin or a silent - 6604
/// false "complete". Reset at the start of every `run`. - 6605
#[derive(Debug, Default)] - 6606
struct RepairState { - 6607
/// Consecutive turns ending with one or more unresolved correctable tool - 6608
/// failures. Resets to 0 when a turn produces no correctable failures. - 6609
consecutive_failed_turns: u32, - 6610
/// Set once the run repair budget is exhausted; the loop must not keep - 6611
/// dispatching for repair after this. - 6612
exhausted: bool, - 6613
} - 6614
- 6615
impl RepairState { - 6616
fn reset(&mut self) { - 6617
*self = Self::default(); - 6618
} - 6619
} - 6620
- 6621
async fn authorize( - 6622
config: &AgentConfig, - 6623
call: &PendingToolCall, - 6624
cwd: &std::path::Path, - 6625
run_call_counts: &std::sync::Mutex<HashMap<String, u32>>, - 6626
tools: &[Arc<dyn Tool>], - 6627
) -> Result<(), String> { - 6628
let input_chars = serde_json::to_string(&call.input) - 6629
.map(|input| input.chars().count()) - 6630
.unwrap_or(MAX_TOOL_INPUT_CHARS.saturating_add(1)); - 6631
if input_chars > MAX_TOOL_INPUT_CHARS { - 6632
return Err(format!( - 6633
"tool call arguments exceed the {}-character safety limit; reduce the arguments and retry", - 6634
MAX_TOOL_INPUT_CHARS - 6635
)); - 6636
} - 6637
if let Some(tool) = tools.iter().find(|tool| tool.name() == call.name) { - 6638
vak_tools::validate_input(&tool.schema(), &call.input)?; - 6639
if let Some(reason) = tool.refusal(&call.input) { - 6640
return Err(reason); - 6641
} - 6642
} - 6643
if config - 6644
.revocation_check - 6645
.as_ref() - 6646
.is_some_and(|check| check(&call.name, &call.input)) - 6647
{ - 6648
return Err(format!( - 6649
"capability `{}` was revoked during this turn", - 6650
call.name - 6651
)); - 6652
} - 6653
let key = format!( - 6654
"{}\u{0}{}", - 6655
call.name, - 6656
serde_json::to_string(&call.input).unwrap_or_default() - 6657
); - 6658
let n = { - 6659
let mut counts = run_call_counts - 6660
.lock() - 6661
.unwrap_or_else(std::sync::PoisonError::into_inner); - 6662
let entry = counts.entry(key).or_insert(0); - 6663
*entry += 1; - 6664
*entry - 6665
}; - 6666
let decision = if let Some(decision) = repeat_guard_decision(&call.name, n) { - 6667
decision - 6668
} else if let Some(engine) = &config.permission { - 6669
engine.evaluate(&call.name, &call.input, config.mode, cwd) - 6670
} else { - 6671
Decision::Allow - 6672
}; - 6673
match decision { - 6674
Decision::Allow => Ok(()), - 6675
Decision::Deny { reason } => Err(reason), - 6676
Decision::Ask { reason, source } => { - 6677
if auto_approve( - 6678
config.approval_mode, - 6679
source, - 6680
&call.name, - 6681
&call.input, - 6682
config.mode, - 6683
config.sandbox.is_some(), - 6684
cwd, - 6685
) { - 6686
if config - 6687
.revocation_check - 6688
.as_ref() - 6689
.is_some_and(|check| check(&call.name, &call.input)) - 6690
{ - 6691
return Err(format!( - 6692
"capability `{}` was revoked during this turn", - 6693
call.name - 6694
)); - 6695
} - 6696
return Ok(()); - 6697
} - 6698
// A live envelope pre-authorizes what it covers — within the - 6699
// authority already in force, never beyond it. An ask a rule or - 6700
// the circuit breaker raised is not one a grant stands in for. - 6701
if !matches!(source, AskSource::Rule | AskSource::CircuitBreaker) - 6702
&& config - 6703
.envelope_check - 6704
.as_ref() - 6705
.and_then(|check| check(&call.name, &call.input)) - 6706
.is_some() - 6707
{ - 6708
if config - 6709
.revocation_check - 6710
.as_ref() - 6711
.is_some_and(|check| check(&call.name, &call.input)) - 6712
{ - 6713
return Err(format!( - 6714
"capability `{}` was revoked during this turn", - 6715
call.name - 6716
)); - 6717
} - 6718
return Ok(()); - 6719
} - 6720
match &config.approver { - 6721
Some(a) - 6722
if a.approve(&call.name, &args_preview(&call.input), &reason) - 6723
.await => - 6724
{ - 6725
if config - 6726
.revocation_check - 6727
.as_ref() - 6728
.is_some_and(|check| check(&call.name, &call.input)) - 6729
{ - 6730
Err(format!( - 6731
"capability `{}` was revoked while approval was pending", - 6732
call.name - 6733
)) - 6734
} else { - 6735
Ok(()) - 6736
} - 6737
} - 6738
// Say who actually refused. On an unattended surface the - 6739
// gate was never put to a person, and reporting it as - 6740
// "denied by user" sent the model looking for a different - 6741
// tool — and the operator looking for a user who had done - 6742
// nothing — instead of naming the policy that decided. - 6743
Some(a) if !a.answerable() => Err(format!( - 6744
"unattended surface: {reason}. No approver is configured to \ - 6745
answer it, so this capability cannot be used on this turn" - 6746
)), - 6747
Some(_) => Err(format!("denied by user: {reason}")), - 6748
None => Err(format!("{reason} (no approver available)")), - 6749
} - 6750
} - 6751
} - 6752
} - 6753
- 6754
enum ToolRunOutput { - 6755
Ok(String), - 6756
Err(String), - 6757
} - 6758
- 6759
fn normalize_response_tool_uses( - 6760
response: &mut AssistantMessage, - 6761
tools: &[Arc<dyn Tool>], - 6762
mcp_index: &std::collections::HashMap<String, String>, - 6763
) { - 6764
for block in &mut response.content { - 6765
let ContentBlock::ToolUse { id, name, input } = block else { - 6766
continue; - 6767
}; - 6768
let call = normalize_schema_wrapper( - 6769
normalize_mcp_call( - 6770
normalize_tool_call(PendingToolCall { - 6771
id: id.clone(), - 6772
name: name.clone(), - 6773
input: input.clone(), - 6774
}), - 6775
mcp_index, - 6776
), - 6777
tools, - 6778
); - 6779
*name = call.name; - 6780
*input = call.input; - 6781
} - 6782
} - 6783
- 6784
/// If `response` carries no structured `tool_use` block but its text holds - 6785
/// one or more explicit tool-call envelopes for a tool loaded this turn -- - 6786
/// `<tool_call>{json}</tool_call>` or a fenced block tagged exactly - 6787
/// `tool_call`/`tool_use` -- rewrites `response` in place: each envelope's - 6788
/// exact span is excised from the text (any surrounding prose survives), - 6789
/// replaced by a real `ContentBlock::ToolUse` with a fresh id, and - 6790
/// `stop_reason` becomes `ToolUse`, so the ledger records a valid - 6791
/// tool_use/tool_result pair instead of the envelope text. Called BEFORE - 6792
/// the response is appended to the ledger. - 6793
/// - 6794
/// Prose CODE EXAMPLES are deliberately not recognised here any more -- a - 6795
/// ```` ```bash ```` fence, `bash -c "..."`, `bash(command=...)` and - 6796
/// `write(path=..., content=...)` used to be executed as if the model had - 6797
/// asked for them, even when written only to ILLUSTRATE a command rather - 6798
/// than invoke it (live: a prose answer's ```bash example ran and created - 6799
/// a file the model never asked to create). The envelope path is the one - 6800
/// fallback that remains, because unlike a shell fence it is unambiguous: - 6801
/// nothing else in ordinary prose looks like `<tool_call>{"name": - 6802
/// ...}</tool_call>`, and it is gated further by naming a tool actually - 6803
/// loaded this turn. - 6804
fn extract_tool_calls( - 6805
response: &mut AssistantMessage, - 6806
loaded: &[Arc<dyn Tool>], - 6807
) -> Vec<PendingToolCall> { - 6808
let structured: Vec<PendingToolCall> = response - 6809
.content - 6810
.iter() - 6811
.filter_map(|b| match b { - 6812
ContentBlock::ToolUse { id, name, input } => Some(PendingToolCall { - 6813
id: id.clone(), - 6814
name: name.clone(), - 6815
input: input.clone(), - 6816
}), - 6817
_ => None, - 6818
}) - 6819
.collect(); - 6820
if !structured.is_empty() { - 6821
return structured; - 6822
} - 6823
let text = response.text_content(); - 6824
if text.trim().is_empty() { - 6825
return Vec::new(); - 6826
} - 6827
let envelopes = parse_tool_call_envelopes(&text, loaded); - 6828
if envelopes.is_empty() { - 6829
return Vec::new(); - 6830
} - 6831
let mut remaining = text; - 6832
for (_, span) in envelopes.iter().rev() { - 6833
remaining.replace_range(span.clone(), ""); - 6834
} - 6835
let remaining = remaining.trim().to_string(); - 6836
response - 6837
.content - 6838
.retain(|block| !matches!(block, ContentBlock::Text { .. })); - 6839
if !remaining.is_empty() { - 6840
response.content.insert(0, ContentBlock::text(remaining)); - 6841
} - 6842
let calls: Vec<PendingToolCall> = envelopes.into_iter().map(|(call, _)| call).collect(); - 6843
for call in &calls { - 6844
response.content.push(ContentBlock::ToolUse { - 6845
id: call.id.clone(), - 6846
name: call.name.clone(), - 6847
input: call.input.clone(), - 6848
}); - 6849
} - 6850
response.stop_reason = StopReason::ToolUse; - 6851
calls - 6852
} - 6853
- 6854
/// Every recognised envelope in `text` that names a tool loaded this turn, - 6855
/// paired with the exact byte span (open tag through close tag) it - 6856
/// occupies so the caller can excise just the envelope and keep any - 6857
/// surrounding prose. - 6858
fn parse_tool_call_envelopes( - 6859
text: &str, - 6860
loaded: &[Arc<dyn Tool>], - 6861
) -> Vec<(PendingToolCall, std::ops::Range<usize>)> { - 6862
extract_tool_call_blocks(text) - 6863
.into_iter() - 6864
.filter_map(|(body, span)| { - 6865
let val = serde_json::from_str::<serde_json::Value>(&body).ok()?; - 6866
json_to_tool_call(&val, loaded).map(|call| (call, span)) - 6867
}) - 6868
.collect() - 6869
} - 6870
- 6871
/// Scans `text` for `<tool_call>...</tool_call>` and fenced blocks tagged - 6872
/// exactly `tool_call` or `tool_use`, returning each one's trimmed JSON - 6873
/// body alongside the full span it occupies, in the order they appear. A - 6874
/// block whose body has neither `"name"` nor `"tool"` is skipped before it - 6875
/// ever reaches JSON parsing -- not every fenced block a model writes is a - 6876
/// call. - 6877
fn extract_tool_call_blocks(text: &str) -> Vec<(String, std::ops::Range<usize>)> { - 6878
let mut blocks = Vec::new(); - 6879
for tag in ["<tool_call>", "```tool_call", "```tool_use"] { - 6880
let mut cursor = 0; - 6881
while let Some(start_idx) = text[cursor..].find(tag) { - 6882
let open_start = cursor + start_idx; - 6883
let body_start = open_start + tag.len(); - 6884
let close_tag = if tag.starts_with('<') { - 6885
"</tool_call>" - 6886
} else { - 6887
"```" - 6888
}; - 6889
let Some(end_idx) = text[body_start..].find(close_tag) else { - 6890
break; - 6891
}; - 6892
let body_end = body_start + end_idx; - 6893
let close_end = body_end + close_tag.len(); - 6894
let body = text[body_start..body_end].trim().to_string(); - 6895
if body.contains("\"name\"") || body.contains("\"tool\"") { - 6896
blocks.push((body, open_start..close_end)); - 6897
} - 6898
cursor = close_end; - 6899
} - 6900
} - 6901
blocks.sort_by_key(|(_, span)| span.start); - 6902
blocks - 6903
} - 6904
- 6905
/// Builds a call from an envelope's parsed JSON body: `name` (or `tool`) - 6906
/// must name a tool loaded this turn, or the envelope is ignored -- - 6907
/// otherwise this fallback could invoke anything a model happened to spell - 6908
/// out. `arguments`/`input`/`parameters`, when present, must be a JSON - 6909
/// object; absent defaults to `{}`, still a valid, argument-less call. - 6910
fn json_to_tool_call(val: &serde_json::Value, loaded: &[Arc<dyn Tool>]) -> Option<PendingToolCall> { - 6911
let name = val - 6912
.get("name") - 6913
.or_else(|| val.get("tool")) - 6914
.and_then(|v| v.as_str())?; - 6915
if !loaded.iter().any(|tool| tool.name() == name) { - 6916
return None; - 6917
} - 6918
let input = match val - 6919
.get("arguments") - 6920
.or_else(|| val.get("input")) - 6921
.or_else(|| val.get("parameters")) - 6922
{ - 6923
Some(value) if value.is_object() => value.clone(), - 6924
Some(_) => return None, - 6925
None => serde_json::json!({}), - 6926
}; - 6927
Some(PendingToolCall { - 6928
id: format!("call_txt_{:08x}", rand_jitter(u64::MAX)), - 6929
name: name.to_string(), - 6930
input, - 6931
}) - 6932
} - 6933
- 6934
fn backoff_delay(attempt: u32, retry_after_secs: Option<u64>, base_ms: u64) -> std::time::Duration { - 6935
if let Some(secs) = retry_after_secs { - 6936
return std::time::Duration::from_secs(secs.max(1)); - 6937
} - 6938
let exp = base_ms.saturating_mul(1u64 << (attempt - 1).min(6)); - 6939
let jitter = rand_jitter(exp); - 6940
std::time::Duration::from_millis((exp / 2).max(1).saturating_add(jitter).min(30_000)) - 6941
} - 6942
- 6943
fn rand_jitter(ms: u64) -> u64 { - 6944
use std::sync::atomic::{AtomicU64, Ordering}; - 6945
static STATE: AtomicU64 = AtomicU64::new(0); - 6946
let x = STATE - 6947
.fetch_add(0x9E3779B97F4A7C15, Ordering::Relaxed) - 6948
.wrapping_add(0x9E3779B97F4A7C15); - 6949
(x >> 33) % ms.max(2) - 6950
} - 6951
- 6952
/// Extracts the `StreamEvent` back out of a failed `try_send`'s returned - 6953
/// `AgentEvent` -- every event the stream-forwarding loop sends is - 6954
/// `AgentEvent::Stream`, but `TrySendError`'s payload is generic over the - 6955
/// channel's whole message type. `None` is unreachable in practice (this is - 6956
/// only ever called on a value this module itself just wrapped) but is - 6957
/// handled as a silent no-op rather than assumed, since asserting it would - 6958
/// mean panicking on a channel error. - 6959
fn into_stream_event(event: AgentEvent) -> Option<StreamEvent> { - 6960
match event { - 6961
AgentEvent::Stream(ev) => Some(ev), - 6962
_ => None, - 6963
} - 6964
} - 6965
- 6966
/// Drains `pending` into `events` oldest-first without blocking (§6 below). - 6967
/// Returns `true` once the listener is confirmed gone (`Closed`); a - 6968
/// still-full channel (`Full`) just leaves the rest queued for the next - 6969
/// call -- backpressure is not the same as nobody listening. - 6970
fn drain_pending_stream_events( - 6971
pending: &mut VecDeque<StreamEvent>, - 6972
events: &mpsc::Sender<AgentEvent>, - 6973
) -> bool { - 6974
while let Some(event) = pending.pop_front() { - 6975
match events.try_send(AgentEvent::Stream(event)) { - 6976
Ok(()) => {} - 6977
Err(mpsc::error::TrySendError::Full(sent)) => { - 6978
if let Some(event) = into_stream_event(sent) { - 6979
pending.push_front(event); - 6980
} - 6981
return false; - 6982
} - 6983
Err(mpsc::error::TrySendError::Closed(_)) => return true, - 6984
} - 6985
} - 6986
false - 6987
} - 6988
- 6989
/// Queues one stream event for forwarding, merging it into the last still- - 6990
/// pending event when `StreamEvent::try_merge` allows it, so a slow - 6991
/// listener's backlog stays one entry per in-progress block instead of - 6992
/// growing one entry per delta (docs/design/68-context-engine.md §6: - 6993
/// lossless streaming -- coalesce under backpressure, never drop). - 6994
fn queue_stream_event(pending: &mut VecDeque<StreamEvent>, event: StreamEvent) { - 6995
match pending.back_mut() { - 6996
Some(last) => { - 6997
if let Some(event) = last.try_merge(event) { - 6998
pending.push_back(event); - 6999
} - 7000
}
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.