- 3679
continue; - 3680
}; - 3681
match std::fs::read_to_string(&resolved) { - 3682
Ok(content) if content.contains(pattern) => { - 3683
vak_session::types::CriterionResult::Passed { - 3684
evidence: format!("file_contains:{}", path.display()), - 3685
} - 3686
} - 3687
Ok(_) => vak_session::types::CriterionResult::Failed { - 3688
reason: format!("pattern not found in {}", path.display()), - 3689
}, - 3690
Err(error) => vak_session::types::CriterionResult::Failed { - 3691
reason: format!("cannot read {}: {error}", path.display()), - 3692
}, - 3693
} - 3694
} - 3695
vak_session::types::CriterionKind::Semantic => continue, - 3696
vak_session::types::CriterionKind::ToolSucceeded { tool } => { - 3697
let mut succeeded = false; - 3698
for item_id in self.criterion_item_ids(projection, &criterion.criterion_id) { - 3699
if self - 3700
.tool_succeeded_for_item(projection, &item_id, tool) - 3701
.await - 3702
{ - 3703
succeeded = true; - 3704
break; - 3705
} - 3706
} - 3707
if succeeded { - 3708
vak_session::types::CriterionResult::Passed { - 3709
evidence: format!("tool_succeeded:{tool}"), - 3710
} - 3711
} else { - 3712
vak_session::types::CriterionResult::Unknown { - 3713
reason: format!("no successful '{tool}' tool result exists yet"), - 3714
} - 3715
} - 3716
} - 3717
vak_session::types::CriterionKind::FlowCompleted { flow } => { - 3718
if self.criterion_item_ids(projection, &criterion.criterion_id).into_iter().any( - 3719
|item_id| { - 3720
projection.items.get(&item_id).is_some_and(|item| { - 3721
item.evidence.iter().any(|evidence| { - 3722
matches!(evidence, vak_session::types::EvidenceRef::FlowNode { flow: name, node_id, .. } if name == flow && node_id == "__flow_completed__") - 3723
}) - 3724
}) - 3725
}, - 3726
) { - 3727
vak_session::types::CriterionResult::Passed { - 3728
evidence: format!("flow_completed:{flow}"), - 3729
} - 3730
} else { - 3731
vak_session::types::CriterionResult::Unknown { - 3732
reason: format!("flow '{flow}' has no linked completion evidence"), - 3733
} - 3734
} - 3735
} - 3736
vak_session::types::CriterionKind::ExternalReceipt { integration } => { - 3737
if self.criterion_item_ids(projection, &criterion.criterion_id).into_iter().any( - 3738
|item_id| { - 3739
projection.items.get(&item_id).is_some_and(|item| { - 3740
item.evidence.iter().any(|evidence| { - 3741
matches!(evidence, vak_session::types::EvidenceRef::ExternalOperation { integration: name, .. } if name == integration) - 3742
}) - 3743
}) - 3744
}, - 3745
) { - 3746
vak_session::types::CriterionResult::Passed { - 3747
evidence: format!("external_receipt:{integration}"), - 3748
} - 3749
} else { - 3750
vak_session::types::CriterionResult::Unknown { - 3751
reason: format!("integration '{integration}' has no receipt evidence"), - 3752
} - 3753
} - 3754
} - 3755
}; - 3756
self.record_managed_criterion(criterion, result).await; - 3757
} - 3758
let semantic: Vec<(vak_session::types::WorkCriterion, String)> = criteria - 3759
.iter() - 3760
.filter_map(|criterion| match &criterion.kind { - 3761
vak_session::types::CriterionKind::Semantic => Some(( - 3762
criterion.clone(), - 3763
format!("[{}] {}", criterion.criterion_id, criterion.statement), - 3764
)), - 3765
_ => None, - 3766
}) - 3767
.collect(); - 3768
if semantic.is_empty() || cancel.is_cancelled() { - 3769
return; - 3770
} - 3771
let judge_criteria: Vec<String> = semantic.iter().map(|(_, text)| text.clone()).collect(); - 3772
match self.run_judge(&judge_criteria, cancel, events).await { - 3773
Ok(verdicts) => { - 3774
for (criterion, expected) in semantic { - 3775
let result = verdicts - 3776
.iter() - 3777
.find(|verdict| verdict.criterion == expected) - 3778
.map(|verdict| match verdict.verdict.as_str() { - 3779
"pass" => vak_session::types::CriterionResult::Passed { - 3780
evidence: verdict.evidence.clone(), - 3781
}, - 3782
"fail" => vak_session::types::CriterionResult::Failed { - 3783
reason: verdict.evidence.clone(), - 3784
}, - 3785
_ => vak_session::types::CriterionResult::Unknown { - 3786
reason: verdict.evidence.clone(), - 3787
}, - 3788
}) - 3789
.unwrap_or_else(|| vak_session::types::CriterionResult::Unknown { - 3790
reason: "judge returned no verdict for this criterion".into(), - 3791
}); - 3792
self.record_managed_criterion(&criterion, result).await; - 3793
} - 3794
} - 3795
Err(reason) => { - 3796
for (criterion, _) in semantic { - 3797
self.record_managed_criterion( - 3798
&criterion, - 3799
vak_session::types::CriterionResult::Unknown { - 3800
reason: reason.clone(), - 3801
}, - 3802
) - 3803
.await; - 3804
} - 3805
} - 3806
} - 3807
} - 3808
- 3809
async fn record_managed_criterion( - 3810
&self, - 3811
criterion: &vak_session::types::WorkCriterion, - 3812
result: vak_session::types::CriterionResult, - 3813
) { - 3814
let contract_id = self - 3815
.session - 3816
.lock() - 3817
.await - 3818
.work_projection() - 3819
.ok() - 3820
.flatten() - 3821
.map(|projection| { - 3822
( - 3823
projection.contract.contract_id, - 3824
projection.contract.revision, - 3825
) - 3826
}); - 3827
if let Some((contract_id, revision)) = contract_id { - 3828
let _ = self - 3829
.session - 3830
.lock() - 3831
.await - 3832
.append_work(vak_session::types::WorkEvent { - 3833
contract_id, - 3834
revision, - 3835
kind: vak_session::types::WorkEventKind::VerificationRecorded { - 3836
criterion_id: criterion.criterion_id.clone(), - 3837
result, - 3838
}, - 3839
}); - 3840
} - 3841
} - 3842
- 3843
fn criterion_item_ids( - 3844
&self, - 3845
projection: &vak_session::work::WorkProjection, - 3846
criterion_id: &str, - 3847
) -> Vec<String> { - 3848
projection - 3849
.contract - 3850
.items - 3851
.iter() - 3852
.filter(|item| item.criterion_ids.iter().any(|id| id == criterion_id)) - 3853
.map(|item| item.item_id.clone()) - 3854
.collect() - 3855
} - 3856
- 3857
async fn tool_succeeded_for_item( - 3858
&self, - 3859
projection: &vak_session::work::WorkProjection, - 3860
item_id: &str, - 3861
tool: &str, - 3862
) -> bool { - 3863
let Some(item) = projection.items.get(item_id) else { - 3864
return false; - 3865
}; - 3866
let tool_result_ids: std::collections::HashSet<&str> = item - 3867
.evidence - 3868
.iter() - 3869
.filter_map(|evidence| match evidence { - 3870
vak_session::types::EvidenceRef::ToolResult { tool_use_id, .. } => { - 3871
Some(tool_use_id.as_str()) - 3872
} - 3873
_ => None, - 3874
}) - 3875
.collect(); - 3876
if tool_result_ids.is_empty() { - 3877
return false; - 3878
} - 3879
let session = self.session.lock().await; - 3880
let mut tool_names = std::collections::HashMap::new(); - 3881
for entry in session.chain_to_root() { - 3882
if let vak_session::types::EntryPayload::Message(record) = &entry.payload { - 3883
for block in &record.message.content { - 3884
if let vak_llm::ContentBlock::ToolUse { id, name, .. } = block { - 3885
tool_names.insert(id.as_str(), name.as_str()); - 3886
} - 3887
} - 3888
} - 3889
} - 3890
session.chain_to_root().iter().any(|entry| { - 3891
let vak_session::types::EntryPayload::Message(record) = &entry.payload else { - 3892
return false; - 3893
}; - 3894
record.message.content.iter().any(|block| { - 3895
matches!( - 3896
block, - 3897
vak_llm::ContentBlock::ToolResult { - 3898
tool_use_id, - 3899
is_error: false, - 3900
.. - 3901
} if tool_result_ids.contains(tool_use_id.as_str()) - 3902
&& tool_names.get(tool_use_id.as_str()).copied() == Some(tool) - 3903
) - 3904
}) - 3905
}) - 3906
} - 3907
- 3908
/// Goal audit gate (Phase H): runs when the model claims completion. - 3909
/// Some(reason) rejects the claim and continues the run; None lets it - 3910
/// end. Completion is recorded from audit, never self-report — and the - 3911
/// audit budget is capped so this can never trap a run. - 3912
async fn goal_gate( - 3913
&mut self, - 3914
_response: &AssistantMessage, - 3915
cancel: &CancellationToken, - 3916
events: &mpsc::Sender<AgentEvent>, - 3917
) -> Option<String> { - 3918
let goal_active = self.active_goal.is_some(); - 3919
if !goal_active { - 3920
return None; - 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 {
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.