- 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(), - 5116
) - 5117
.await; - 5118
let text = match evidence { - 5119
Some(found) => format!( - 5120
"I could not turn what I found into a card that actually answers this, so here is \ - 5121
the raw result instead:\n\n{found}" - 5122
), - 5123
None => "I could not produce a card that matches what was asked, and had nothing else \ - 5124
to fall back on this turn. Could you rephrase the question?" - 5125
.to_string(), - 5126
}; - 5127
TurnOutcome::Completed { - 5128
response: AssistantMessage { - 5129
content: vec![ContentBlock::text(text)], - 5130
stop_reason: StopReason::EndTurn, - 5131
usage: Usage::default(), - 5132
model: self.config.model.clone(), - 5133
response_id: None, - 5134
}, - 5135
} - 5136
} - 5137
- 5138
/// The turn's answer when the model kept re-emitting a card it had - 5139
/// already shown: the card stands, the loop stops paying for acks. - 5140
async fn card_repeat_outcome(&self) -> TurnOutcome { - 5141
self.record_activity( - 5142
vak_session::ActivityKind::Diagnostic, - 5143
vak_session::ActivityStatus::Succeeded, - 5144
"card-repeat-exhausted".into(), - 5145
Some(format!( - 5146
"{CARD_REPEAT_EXHAUSTION_THRESHOLD} consecutive steps re-emitted an already-shown card; \ - 5147
the turn closed on that card as its answer" - 5148
)), - 5149
std::collections::BTreeMap::new(), - 5150
) - 5151
.await; - 5152
TurnOutcome::Completed { - 5153
response: AssistantMessage { - 5154
content: Vec::new(), - 5155
stop_reason: StopReason::EndTurn, - 5156
usage: Usage::default(), - 5157
model: self.config.model.clone(), - 5158
response_id: None, - 5159
}, - 5160
} - 5161
} - 5162
- 5163
async fn execute_batch_calls( - 5164
&self, - 5165
calls: Vec<PendingToolCall>, - 5166
cancel: &CancellationToken, - 5167
events: &mpsc::Sender<AgentEvent>, - 5168
) -> Vec<(String, ToolRunOutput)> { - 5169
let index = self - 5170
.config - 5171
.mcp_tool_index - 5172
.lock() - 5173
.unwrap_or_else(std::sync::PoisonError::into_inner) - 5174
.clone(); - 5175
let calls = calls - 5176
.into_iter() - 5177
.map(normalize_tool_call) - 5178
.map(|call| normalize_mcp_call(call, &index)) - 5179
.map(|call| normalize_schema_wrapper(call, &self.config.tools)) - 5180
.collect::<Vec<_>>(); - 5181
let n = calls.len(); - 5182
let cwd = self - 5183
.session - 5184
.lock() - 5185
.await - 5186
.header() - 5187
.map(|h| h.contract_cwd()) - 5188
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| ".".into())); - 5189
let sandbox = self.config.sandbox.clone(); - 5190
let hooks = self.config.hooks.clone(); - 5191
let hook_recorder = self.config.hook_recorder.clone(); - 5192
let tool_activity_recorder = self.config.tool_activity_recorder.clone(); - 5193
let session_id = self - 5194
.session - 5195
.lock() - 5196
.await - 5197
.header() - 5198
.map(|h| h.session_id.clone()) - 5199
.unwrap_or_default(); - 5200
let agent_id = self - 5201
.session - 5202
.lock() - 5203
.await - 5204
.header() - 5205
.and_then(|h| h.agent.as_ref().map(|a| a.id.clone())); - 5206
let skill_names = self - 5207
.session - 5208
.lock() - 5209
.await - 5210
.header() - 5211
.map(|h| { - 5212
h.contract - 5213
.capabilities - 5214
.iter() - 5215
.filter(|capability| { - 5216
capability.kind == vak_session::types::CapabilityKind::Skill - 5217
}) - 5218
.map(|capability| capability.name.clone()) - 5219
.collect::<Vec<_>>() - 5220
}) - 5221
.unwrap_or_default(); - 5222
- 5223
let mut authz: Vec<Result<(), String>> = Vec::with_capacity(n); - 5224
for call in &calls { - 5225
authz.push( - 5226
authorize( - 5227
&self.config, - 5228
call, - 5229
&cwd, - 5230
&self.run_call_counts, - 5231
&self.config.tools, - 5232
) - 5233
.await, - 5234
); - 5235
} - 5236
let ids: Vec<String> = calls.iter().map(|c| c.id.clone()).collect(); - 5237
- 5238
if !self.config.parallel_tools - 5239
|| n == 1 - 5240
|| self.config.work_mode == WorkMode::Managed - 5241
|| calls.iter().any(|call| call.name == "work") - 5242
|| batch_has_file_dependency(&calls) - 5243
{ - 5244
let mut out = Vec::with_capacity(n); - 5245
for (call, verdict) in calls.into_iter().zip(authz) { - 5246
match verdict { - 5247
Err(reason) => out.push((call.id, ToolRunOutput::Err(reason))), - 5248
Ok(()) => { - 5249
if cancel.is_cancelled() { - 5250
out.push((call.id, ToolRunOutput::Err("cancelled".into()))); - 5251
continue; - 5252
} - 5253
let managed_item = (self.config.work_mode == WorkMode::Managed - 5254
&& call.name != "work") - 5255
.then(|| call.name.clone()); - 5256
let managed_item = match managed_item { - 5257
Some(tool) => { - 5258
self.begin_managed_tool_item( - 5259
&tool, - 5260
(call.name == "task" || call.name == "flow") - 5261
.then(|| { - 5262
Some(( - 5263
call.input.get("contract_id")?.as_str()?, - 5264
call.input.get("work_item_id")?.as_str()?, - 5265
)) - 5266
}) - 5267
.flatten(), - 5268
call.input.get("flow").and_then(|value| value.as_str()), - 5269
) - 5270
.await - 5271
} - 5272
None => None, - 5273
}; - 5274
if self.config.work_mode == WorkMode::Managed - 5275
&& call.name != "work" - 5276
&& managed_item.is_none() - 5277
{ - 5278
out.push(( - 5279
call.id, - 5280
ToolRunOutput::Err( - 5281
"managed execution requires a compatible Ready or Running work item" - 5282
.into(), - 5283
), - 5284
)); - 5285
continue; - 5286
} - 5287
out.push(if call.name == "work" { - 5288
let result = self.execute_work_call(&call.input).await; - 5289
self.emit_work_state(events).await; - 5290
(call.id.clone(), result) - 5291
} else { - 5292
let owns_lifecycle = call.name == "task" || call.name == "flow"; - 5293
let (returned_id, result) = execute_one( - 5294
call, - 5295
&self.config.tools, - 5296
&cwd, - 5297
&session_id, - 5298
agent_id.as_deref(),
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.