- 3921
} - 3922
- 3923
let mut findings = String::new(); - 3924
- 3925
// 1) Regression obligations: everything proven green must stay so. - 3926
for cmd in self.obligations.clone() { - 3927
if cancel.is_cancelled() { - 3928
return None; - 3929
} - 3930
match self.run_audit_command(&cmd, cancel).await { - 3931
Ok(()) => {} - 3932
Err(err) => { - 3933
findings.push_str(&format!( - 3934
"REGRESSION: previously-green command failed now:\n $ {cmd}\n {err}\n" - 3935
)); - 3936
} - 3937
} - 3938
} - 3939
- 3940
let criteria = self - 3941
.active_goal - 3942
.as_ref() - 3943
.map(|g| g.criteria.clone()) - 3944
.unwrap_or_default(); - 3945
- 3946
// 2) Deterministic shell criteria. - 3947
let mut judged_criteria: Vec<String> = Vec::new(); - 3948
for criterion in &criteria { - 3949
if goal::is_shell_criterion(criterion) { - 3950
let cmd = goal::shell_command(criterion); - 3951
match self.run_audit_command(cmd, cancel).await { - 3952
Ok(()) => {} - 3953
Err(err) => { - 3954
findings.push_str(&format!("CRITERION FAILED: {criterion}\n {err}\n")); - 3955
} - 3956
} - 3957
judged_criteria.push(criterion.clone()); - 3958
} - 3959
} - 3960
- 3961
// 3) Judge call for remaining free-text criteria — only worth a - 3962
// model dispatch when deterministic checks already passed. - 3963
let text_criteria: Vec<String> = criteria - 3964
.iter() - 3965
.filter(|c| !goal::is_shell_criterion(c)) - 3966
.cloned() - 3967
.collect(); - 3968
if !text_criteria.is_empty() && findings.is_empty() { - 3969
match self.run_judge(&text_criteria, cancel, events).await { - 3970
Ok(verdicts) => { - 3971
for v in verdicts { - 3972
if v.verdict != "pass" { - 3973
findings.push_str(&format!( - 3974
"JUDGE {}: evidence: {}\n", - 3975
v.verdict.to_uppercase(), - 3976
v.evidence - 3977
)); - 3978
} - 3979
} - 3980
} - 3981
Err(e) => { - 3982
// Fail closed: an unavailable auditor cannot confirm done. - 3983
findings.push_str(&format!("AUDIT UNAVAILABLE: {e}\n")); - 3984
} - 3985
} - 3986
} - 3987
- 3988
if findings.is_empty() { - 3989
let mut session = self.session.lock().await; - 3990
let goal_snapshot = self.active_goal.clone(); - 3991
let sid = session - 3992
.header() - 3993
.map(|h| h.session_id.clone()) - 3994
.unwrap_or_default(); - 3995
if let Some(g) = goal_snapshot.as_ref() { - 3996
let _ = session.append_goal(vak_session::types::GoalEntry { - 3997
goal_id: format!("goal-{sid}"), - 3998
objective: g.objective.clone(), - 3999
criteria: g.criteria.clone(), - 4000
status: vak_session::types::GoalStatus::Done { audited: true }, - 4001
}); - 4002
} - 4003
drop(session); - 4004
self.active_goal = None; - 4005
return None; - 4006
} - 4007
- 4008
let audits_left = self - 4009
.active_goal - 4010
.as_mut() - 4011
.map(|g| { - 4012
g.audits_left = g.audits_left.saturating_sub(1); - 4013
g.audits_left - 4014
}) - 4015
.unwrap_or(0); - 4016
if audits_left > 0 { - 4017
Some(format!( - 4018
"[goal-audit] not verified yet:\n{findings}\nAddress these and finish again." - 4019
)) - 4020
} else { - 4021
// Budget exhausted: degrade to Unverified rather than trapping. - 4022
let mut session = self.session.lock().await; - 4023
let goal_snapshot = self.active_goal.clone(); - 4024
let sid = session - 4025
.header() - 4026
.map(|h| h.session_id.clone()) - 4027
.unwrap_or_default(); - 4028
if let Some(g) = goal_snapshot.as_ref() { - 4029
let _ = session.append_goal(vak_session::types::GoalEntry { - 4030
goal_id: format!("goal-{sid}"), - 4031
objective: g.objective.clone(), - 4032
criteria: g.criteria.clone(), - 4033
status: vak_session::types::GoalStatus::Unverified { - 4034
reason: findings.clone(), - 4035
}, - 4036
}); - 4037
} - 4038
drop(session); - 4039
self.active_goal = None; - 4040
None - 4041
} - 4042
} - 4043
- 4044
/// Runs one brokered bash command for auditing; Err = failure text. - 4045
async fn run_audit_command(&self, cmd: &str, cancel: &CancellationToken) -> Result<(), String> { - 4046
let tool = self - 4047
.config - 4048
.tools - 4049
.iter() - 4050
.find(|t| t.name() == "bash") - 4051
.ok_or_else(|| "no bash tool available for verification".to_string())?; - 4052
let cwd = self - 4053
.session - 4054
.lock() - 4055
.await - 4056
.header() - 4057
.map(|h| h.contract_cwd()) - 4058
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| ".".into())); - 4059
authorize( - 4060
&self.config, - 4061
&PendingToolCall { - 4062
id: "managed-verification".into(), - 4063
name: "bash".into(), - 4064
input: serde_json::json!({"command": cmd}), - 4065
}, - 4066
&cwd, - 4067
&self.run_call_counts, - 4068
&self.config.tools, - 4069
) - 4070
.await?; - 4071
let Some(sandbox) = self - 4072
.config - 4073
.sandbox - 4074
.as_ref() - 4075
.and_then(|sandbox| sandbox.read_only_variant()) - 4076
else { - 4077
return Err("shell verification requires a read-only sandbox".into()); - 4078
}; - 4079
let ctx = vak_tools::ToolContext { - 4080
cwd, - 4081
cancel: cancel.child_token(), - 4082
sandbox: Some(sandbox), - 4083
sandbox_sink: None, - 4084
agent_id: None, - 4085
new_documents: Vec::new(), - 4086
}; - 4087
let input = serde_json::json!({ "command": cmd }); - 4088
let out = match tokio::time::timeout( - 4089
std::time::Duration::from_secs(120), - 4090
tool.execute(&input, &ctx), - 4091
) - 4092
.await - 4093
{ - 4094
Ok(o) => o, - 4095
Err(_) => return Err("timed out after 120s".to_string()), - 4096
}; - 4097
if out.is_error { - 4098
Err(vak_session::bash_digest(&out.content)) - 4099
} else { - 4100
Ok(()) - 4101
} - 4102
} - 4103
- 4104
/// One skeptical judge dispatch over the transcript digest. - 4105
async fn run_judge( - 4106
&mut self, - 4107
text_criteria: &[String], - 4108
cancel: &CancellationToken, - 4109
events: &mpsc::Sender<AgentEvent>, - 4110
) -> Result<Vec<goal::CriterionVerdict>, String> { - 4111
// config.model reflects the per-turn effective route set by run_turn_inner. - 4112
let model = self.config.model.clone(); - 4113
let digest = { - 4114
let session = self.session.lock().await; - 4115
let msgs = session.derive_messages(); - 4116
goal::transcript_digest(&msgs, 24_000) - 4117
}; - 4118
let objective = self - 4119
.active_goal - 4120
.as_ref() - 4121
.map(|g| g.objective.clone()) - 4122
.unwrap_or_default(); - 4123
// MEA: environment facts over transcript claims. Unavailable delta - 4124
// is stated as such to the judge (UNKNOWN, never fabricated). - 4125
let workspace_delta = match &self.config.workspace_delta { - 4126
Some(p) => p.summary().ok(), - 4127
None => None, - 4128
}; - 4129
let req = goal::audit_request( - 4130
&model, - 4131
goal::audit_prompt( - 4132
&objective, - 4133
text_criteria, - 4134
&digest, - 4135
workspace_delta.as_deref(), - 4136
), - 4137
); - 4138
let mut ledger = StepLedger::new( - 4139
WorkPurpose::Verify, - 4140
self.provider.name(), - 4141
&model, - 4142
self.config.dispatch_ceiling.min(4), - 4143
); - 4144
let reply = self - 4145
.complete_with_reliability(&req, cancel, events, false, &mut ledger) - 4146
.await - 4147
.map_err(|e| e.to_string())?; - 4148
let _ = self - 4149
.session - 4150
.lock() - 4151
.await - 4152
.append_receipt(ledger.take_receipt()); - 4153
goal::parse_verdicts(&reply.text_content()) - 4154
} - 4155
- 4156
/// Reset-with-handoff rescue: one structured summary replaces the whole - 4157
/// projection; returns the handoff markdown. - 4158
async fn write_handoff( - 4159
&self, - 4160
est_tokens: u64, - 4161
original_prompt: &str, - 4162
cancel: &CancellationToken, - 4163
events: &mpsc::Sender<AgentEvent>, - 4164
) -> Result<String, LlmError> { - 4165
// config.model reflects the per-turn effective route set by run_turn_inner. - 4166
let model = self.config.model.clone(); - 4167
let digest = { - 4168
let session = self.session.lock().await; - 4169
let msgs = session.derive_messages(); - 4170
goal::transcript_digest(&msgs, 20_000) - 4171
}; - 4172
let objective_line = if self.active_goal.is_some() || !original_prompt.is_empty() { - 4173
format!("Original task: {original_prompt}\n\n") - 4174
} else { - 4175
String::new() - 4176
}; - 4177
let req = goal::handoff_request(&model, format!("{objective_line}{digest}")); - 4178
let mut ledger = StepLedger::new( - 4179
WorkPurpose::Summarize, - 4180
self.provider.name(), - 4181
&model, - 4182
self.config.dispatch_ceiling.min(3), - 4183
); - 4184
let _ = events - 4185
.send(AgentEvent::ContextCompacting { - 4186
estimated_tokens: est_tokens, - 4187
}) - 4188
.await; - 4189
let reply = self - 4190
.complete_with_reliability(&req, cancel, events, false, &mut ledger) - 4191
.await?; - 4192
let _ = self - 4193
.session - 4194
.lock() - 4195
.await - 4196
.append_receipt(ledger.take_receipt()); - 4197
Ok(reply.text_content()) - 4198
} - 4199
- 4200
/// Internal premature-completion gate. Returns a continuation reason - 4201
/// when the stop policy fires and budget remains. - 4202
async fn stop_gate( - 4203
&self, - 4204
prompt: &str, - 4205
response: &AssistantMessage, - 4206
receipts: &stop_policy::ReceiptSummary, - 4207
verification_stale: bool, - 4208
blocks_left: &mut u32, - 4209
user_completion_released: bool, - 4210
) -> Option<String> { - 4211
let policy = self.config.stop_policy.as_ref()?; - 4212
if user_completion_released { - 4213
return None; - 4214
} - 4215
let reason = policy.evaluate_receipts( - 4216
prompt, - 4217
&response.text_content(), - 4218
self.config.outcome.as_ref(), - 4219
receipts, - 4220
verification_stale, - 4221
)?; - 4222
if !matches!(reason, BlockReason::UserCompletionRequired) && *blocks_left == 0 { - 4223
return None; - 4224
} - 4225
if matches!(reason, BlockReason::UserCompletionRequired) { - 4226
return Some(reason.message()); - 4227
} - 4228
*blocks_left -= 1; - 4229
Some(reason.message()) - 4230
} - 4231
- 4232
/// Appends the continue nudge (model-visible => logged) and reports - 4233
/// whether the loop may continue within max_turns. - 4234
async fn guard_continue( - 4235
&mut self, - 4236
reason: String, - 4237
events: &mpsc::Sender<AgentEvent>, - 4238
turn: usize, - 4239
) -> bool { - 4240
if turn + 1 >= self.config.max_turns { - 4241
return false; - 4242
} - 4243
let _ = events - 4244
.send(AgentEvent::StopHookContinuation { - 4245
reason: reason.clone(), - 4246
}) - 4247
.await; - 4248
let _ = self - 4249
.session - 4250
.lock() - 4251
.await - 4252
.append_message(MessageRecord::control( - 4253
vak_intent::control::ControlKind::StopGuard, - 4254
format!("[stop-guard]: {reason}\nPlease continue."), - 4255
)); - 4256
let _ = events.send(AgentEvent::DraftDiscarded { turn }).await; - 4257
true - 4258
} - 4259
- 4260
/// Recovers the session ledger after a run (server/API consumers). - 4261
#[allow(clippy::panic)] - 4262
pub async fn into_session(self) -> SessionLog { - 4263
let mut config = self.config; - 4264
config.tools.clear(); - 4265
config.flow_dispatcher = None; - 4266
drop(config); - 4267
match Arc::try_unwrap(self.session) { - 4268
Ok(session) => session.into_inner(), - 4269
Err(_) => { - 4270
panic!("managed flow dispatcher retained the session after the agent stopped") - 4271
} - 4272
} - 4273
} - 4274
- 4275
/// The path a successful call of `name` with `input` delivers for - 4276
/// review, when `name` is one of this agent's tools - 4277
/// (`Tool::delivered_file`). - 4278
fn delivered_file(&self, name: &str, input: &Value) -> Option<String> { - 4279
self.config - 4280
.tools - 4281
.iter() - 4282
.find(|tool| tool.name() == name) - 4283
.and_then(|tool| tool.delivered_file(input)) - 4284
} - 4285
- 4286
/// Whether `name` is one of this agent's tools and it declares that a - 4287
/// successful call presents a card (`Tool::presents_cards`). - 4288
fn tool_presents_cards(&self, name: &str) -> bool { - 4289
self.config - 4290
.tools - 4291
.iter() - 4292
.any(|tool| tool.name() == name && tool.presents_cards()) - 4293
} - 4294
- 4295
/// Appends the model's response and returns its ledger entry id, so - 4296
/// callers that need to attribute something back to this exact message - 4297
/// (a fence-path `Presentation`, docs/design/68-context-engine.md §10) - 4298
/// don't have to re-derive it. `None` only on a session write failure. - 4299
async fn append_assistant(&self, response: &AssistantMessage) -> Option<String> { - 4300
let mut session = self.session.lock().await; - 4301
session - 4302
.append_message(MessageRecord { - 4303
message: response.clone().into_message(), - 4304
meta: Some(MessageMeta { - 4305
model: Some(response.model.clone()), - 4306
stop_reason: Some(format!("{:?}", response.stop_reason).to_lowercase()), - 4307
usage: Some(response.usage.clone()), - 4308
control: None, - 4309
..Default::default() - 4310
}), - 4311
}) - 4312
.ok() - 4313
.map(|entry| entry.id) - 4314
} - 4315
- 4316
async fn record_activity( - 4317
&self, - 4318
kind: vak_session::ActivityKind, - 4319
status: vak_session::ActivityStatus, - 4320
label: String, - 4321
detail: Option<String>, - 4322
data: std::collections::BTreeMap<String, String>, - 4323
) { - 4324
let now = chrono::Utc::now(); - 4325
let activity = vak_session::ActivityRecord { - 4326
activity_id: format!( - 4327
"activity-{}", - 4328
now.timestamp_nanos_opt() - 4329
.unwrap_or_else(|| now.timestamp_micros() * 1_000) - 4330
), - 4331
turn: None, - 4332
kind, - 4333
status, - 4334
label, - 4335
detail, - 4336
data, - 4337
}; - 4338
let _ = self.session.lock().await.append_activity(activity); - 4339
} - 4340
- 4341
/// Serializes a `CapacityProfile` and its key into an `Activity`'s - 4342
/// `data` map, the shape `SessionLog::latest_capacity_profile` reads - 4343
/// back (docs/design/68-context-engine.md §1, §4). - 4344
fn capacity_activity_data( - 4345
&self, - 4346
profile: &CapacityProfile, - 4347
) -> std::collections::BTreeMap<String, String> { - 4348
let mut data = std::collections::BTreeMap::new(); - 4349
let key = self - 4350
.config - 4351
.capacity_key - 4352
.clone() - 4353
.unwrap_or_else(|| capacity::ProfileKey { - 4354
provider: self.provider.name().to_string(), - 4355
model: self.config.model.clone(), - 4356
quantisation: profile.provenance.quantisation.clone(), - 4357
}); - 4358
if let Ok(key_json) = serde_json::to_string(&key) { - 4359
data.insert("key".into(), key_json); - 4360
} - 4361
if let Ok(profile_json) = serde_json::to_string(profile) { - 4362
data.insert("profile".into(), profile_json); - 4363
} - 4364
data - 4365
} - 4366
- 4367
/// Folds one turn's real usage into `self.config.capacity` (§1 - 4368
/// "Feedback") and records a `capacity-feedback` activity when a field - 4369
/// moved by more than `CAPACITY_FEEDBACK_CHANGE_THRESHOLD`. A no-op - 4370
/// when no profile was wired in for this run. - 4371
async fn record_capacity_usage_feedback( - 4372
&mut self, - 4373
request: &ChatRequest, - 4374
usage: &Usage, - 4375
first_token_latency_ms: Option<u64>, - 4376
) { - 4377
let Some(profile) = self.config.capacity.as_mut() else { - 4378
return; - 4379
}; - 4380
let before = profile.clone(); - 4381
let chars_sent = chat_request_chars(request); - 4382
let cache_miss = usage.cache_read_input_tokens.unwrap_or(0) == 0; - 4383
profile.observe_usage(chars_sent, usage, first_token_latency_ms, cache_miss); - 4384
let after = profile.clone(); - 4385
let delta = capacity_feedback_delta(&before, &after); - 4386
if delta.is_empty() { - 4387
return; - 4388
} - 4389
let mut data = self.capacity_activity_data(&after); - 4390
data.extend(delta); - 4391
self.record_activity( - 4392
vak_session::ActivityKind::CapacityFeedback, - 4393
vak_session::ActivityStatus::Succeeded, - 4394
"Capacity profile updated from usage".into(), - 4395
None, - 4396
data, - 4397
) - 4398
.await; - 4399
} - 4400
- 4401
/// Folds an explicit-instruction failure (required card not emitted, - 4402
/// required tool not called, a stop-policy block — the runtime already - 4403
/// classifies each) into `self.config.capacity` (§1 "Horizon - 4404
/// tightening") and records the change. A no-op when no profile was - 4405
/// wired in, or when the request was far from the current horizon. - 4406
async fn record_capacity_instruction_failure(&mut self, request_tokens: u64) { - 4407
let Some(profile) = self.config.capacity.as_mut() else { - 4408
return; - 4409
}; - 4410
let before = profile.clone(); - 4411
profile.observe_instruction_failure(request_tokens); - 4412
let after = profile.clone(); - 4413
if before.instruction_horizon.tokens == after.instruction_horizon.tokens { - 4414
return; - 4415
} - 4416
let mut data = self.capacity_activity_data(&after); - 4417
data.extend(capacity_feedback_delta(&before, &after)); - 4418
self.record_activity( - 4419
vak_session::ActivityKind::CapacityFeedback, - 4420
vak_session::ActivityStatus::Succeeded, - 4421
"Capacity horizon lowered by instruction-following feedback".into(), - 4422
Some(format!("request was {request_tokens} tokens")), - 4423
data, - 4424
) - 4425
.await; - 4426
} - 4427
- 4428
/// Model drift (docs/design/68-context-engine.md §7): the step's own - 4429
/// evidence, not a guess. The only trigger is a final answer that - 4430
/// exactly matches a prior turn's card narration — a verbatim repeat of - 4431
/// a past answer instead of addressing the current one. Returns the - 4432
/// drift reason for the nudge, or `None`. - 4433
/// - 4434
/// A tool's declared capability DOMAINS (e.g. `bash` → `code-exec`) - 4435
/// used to be compared against the current reading's open-vocabulary - 4436
/// SUBJECT domains (e.g. `engineering`) and flagged as drift whenever - 4437
/// they were disjoint — which they almost always were, since the two - 4438
/// vocabularies describe different things: a subject domain is never - 4439
/// supposed to gate control flow (docs/design/47-commitment-kernel.md). - 4440
/// Live: every `bash` call inside an ordinary engineering turn was - 4441
/// flagged ("serves none of the current directive's domains - 4442
/// (engineering)"), and three in a row ended the turn outright. - 4443
async fn detect_model_drift( - 4444
&self, - 4445
response: &AssistantMessage, - 4446
calls: &[PendingToolCall], - 4447
) -> Option<String> { - 4448
if !calls.is_empty() { - 4449
return None; - 4450
} - 4451
let past_narrations: Vec<String> = { - 4452
let session = self.session.lock().await; - 4453
TurnIndex::from_log(&session) - 4454
.turns - 4455
.iter() - 4456
.filter_map(|turn| turn.card.as_ref()) - 4457
.map(|card| card.answered.narration.trim().to_string()) - 4458
.filter(|narration| !narration.is_empty()) - 4459
.collect() - 4460
}; - 4461
let text = response.text_content(); - 4462
let trimmed = text.trim(); - 4463
if !trimmed.is_empty() && past_narrations.iter().any(|n| n == trimmed) { - 4464
return Some( - 4465
"repeated a previous turn's answer verbatim instead of addressing the current directive" - 4466
.to_string(), - 4467
); - 4468
} - 4469
None - 4470
} - 4471
- 4472
/// The degraded, honest completion returned when model drift exhausts - 4473
/// its steering budget (docs/design/68-context-engine.md §7) — same - 4474
/// shape as `degraded_outcome`'s tool-repair exhaustion - 4475
/// (docs/design/15-reliability.md), a different diagnostic label. - 4476
async fn degraded_drift_outcome(&self, last_reason: &str) -> TurnOutcome { - 4477
self.record_activity( - 4478
vak_session::ActivityKind::Diagnostic, - 4479
vak_session::ActivityStatus::Failed, - 4480
"model-drift-exhausted".into(), - 4481
Some(format!( - 4482
"three consecutive steps served a different directive than the current one; \ - 4483
run stopped rather than continuing to answer the wrong request. Last: {last_reason}" - 4484
)), - 4485
std::collections::BTreeMap::new(), - 4486
) - 4487
.await; - 4488
let response = AssistantMessage { - 4489
content: vec![ContentBlock::text( - 4490
"I kept drifting away from your current request across several steps and \ - 4491
could not stay on it within this turn's steering budget. Please restate what \ - 4492
you need now, or narrow the request, and I will address it directly." - 4493
.to_string(), - 4494
)], - 4495
stop_reason: StopReason::EndTurn, - 4496
usage: Usage::default(), - 4497
model: self.config.model.clone(), - 4498
response_id: None, - 4499
}; - 4500
TurnOutcome::Completed { response } - 4501
} - 4502
- 4503
/// One provider completion with watchdog, retry/backoff (honoring - 4504
/// Retry-After), circuit breaker, dispatch-ceiling enforcement, and - 4505
/// per-attempt receipt recording. When `forward` is true, stream - 4506
/// deltas are forwarded to `events`; otherwise they are drained. - 4507
async fn complete_with_reliability( - 4508
&self, - 4509
request: &ChatRequest, - 4510
cancel: &CancellationToken, - 4511
events: &mpsc::Sender<AgentEvent>, - 4512
forward: bool, - 4513
ledger: &mut StepLedger, - 4514
) -> Result<AssistantMessage, LlmError> { - 4515
// Frozen-ladder walk (Phase B): `ladder` holds FALLBACK legs; - 4516
// the primary provider/model always walks first. Ceiling, - 4517
// receipt, and endurance budget are shared across ALL legs -- - 4518
// walking the ladder is contract execution, never a switch. - 4519
let mut legs: Vec<(Arc<dyn Provider>, String)> = - 4520
Vec::with_capacity(1 + self.config.ladder.len()); - 4521
legs.push((self.provider.clone(), request.model.clone())); - 4522
legs.extend(self.config.ladder.iter().cloned()); - 4523
let mut leg_req = request.clone(); - 4524
let mut last_err: Option<LlmError> = None; - 4525
- 4526
'legs: for (li, (provider_arc, model)) in legs.iter().enumerate() { - 4527
leg_req.model = model.clone(); - 4528
let breaker_key = provider_arc.circuit_key(); - 4529
if let Some(breaker) = &self.config.circuit_breaker - 4530
&& let Err(open) = breaker.check_key(&breaker_key) - 4531
{ - 4532
last_err = Some(LlmError::Network(open.to_string())); - 4533
continue 'legs; - 4534
} - 4535
let route_provider = if li == 0 { - 4536
self.config - 4537
.provider_name - 4538
.clone() - 4539
.unwrap_or_else(|| provider_arc.name().to_string()) - 4540
} else { - 4541
self.config - 4542
.ladder_provider_names - 4543
.get(li - 1) - 4544
.cloned() - 4545
.unwrap_or_else(|| provider_arc.name().to_string()) - 4546
}; - 4547
ledger.receipt.stamp_leg(&route_provider, model); - 4548
// Per-leg tool inclusion (docs/design/68-context-engine.md §5): - 4549
// only Anthropic legs get the deferred schemas (withheld from - 4550
// the prefix there via `defer_loading`); every other provider - 4551
// sees core only, since its tool index is already in the - 4552
// prefix and `find_tools` is how it reaches the rest. - 4553
leg_req.tools = tools_for_leg(&request.tools, &route_provider); - 4554
if li > 0 && forward { - 4555
self.record_activity( - 4556
vak_session::ActivityKind::RouteFallback, - 4557
vak_session::ActivityStatus::Running, - 4558
"Route fallback".into(), - 4559
Some(format!("{}/{}", route_provider, model)), - 4560
[ - 4561
("provider".into(), route_provider.clone()), - 4562
("model".into(), model.clone()), - 4563
] - 4564
.into(), - 4565
) - 4566
.await; - 4567
let _ = events - 4568
.send(AgentEvent::RouteFallback { - 4569
to_provider: route_provider, - 4570
to_model: model.clone(), - 4571
}) - 4572
.await; - 4573
} - 4574
let mut attempt: u32 = 0; - 4575
loop { - 4576
if cancel.is_cancelled() { - 4577
return Err(LlmError::Aborted { partial: None }); - 4578
} - 4579
// Budget admission precedes every paid dispatch (Phase D). A - 4580
// denial becomes one bounded budget Ask; refusal -- or no - 4581
// approver, which is the unattended case -- fails the step - 4582
// permanently (never retried, never breaker-tripping). - 4583
if let Some(gate) = &self.config.spend_gate { - 4584
let session_id = self - 4585
.session - 4586
.lock() - 4587
.await - 4588
.header() - 4589
.map(|h| h.session_id.clone()) - 4590
.unwrap_or_default(); - 4591
let profile = self.effective_capacity_profile(); - 4592
let est_input = profile.estimate_tokens( - 4593
messages_chars(&leg_req.messages) - 4594
+ prefix_chars(leg_req.system.as_deref().unwrap_or(""), &leg_req.tools), - 4595
); - 4596
let check = SpendCheck { - 4597
model, - 4598
provider: provider_arc.name(), - 4599
session_id: &session_id, - 4600
est_input_tokens: est_input, - 4601
planned_output_tokens: self.config.max_output, - 4602
}; - 4603
if let Err(reason) = gate.authorize(&check).await { - 4604
let approved = match &self.config.approver { - 4605
Some(a) => { - 4606
a.approve( - 4607
"finops-budget", - 4608
&args_preview(&serde_json::json!({ - 4609
"model": model, - 4610
"reason": reason, - 4611
})), - 4612
&reason, - 4613
) - 4614
.await - 4615
} - 4616
None => false, - 4617
}; - 4618
if !approved { - 4619
return Err(LlmError::InvalidRequest(format!( - 4620
"budget admission denied: {reason}" - 4621
))); - 4622
} - 4623
// Raise-cap-once: the rest of THIS run is admitted. - 4624
gate.on_budget_approved(); - 4625
} - 4626
} - 4627
// Ceiling check happens before every paid dispatch; exhaustion - 4628
// surfaces as a plain error that the endurance loop treats as - 4629
// fail-closed (never transient). - 4630
if let Err(c) = ledger.budget.consume() { - 4631
return Err(LlmError::Network(c.to_string())); - 4632
} - 4633
let reason = if li > 0 && attempt == 0 { - 4634
AttemptReason::RouteFallback - 4635
} else if attempt == 0 { - 4636
AttemptReason::Initial - 4637
} else { - 4638
AttemptReason::Retry - 4639
}; - 4640
let started = std::time::Instant::now(); - 4641
- 4642
let provider_for_stream = provider_arc.clone(); - 4643
let req_for_stream = leg_req.clone(); - 4644
// A per-attempt CHILD token: the run's own `cancel` still - 4645
// propagates DOWN into it (a real user abort still ends - 4646
// this attempt immediately), but cancelling this one never - 4647
// propagates back UP. A watchdog timeout below cancels only - 4648
// this token, so a hung provider is abandoned without - 4649
// aborting the whole run. - 4650
let attempt_cancel = cancel.child_token(); - 4651
let stream_cancel = attempt_cancel.clone(); - 4652
let step = async move { - 4653
let mut stream = provider_for_stream - 4654
.stream(req_for_stream, stream_cancel) - 4655
.await?; - 4656
// First-token latency, wall clock from just before - 4657
// `.stream()` was called to the first event off the - 4658
// wire. Ollama fills `usage.prefill_ms` itself - 4659
// (preferred when present); this is what lets every - 4660
// other provider feed `prefill_tps` too - 4661
// (docs/design/68-context-engine.md §1 "Feedback"). - 4662
let mut first_token_ms: Option<u64> = None; - 4663
// Coalesce, never drop (clients render deltas, so a - 4664
// dropped one loses text): a still-full channel means - 4665
// the listener is behind, not gone, so anything that - 4666
// does not fit is queued -- merging consecutive deltas - 4667
// of the same kind/index via `StreamEvent::try_merge` - 4668
// -- and retried ahead of the next event, oldest - 4669
// first. Awaiting a send here would let a listener - 4670
// that stopped reading stall a step the provider had - 4671
// already finished, until the watchdog, so this never - 4672
// blocks mid-stream. - 4673
let mut pending: VecDeque<StreamEvent> = VecDeque::new(); - 4674
while let Some(ev) = futures::StreamExt::next(&mut stream).await { - 4675
if first_token_ms.is_none() { - 4676
first_token_ms = Some(started.elapsed().as_millis() as u64); - 4677
} - 4678
if !forward { - 4679
continue; - 4680
} - 4681
if drain_pending_stream_events(&mut pending, events) { - 4682
cancel.cancel(); - 4683
continue; - 4684
} - 4685
if !pending.is_empty() { - 4686
queue_stream_event(&mut pending, ev); - 4687
continue; - 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!(
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.