- 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!( - 4921
r#"{{"type":"invalid_arguments","message":"no evidence {id}"}}"# - 4922
)), - 4923
}, - 4924
} - 4925
} - 4926
- 4927
async fn execute_batch_inner( - 4928
&self, - 4929
calls: Vec<PendingToolCall>, - 4930
cancel: &CancellationToken, - 4931
events: &mpsc::Sender<AgentEvent>, - 4932
) -> Vec<(String, ToolRunOutput)> { - 4933
let key = |call: &PendingToolCall| { - 4934
format!( - 4935
"{}\u{0}{}", - 4936
call.name, - 4937
serde_json::to_string(&call.input).unwrap_or_default() - 4938
) - 4939
}; - 4940
let mut shown = self - 4941
.presented_cards - 4942
.lock() - 4943
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4944
.clone(); - 4945
let mut repeats: Vec<String> = Vec::new(); - 4946
let mut answered: Vec<(String, ToolRunOutput)> = Vec::new(); - 4947
let mut live_cards: HashMap<String, String> = HashMap::new(); - 4948
let mut live_deliveries: HashMap<String, (String, String)> = HashMap::new(); - 4949
let mut live = Vec::with_capacity(calls.len()); - 4950
for call in calls { - 4951
let mut deliveries = self - 4952
.deliveries - 4953
.lock() - 4954
.unwrap_or_else(std::sync::PoisonError::into_inner); - 4955
if let Some(path) = self.delivered_file(&call.name, &call.input) { - 4956
let delivery_key = key(&call); - 4957
if let Some(first) = deliveries.results.get(&delivery_key) { - 4958
answered.push(( - 4959
call.id.clone(), - 4960
ToolRunOutput::Ok(format!("{DRAFT_REPEAT_ACK}\n{first}")), - 4961
)); - 4962
continue; - 4963
} - 4964
live_deliveries.insert(call.id.clone(), (delivery_key, path)); - 4965
} - 4966
if vak_tools::canonical_tool_name(&call.name) == "bash" - 4967
&& let Some(command) = call.input.get("command").and_then(Value::as_str) - 4968
&& let Some(delivered) = copies_delivered_draft(command, &deliveries.paths).cloned() - 4969
{ - 4970
answered.push(( - 4971
call.id.clone(), - 4972
ToolRunOutput::Ok(format!( - 4973
"{DRAFT_COPY_REFUSED} {delivered} is a draft waiting for the person's review, and copying it out of .vak/scratch/ would skip that review. It reaches the workspace when they accept it. Answer with one sentence saying what you changed." - 4974
)), - 4975
)); - 4976
continue; - 4977
} - 4978
if self.tool_presents_cards(&call.name) - 4979
&& let Some(path) = call - 4980
.input - 4981
.pointer("/payload/artifact_path") - 4982
.and_then(Value::as_str) - 4983
&& let Some(delivered) = deliveries - 4984
.paths - 4985
.iter() - 4986
.find(|delivered| same_workspace_path(delivered, path)) - 4987
.cloned() - 4988
{ - 4989
deliveries.withheld_cards.insert(call.id.clone()); - 4990
answered.push(( - 4991
call.id.clone(), - 4992
ToolRunOutput::Ok(format!( - 4993
"{WITHHELD_CARD_ACK} {delivered} is already in front of the person as a draft they review with its change list, and a card cannot show it better. Do not present it again; answer with one sentence saying what you changed." - 4994
)), - 4995
)); - 4996
continue; - 4997
} - 4998
drop(deliveries); - 4999
if self.tool_presents_cards(&call.name) { - 5000
let card_key = key(&call); - 5001
if !shown.insert(card_key.clone()) { - 5002
repeats.push(call.id.clone()); - 5003
continue; - 5004
} - 5005
live_cards.insert(call.id.clone(), card_key); - 5006
} - 5007
live.push(call); - 5008
} - 5009
let mut results = self.execute_batch_calls(live, cancel, events).await; - 5010
{ - 5011
let mut presented = self - 5012
.presented_cards - 5013
.lock() - 5014
.unwrap_or_else(std::sync::PoisonError::into_inner); - 5015
for (id, output) in &results { - 5016
if let (Some(card_key), ToolRunOutput::Ok(_)) = (live_cards.get(id), output) { - 5017
presented.insert(card_key.clone()); - 5018
} - 5019
} - 5020
} - 5021
{ - 5022
let mut deliveries = self - 5023
.deliveries - 5024
.lock() - 5025
.unwrap_or_else(std::sync::PoisonError::into_inner); - 5026
for (id, output) in &results { - 5027
if let (Some((delivery_key, path)), ToolRunOutput::Ok(text)) = - 5028
(live_deliveries.get(id), output) - 5029
{ - 5030
deliveries - 5031
.results - 5032
.insert(delivery_key.clone(), text.clone()); - 5033
deliveries.paths.insert(path.clone()); - 5034
} - 5035
} - 5036
} - 5037
results.extend(answered); - 5038
results.extend( - 5039
repeats - 5040
.into_iter() - 5041
.map(|id| (id, ToolRunOutput::Ok(CARD_REPEAT_ACK.into()))), - 5042
); - 5043
results - 5044
} - 5045
- 5046
/// The turn's answer when a current value was asked for and nothing - 5047
/// was retrieved after the one repair (docs/design/68 §7): an honest - 5048
/// statement naming the last figure this conversation recorded and - 5049
/// when, never that figure presented as current. - 5050
async fn stale_data_outcome(&self) -> TurnOutcome { - 5051
let last_known = { - 5052
let session = self.session.lock().await; - 5053
let turn_id = session.latest_directive_entry_id(); - 5054
let entries = session.chain_to_root(); - 5055
session - 5056
.presentations() - 5057
.into_iter() - 5058
.rev() - 5059
.find(|(_, record)| Some(record.turn_id.as_str()) != turn_id.as_deref()) - 5060
.map(|(id, record)| { - 5061
let when = entries - 5062
.iter() - 5063
.find(|entry| entry.id == id) - 5064
.map(|entry| entry.ts.format("%Y-%m-%d %H:%M UTC").to_string()) - 5065
.unwrap_or_else(|| "an earlier turn".to_string()); - 5066
format!( - 5067
" The most recent figure in this conversation was recorded at {when}: {}.", - 5068
record.identity_digest - 5069
) - 5070
}) - 5071
.unwrap_or_default() - 5072
}; - 5073
self.record_activity( - 5074
vak_session::ActivityKind::Diagnostic, - 5075
vak_session::ActivityStatus::Failed, - 5076
"stale-data-refused".into(), - 5077
Some( - 5078
"a current value was asked for, nothing was retrieved this turn after one repair, \ - 5079
and no carried-over figure was presented as current" - 5080
.into(), - 5081
), - 5082
std::collections::BTreeMap::new(), - 5083
) - 5084
.await; - 5085
TurnOutcome::Completed { - 5086
response: AssistantMessage { - 5087
content: vec![ContentBlock::text(format!( - 5088
"I could not retrieve a current value on this turn, so I am not presenting a \ - 5089
carried-over figure as current.{last_known} Ask again when a retrieval tool \ - 5090
is available, or ask for the last known figure explicitly." - 5091
))], - 5092
stop_reason: StopReason::EndTurn, - 5093
usage: Usage::default(), - 5094
model: self.config.model.clone(), - 5095
response_id: None, - 5096
}, - 5097
} - 5098
} - 5099
- 5100
/// The turn's answer when a card repeatedly refuses to be about what was - 5101
/// asked (docs/design/68 §7): honest text instead of a wrong card. When - 5102
/// this turn's own retrieval actually succeeded, its raw result is - 5103
/// quoted rather than left out — the evidence exists, only the card - 5104
/// built from it did not. - 5105
async fn topic_mismatch_outcome(&self, evidence: Option<String>) -> TurnOutcome { - 5106
self.record_activity( - 5107
vak_session::ActivityKind::Diagnostic, - 5108
vak_session::ActivityStatus::Failed, - 5109
"topic-mismatch-refused".into(), - 5110
Some( - 5111
"a card was offered twice whose content did not match the directive; the turn \ - 5112
closed on an honest statement instead of showing a wrong card" - 5113
.into(), - 5114
), - 5115
std::collections::BTreeMap::new(),
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.