- 6638
surface.core.iter().map(|def| def.name.clone()).collect(); - 6639
let deferred_tool_names: Vec<String> = surface - 6640
.deferred - 6641
.iter() - 6642
.map(|def| def.name.clone()) - 6643
.collect(); - 6644
let unpredicted_tools = surface.unpredicted.clone(); - 6645
let mut defs: Vec<vak_llm::ToolDefinition> = - 6646
Vec::with_capacity(surface.core.len() + surface.deferred.len() + 1); - 6647
defs.push(find_tools_def); - 6648
defs.extend(surface.core); - 6649
defs.extend(surface.deferred.into_iter().map(|def| def.deferred())); - 6650
cfg.tool_definitions = Some(defs); - 6651
cfg.tools = tools; - 6652
if cfg.work_mode == WorkMode::Managed && turn_capabilities.flow_admitted { - 6653
cfg.flow_dispatcher = Some(Arc::new(CoreFlowDispatcher { - 6654
core: self.clone(), - 6655
tools: cfg.tools.clone(), - 6656
system_prompt: cfg.system_prefix.clone(), - 6657
})); - 6658
} - 6659
let selected_ids: std::collections::BTreeSet<String> = turn_capabilities - 6660
.descriptors - 6661
.iter() - 6662
.map(|descriptor| format!("{:?}:{}", descriptor.kind, descriptor.name)) - 6663
.collect(); - 6664
let all_ids: std::collections::BTreeSet<String> = cap_set - 6665
.all() - 6666
.map(|capability| format!("{:?}:{}", capability.id.kind, capability.id.name)) - 6667
.collect(); - 6668
let tool_schemas = cfg - 6669
.tool_definitions - 6670
.as_ref() - 6671
.map(|definitions| { - 6672
definitions - 6673
.iter() - 6674
.filter_map(|definition| serde_json::to_value(definition).ok()) - 6675
.collect() - 6676
}) - 6677
.unwrap_or_default(); - 6678
// Per-tool declared domains, so a later projection can derive - 6679
// delivery signals from what the capability declared it serves - 6680
// rather than from its name (docs/design/68 §9's - 6681
// `SignalContext.domains` note). Covers MCP servers and their - 6682
// discovered tools too, not just built-ins. - 6683
let mut tool_domains: std::collections::BTreeMap<String, Vec<String>> = - 6684
std::collections::BTreeMap::new(); - 6685
for capability in cap_set.all() { - 6686
let labels = capability.serves.labels(); - 6687
if labels.is_empty() { - 6688
continue; - 6689
} - 6690
tool_domains.insert(capability.id.name.clone(), labels.clone()); - 6691
if capability.id.kind == CapabilityKind::McpServer - 6692
&& let Some(inventory) = capability - 6693
.configuration - 6694
.get("tools") - 6695
.and_then(|t| t.as_array()) - 6696
{ - 6697
for tool in inventory { - 6698
if let Some(name) = tool.get("name").and_then(|n| n.as_str()) { - 6699
tool_domains - 6700
.entry(name.to_string()) - 6701
.or_insert_with(|| labels.clone()); - 6702
} - 6703
} - 6704
} - 6705
} - 6706
if let Err(error) = - 6707
session.append_turn_capabilities(vak_session::types::TurnCapabilitiesBound { - 6708
epoch: cap_set.epoch, - 6709
capability_ids: selected_ids.iter().cloned().collect(), - 6710
excluded_ids: all_ids.difference(&selected_ids).cloned().collect(), - 6711
system_prompt: cfg.system_prefix.clone(), - 6712
tool_schemas, - 6713
core_tool_names, - 6714
deferred_tool_names, - 6715
tool_index: tool_catalogue, - 6716
tool_domains, - 6717
}) - 6718
{ - 6719
return Err(CoreError::Session(error)); - 6720
} - 6721
// Hooks come from TurnCapabilities: the same admission as every - 6722
// other kind, read by the shared `hook_def` reader. - 6723
let hooks: std::sync::Arc<Vec<vak_hooks::HookDef>> = - 6724
std::sync::Arc::new(turn_capabilities.hooks); - 6725
cfg.hooks = Some(hooks.clone()); - 6726
let shared_root = self.shared_capability_root(); - 6727
let plugin_hooks: Vec<_> = self - 6728
.capability_roots() - 6729
.into_iter() - 6730
.flat_map(|root| { - 6731
vak_plugin::PluginStore::new(root.path) - 6732
.enabled_hooks() - 6733
.unwrap_or_default() - 6734
}) - 6735
.collect(); - 6736
let shared_home = shared_root; - 6737
let workspace_home = self.inner.cwd.join(".vak"); - 6738
let activity_ledger = finops::ActivityLedger::new(&self.sessions_home()); - 6739
let activity_session = session.header().map(|h| h.session_id.clone()); - 6740
cfg.hook_recorder = Some(Arc::new( - 6741
move |hook: &vak_hooks::HookDef, success: bool, duration_ms: u64| { - 6742
let plugin_name = plugin_hooks - 6743
.iter() - 6744
.find(|(_, candidate)| candidate.command == hook.command) - 6745
.map(|(plugin, _)| plugin.name.clone()); - 6746
let _ = activity_ledger.append(&finops::ActivityRow { - 6747
ts: chrono::Utc::now(), - 6748
kind: "hook".into(), - 6749
name: hook.event.as_str().into(), - 6750
success, - 6751
duration_ms: Some(duration_ms), - 6752
session_id: activity_session.clone(), - 6753
plugin: plugin_name, - 6754
}); - 6755
if let Some((plugin, _)) = plugin_hooks - 6756
.iter() - 6757
.find(|(_, candidate)| candidate.command == hook.command) - 6758
{ - 6759
let store = vak_plugin::PluginStore::new(match plugin.scope { - 6760
vak_plugin::InstallScope::User => &shared_home, - 6761
vak_plugin::InstallScope::Workspace => &workspace_home, - 6762
}); - 6763
let _ = store.record_invocation( - 6764
&plugin.trace_id, - 6765
&plugin.name, - 6766
&format!("hook:{}", hook.event.as_str()), - 6767
success, - 6768
); - 6769
} - 6770
}, - 6771
)); - 6772
let tool_activity_ledger = finops::ActivityLedger::new(&self.sessions_home()); - 6773
let tool_activity_session = session.header().map(|h| h.session_id.clone()); - 6774
cfg.tool_activity_recorder = Some(Arc::new( - 6775
move |name: &str, args: &serde_json::Value, success: bool, duration_ms: u64| { - 6776
let plugin = name - 6777
.strip_prefix("plugin.") - 6778
.and_then(|rest| rest.split('.').next()) - 6779
.map(str::to_owned); - 6780
let activity_name = if name == "skill" { - 6781
args.get("name") - 6782
.and_then(serde_json::Value::as_str) - 6783
.map_or_else( - 6784
|| "skill/(unknown)".into(), - 6785
|skill| format!("skill/{skill}"), - 6786
) - 6787
} else { - 6788
name.to_owned() - 6789
}; - 6790
let _ = tool_activity_ledger.append(&finops::ActivityRow { - 6791
ts: chrono::Utc::now(), - 6792
kind: if name == "skill" { "skill" } else { "tool" }.into(), - 6793
name: activity_name, - 6794
success, - 6795
duration_ms: Some(duration_ms), - 6796
session_id: tool_activity_session.clone(), - 6797
plugin, - 6798
}); - 6799
}, - 6800
)); - 6801
- 6802
// session-start hooks fire once per run, before any tool or - 6803
// checkpoint activity. A block aborts the run before it starts. - 6804
if hooks - 6805
.iter() - 6806
.any(|h| h.event == vak_hooks::HookEvent::SessionStart) - 6807
{ - 6808
let session_id = session - 6809
.header() - 6810
.map(|h| h.session_id.clone()) - 6811
.unwrap_or_default(); - 6812
let outcome = vak_hooks::run_hooks_with_recorder( - 6813
hooks.clone(), - 6814
vak_hooks::HookEvent::SessionStart, - 6815
&session_id, - 6816
&self.inner.cwd, - 6817
None, - 6818
None, - 6819
&cancel, - 6820
cfg.hook_recorder.as_deref(), - 6821
) - 6822
.await; - 6823
if outcome.blocked { - 6824
let reason = outcome.reason.unwrap_or_else(|| "blocked by hook".into()); - 6825
return Err(CoreError::HookBlocked(format!("session-start: {reason}"))); - 6826
} - 6827
} - 6828
- 6829
// Checkpoint the workspace before any mutation of this run. Skipped - 6830
// only when the engagement is confident nothing will be executed or - 6831
// written (a greeting, a question): a wrong reading there costs a - 6832
// missed checkpoint, so the skip needs the acceptance bar, not the - 6833
// provisional one. - 6834
let expects_effect = engagement.posture.checkpoint_before_effect - 6835
|| resolved_intent.provenance.tier == vak_intent::Tier::General - 6836
|| !resolved_intent - 6837
.reading - 6838
.may_slice_capabilities(self.inner.config.intent.accept_confidence) - 6839
|| admitted_outcome.requires_execution(); - 6840
if let Some(h) = session.header() { - 6841
let seq = self.next_checkpoint_seq(&h.session_id); - 6842
// The session's first checkpoint is always taken: it is the - 6843
// baseline "what has this session changed" is measured against - 6844
// (`ContextProfile::Working`), whatever the first turn was. - 6845
let first_of_session = seq == 0; - 6846
if !expects_effect && !first_of_session { - 6847
// Nothing will be executed or written; skip the capture. - 6848
} else { - 6849
if let Ok((cp, _stats)) = checkpoints::capture( - 6850
&self.inner.cwd, - 6851
&self.sessions_home(), - 6852
&h.session_id, - 6853
seq, - 6854
&format!("turn: {}", prompt.text_content()), - 6855
) { - 6856
let _ = checkpoints::store(&self.sessions_home(), &cp); - 6857
} - 6858
} - 6859
} - 6860
- 6861
// Record the intent before the turn dispatches. The entry carries the - 6862
// exact note the engagement contributes, so the projection the model - 6863
// sees comes from the ledger rather than from a derivation that might - 6864
// read differently on replay (invariant 1: model-visible means - 6865
// logged). A write failure is not fatal — losing the audit row must - 6866
// not lose the user's turn — but it does mean the note does not reach - 6867
// the model either, because both come from the same entry. - 6868
let mut session = session; - 6869
- 6870
// Only an explicit command corrects or replaces the goal; ordinary - 6871
// text adds to it (docs/design/47, control plane). - 6872
let goal_update = session.next_goal_update(&prompt.text_content()); - 6873
if let Err(error) = session.append_goal_update(goal_update) { - 6874
eprintln!("[goal] could not record this request relationship: {error}"); - 6875
} - 6876
- 6877
// A request restated verbatim right after the previous turn is the - 6878
// user saying the previous reading did the wrong thing. That counts - 6879
// against the *previous* reading (misread ledger, I8), not this one. - 6880
if self.inner.config.intent.enabled { - 6881
let chain = session.chain_to_root(); - 6882
// A person's message, whatever metadata it carries (an attachment - 6883
// is metadata); only a runtime nudge is skipped, by its tag. - 6884
let previous_user_text = chain.iter().rev().find_map(|entry| match &entry.payload { - 6885
vak_session::EntryPayload::Message(record) - 6886
if record.message.role == vak_llm::Role::User - 6887
&& record.control_kind().is_none() => - 6888
{ - 6889
Some(record.message.text_content()) - 6890
} - 6891
_ => None, - 6892
}); - 6893
let previous_intent = chain.iter().rev().find_map(|entry| match &entry.payload { - 6894
vak_session::EntryPayload::Intent(record) => Some(record.as_ref().clone()), - 6895
_ => None, - 6896
}); - 6897
let same = |a: &str, b: &str| { - 6898
let norm = |t: &str| { - 6899
t.split_whitespace() - 6900
.collect::<Vec<_>>() - 6901
.join(" ") - 6902
.to_ascii_lowercase() - 6903
}; - 6904
!a.trim().is_empty() && norm(a) == norm(b) - 6905
}; - 6906
if let (Some(previous_text), Some(previous)) = (previous_user_text, previous_intent) - 6907
&& same(&previous_text, &prompt_text) - 6908
&& previous.provenance.tier != vak_intent::Tier::General - 6909
{ - 6910
misread::MisreadLedger::new(&self.sessions_home()).record( - 6911
&previous.reading, - 6912
previous.provenance.tier, - 6913
previous.provenance.resolver_version, - 6914
misread::Outcome::Restated, - 6915
None, - 6916
self.reading_sliced(&previous.reading), - 6917
); - 6918
} - 6919
} - 6920
- 6921
// Durable work earns a commitment of its own before the turn runs, so - 6922
// the episode brackets the work rather than being reconstructed from - 6923
// it afterwards. A ledger failure is logged and dropped: losing the - 6924
// audit row must never cost the user their turn. - 6925
let episodes = session - 6926
.header() - 6927
.map(|header| { - 6928
commitments::begin_episodes( - 6929
&self.sessions_home(), - 6930
&self.inner.config, - 6931
&resolved_intent, - 6932
&episode_plan, - 6933
&prompt.text_content(), - 6934
&header.session_id, - 6935
&self.inner.cwd, - 6936
header - 6937
.conversation - 6938
.as_ref() - 6939
.map(|conversation| conversation.audience_id.as_str()), - 6940
) - 6941
}) - 6942
.unwrap_or_default(); - 6943
// The turn's primary commitment: the first durable strand's. - 6944
let episode = episodes.first().cloned(); - 6945
- 6946
// A live grant pre-authorizes the actions it covers, one gate at a - 6947
// time — only under delegation, and never for a turn with anything - 6948
// irreversible in it, which reaches a human whatever was delegated. - 6949
let enveloped = episode_plan.enveloped_commitments(); - 6950
let irreversible = resolved_intent.reading.stakes == vak_intent::Stakes::Irreversible - 6951
|| resolved_intent - 6952
.strands - 6953
.iter() - 6954
.any(|strand| strand.reading.stakes == vak_intent::Stakes::Irreversible); - 6955
if turn_authority.autonomy == vak_intent::Autonomy::Delegated - 6956
&& !irreversible - 6957
&& !enveloped.is_empty() - 6958
{ - 6959
cfg.envelope_check = Some(intent::envelope_check( - 6960
self.sessions_home(), - 6961
enveloped, - 6962
self.inner.cwd.clone(), - 6963
)); - 6964
} - 6965
- 6966
if self.inner.config.intent.enabled { - 6967
// `ContextProfile::Full`: durable work sees its obligations - 6968
// rendered from the commitment ledger, appended to the intent - 6969
// note so the ledger row carries exactly what the model saw. - 6970
let model_visible = match ( - 6971
engagement.posture.context, - 6972
commitments::prompt_projection(&self.sessions_home(), &episodes), - 6973
) { - 6974
(vak_intent::ContextProfile::Full, Some(projection)) => { - 6975
Some(match resolved_intent.model_visible() { - 6976
Some(note) => format!("{note}\n{projection}"), - 6977
None => projection, - 6978
}) - 6979
} - 6980
_ => resolved_intent.model_visible(), - 6981
}; - 6982
let record = vak_session::types::IntentRecord { - 6983
reading: resolved_intent.reading.clone(), - 6984
strands: resolved_intent.strands.clone(), - 6985
engagement: resolved_intent.engagement.clone(), - 6986
provenance: resolved_intent.provenance.clone(), - 6987
outcome: Some(admitted_outcome.clone()), - 6988
model_visible, - 6989
commitment_id: episode - 6990
.as_ref() - 6991
.map(|episode| episode.commitment_id.clone()), - 6992
strand_commitments: episodes - 6993
.iter() - 6994
.map(|episode| (episode.strand_id.clone(), episode.commitment_id.clone())) - 6995
.collect(), - 6996
}; - 6997
if let Err(error) = session.append_intent(record) { - 6998
return Err(CoreError::Session(error)); - 6999
} - 7000
- 7001
// `ContextProfile::Working` / `Full`: what this session has - 7002
// changed in the workspace so far, rendered into the tail. It - 7003
// is a filesystem observation, so the bytes go into the ledger - 7004
// as an activity first (model-visible means logged) and the - 7005
// tail reads them from there. Measured against the session's - 7006
// first checkpoint; a first turn has nothing to compare. - 7007
if matches!( - 7008
engagement.posture.context, - 7009
vak_intent::ContextProfile::Working | vak_intent::ContextProfile::Full - 7010
) && let Some(header) = session.header() - 7011
&& let Ok(list) = checkpoints::list(&self.sessions_home(), &header.session_id) - 7012
&& let Some(first) = list.iter().map(|cp| cp.seq).min() - 7013
&& let Ok(delta) = checkpoints::delta_summary( - 7014
&self.inner.cwd, - 7015
&self.sessions_home(), - 7016
&header.session_id, - 7017
first, - 7018
8_192, - 7019
) - 7020
&& !delta.contains("workspace unchanged since checkpoint") - 7021
{ - 7022
let _ = session.append_activity(vak_session::ActivityRecord { - 7023
activity_id: format!("workspace-delta-{}", uuid_like()), - 7024
turn: None, - 7025
kind: vak_session::ActivityKind::Diagnostic, - 7026
status: vak_session::ActivityStatus::Succeeded, - 7027
label: "Workspace changes since the session began".into(), - 7028
detail: Some(delta), - 7029
data: std::collections::BTreeMap::from([( - 7030
"section".to_string(), - 7031
SessionLog::WORKSPACE_DELTA_SECTION.to_string(), - 7032
)]), - 7033
}); - 7034
} - 7035
} - 7036
- 7037
// `Defer`: a gate nobody here can answer is parked in the inbox and - 7038
// suspends the commitment instead of merely failing the run. - 7039
if engagement.posture.gate_fallback == vak_intent::GateFallback::Defer - 7040
&& let Some(episode) = &episode - 7041
{ - 7042
let escalation = envelopes - 7043
.get(&episode.strand_id) - 7044
.map(|envelope| envelope.escalation.clone()) - 7045
.unwrap_or_default(); - 7046
cfg.approver = Some(std::sync::Arc::new(intent::DeferringApprover::new( - 7047
cfg.approver.clone(), - 7048
self.shared_data_home(), - 7049
self.sessions_home(), - 7050
sid.clone(), - 7051
episode.commitment_id.clone(), - 7052
escalation, - 7053
))); - 7054
} - 7055
- 7056
let steering = match steering { - 7057
Some(s) => s, - 7058
None => std::sync::Arc::new(vak_agent::SteeringQueues::new()), - 7059
}; - 7060
let admitted_outcome = cfg.outcome.clone(); - 7061
let mut agent = Agent::new(provider, session, cfg); - 7062
if let Some((objective, criteria)) = goal { - 7063
agent.set_goal(objective, criteria); - 7064
} - 7065
let (receipts_before, entries_before) = { - 7066
let s = agent.session.lock().await; - 7067
(s.receipts().len(), s.chain_to_root().len()) - 7068
}; - 7069
let run_cancel = cancel.child_token(); - 7070
let outcome = { - 7071
let run = agent.run_message( - 7072
vak_session::MessageRecord { - 7073
message: prompt, - 7074
meta: prompt_meta, - 7075
}, - 7076
&steering, - 7077
run_cancel.clone(), - 7078
events, - 7079
); - 7080
tokio::pin!(run); - 7081
tokio::select! { - 7082
biased; - 7083
_ = permission_lease.cancelled() => { - 7084
run_cancel.cancel(); - 7085
run.await - 7086
} - 7087
outcome = &mut run => outcome, - 7088
} - 7089
}; - 7090
let mut session = agent.into_session().await; - 7091
- 7092
// Record what the runtime actually produced separately from the - 7093
// earlier intent record. The outcome contract is append-only: a - 7094
// response may exist without satisfying its evidence requirements. - 7095
if self.inner.config.intent.enabled { - 7096
// The answer is its presentations plus its narration - 7097
// (docs/design/68-context-engine.md §10): a turn that emitted a - 7098
// card and no prose still delivered, so the evaluator sees the - 7099
// cards' rendered text alongside whatever text the model wrote. - 7100
let presented_text = { - 7101
let turn_id = session.latest_directive_entry_id(); - 7102
session - 7103
.presentations() - 7104
.into_iter() - 7105
.filter(|(_, record)| Some(record.turn_id.as_str()) == turn_id.as_deref()) - 7106
.map(|(_, record)| { - 7107
format!( - 7108
"{{\"semantic_type\":\"{}\",\"payload\":{}}}", - 7109
record.semantic_type, record.payload - 7110
) - 7111
}) - 7112
.collect::<Vec<_>>() - 7113
.join("\n") - 7114
}; - 7115
let response_text = match &outcome { - 7116
TurnOutcome::Completed { response } => Some(response.text_content()), - 7117
TurnOutcome::Aborted { partial } => { - 7118
partial.as_ref().map(|message| message.text_content()) - 7119
} - 7120
TurnOutcome::Failed { .. } | TurnOutcome::MaxTurnsReached => None, - 7121
} - 7122
.map(|text| { - 7123
if presented_text.is_empty() { - 7124
text - 7125
} else if text.trim().is_empty() { - 7126
presented_text.clone() - 7127
} else { - 7128
format!("{presented_text}\n{text}") - 7129
} - 7130
}); - 7131
let turn = session - 7132
.chain_to_root() - 7133
.iter() - 7134
.filter(|entry| { - 7135
matches!( - 7136
&entry.payload, - 7137
vak_session::types::EntryPayload::Message(record) - 7138
if record.message.role == vak_llm::Role::User - 7139
&& record.control_kind().is_none() - 7140
&& record.message.content.iter().any(|block| { - 7141
matches!(block, vak_llm::ContentBlock::Text { .. }) - 7142
}) - 7143
) - 7144
}) - 7145
.count(); - 7146
let mut outcome_spec = admitted_outcome.unwrap_or_else(|| { - 7147
vak_intent::OutcomeSpec::from_reading( - 7148
prompt_text, - 7149
&resolved_intent.reading, - 7150
resolved_intent.provenance.resolver_version, - 7151
) - 7152
}); - 7153
outcome_spec.evidence_max_age_secs = - 7154
Some(self.inner.config.intent.evidence_max_age_secs); - 7155
let mut tool_calls = - 7156
std::collections::HashMap::<String, (Option<String>, Option<String>)>::new(); - 7157
let mut successful_receipts = std::collections::HashSet::new(); - 7158
let mut written_paths = std::collections::HashSet::<String>::new(); - 7159
let mut successful_effect_inputs = Vec::<String>::new(); - 7160
let mut failed_correctable = std::collections::HashSet::new(); - 7161
for entry in session.chain_to_root() { - 7162
if let vak_session::EntryPayload::Message(record) = &entry.payload { - 7163
if record.message.role == vak_llm::Role::User - 7164
&& record.control_kind().is_none() - 7165
&& record - 7166
.message - 7167
.content - 7168
.iter() - 7169
.any(|block| matches!(block, vak_llm::ContentBlock::Text { .. })) - 7170
{ - 7171
tool_calls.clear(); - 7172
successful_receipts.clear(); - 7173
written_paths.clear(); - 7174
successful_effect_inputs.clear(); - 7175
failed_correctable.clear(); - 7176
} - 7177
for block in &record.message.content { - 7178
match block { - 7179
vak_llm::ContentBlock::ToolUse { id, name, input } => { - 7180
let canonical = vak_tools::canonical_tool_name(name); - 7181
let path = matches!(canonical, "write" | "edit") - 7182
.then(|| input.get("path").and_then(serde_json::Value::as_str)) - 7183
.flatten() - 7184
.map(str::to_ascii_lowercase); - 7185
let effect_input = matches!( - 7186
canonical, - 7187
"write" | "edit" | "apply_patch" | "bash" | "imagegen" - 7188
) - 7189
.then(|| input.to_string().to_ascii_lowercase()); - 7190
tool_calls.insert(id.clone(), (path, effect_input)); - 7191
} - 7192
vak_llm::ContentBlock::ToolResult { - 7193
tool_use_id, - 7194
is_error: false, - 7195
.. - 7196
} if tool_calls.contains_key(tool_use_id) => { - 7197
successful_receipts.insert(tool_use_id.clone()); - 7198
if let Some((Some(path), _)) = tool_calls.get(tool_use_id) { - 7199
written_paths.insert(path.clone()); - 7200
} - 7201
if let Some((_, Some(input))) = tool_calls.get(tool_use_id) { - 7202
successful_effect_inputs.push(input.clone()); - 7203
} - 7204
} - 7205
vak_llm::ContentBlock::ToolResult { - 7206
tool_use_id, - 7207
is_error: true, - 7208
content, - 7209
.. - 7210
} if tool_calls.contains_key(tool_use_id) - 7211
&& vak_tools::ToolErrorKind::classify(content).is_correctable() => - 7212
{ - 7213
failed_correctable.insert(tool_use_id.clone()); - 7214
} - 7215
_ => {} - 7216
} - 7217
} - 7218
} - 7219
} - 7220
let evidence_state = session - 7221
.successful_tool_receipts_for_latest_turn() - 7222
.last() - 7223
.map_or(vak_intent::EvidenceState::None, |(_, recorded_at)| { - 7224
vak_intent::evidence_state_from_age( - 7225
chrono::Utc::now(), - 7226
*recorded_at, - 7227
chrono::Duration::seconds( - 7228
outcome_spec.evidence_max_age_secs.unwrap_or(86_400), - 7229
), - 7230
) - 7231
}); - 7232
let unresolved_correctable = - 7233
!failed_correctable.is_empty() && successful_receipts.is_empty(); - 7234
let mut status = vak_intent::evaluate_response_with_failures( - 7235
response_text.as_deref(), - 7236
matches!( - 7237
outcome, - 7238
TurnOutcome::Failed { .. } | TurnOutcome::MaxTurnsReached - 7239
), - 7240
matches!(outcome, TurnOutcome::Aborted { .. }), - 7241
unresolved_correctable, - 7242
); - 7243
// A named saved-file request needs an observed tool result. A - 7244
// model's sentence saying it wrote the file is not a deliverable. - 7245
if status == vak_intent::OutcomeStatus::Produced - 7246
&& let Some(target) = outcome_spec.saved_file_target() - 7247
&& !written_paths.iter().any(|path| { - 7248
std::path::Path::new(path).file_name() - 7249
== std::path::Path::new(&target).file_name() - 7250
}) - 7251
&& !successful_effect_inputs - 7252
.iter() - 7253
.any(|input| input.contains(&target)) - 7254
{ - 7255
status = vak_intent::OutcomeStatus::Unknown; - 7256
} - 7257
let requirement_evaluations = vak_intent::evaluate_requirements_with_state( - 7258
&outcome_spec, - 7259
response_text.as_deref(), - 7260
evidence_state, - 7261
); - 7262
let completion = - 7263
vak_intent::evaluate_completion(status, &requirement_evaluations, &outcome_spec); - 7264
let human_review = vak_intent::human_review_state(completion); - 7265
let evidence_receipts = successful_receipts - 7266
.into_iter() - 7267
.collect::<Vec<_>>() - 7268
.join(","); - 7269
let _ = session.append_activity(vak_session::types::ActivityRecord { - 7270
activity_id: format!("outcome-evaluation-{}", uuid_like()), - 7271
turn: Some(turn), - 7272
kind: vak_session::types::ActivityKind::Diagnostic, - 7273
status: vak_session::types::ActivityStatus::Succeeded, - 7274
label: "Outcome evaluation".into(), - 7275
detail: Some(format!("primary deliverable: {status:?}")), - 7276
data: std::collections::BTreeMap::from([ - 7277
("status".into(), format!("{status:?}").to_ascii_lowercase()), - 7278
( - 7279
"completion".into(), - 7280
format!("{completion:?}").to_ascii_lowercase(), - 7281
), - 7282
("evidence_receipts".into(), evidence_receipts), - 7283
( - 7284
"evidence_state".into(), - 7285
format!("{evidence_state:?}").to_ascii_lowercase(), - 7286
), - 7287
("human_review".into(), human_review.into()), - 7288
( - 7289
"evidence_max_age_secs".into(), - 7290
outcome_spec - 7291
.evidence_max_age_secs - 7292
.map_or_else(|| "none".into(), |value| value.to_string()), - 7293
), - 7294
( - 7295
"requirements".into(), - 7296
outcome_spec - 7297
.requirements - 7298
.iter() - 7299
.map(|requirement| requirement.id.as_str()) - 7300
.collect::<Vec<_>>() - 7301
.join(","), - 7302
), - 7303
( - 7304
"evaluation".into(), - 7305
serde_json::to_string(&requirement_evaluations) - 7306
.unwrap_or_else(|_| "[]".into()), - 7307
), - 7308
]), - 7309
}); - 7310
} - 7311
- 7312
// Was the reading right? The strongest answer is measured, not - 7313
// guessed: if the engagement withheld a tool and the model then asked - 7314
// for that exact tool, the reading was wrong and we know which lexicon - 7315
// entry to change. Slicing is what makes this observable at all. - 7316
if self.inner.config.intent.enabled - 7317
&& resolved_intent.provenance.tier != vak_intent::Tier::General - 7318
{ - 7319
// Only this turn's own tool calls count. Scanning the whole - 7320
// chain recorded a tool used three turns ago as an escalation - 7321
// against today's reading, and biased every cell downward with - 7322
// session length. - 7323
let attempted: Vec<String> = session - 7324
.chain_to_root() - 7325
.iter() - 7326
.skip(entries_before) - 7327
.flat_map(|entry| match &entry.payload { - 7328
vak_session::EntryPayload::Message(record) => record - 7329
.message - 7330
.content - 7331
.iter() - 7332
.filter_map(|block| match block { - 7333
vak_llm::ContentBlock::ToolUse { name, .. } => Some(name.clone()), - 7334
_ => None, - 7335
}) - 7336
.collect::<Vec<_>>(), - 7337
_ => Vec::new(), - 7338
}) - 7339
.collect(); - 7340
let wanted = misread::escalated_capability(&unpredicted_tools, &attempted); - 7341
let outcome = match (&wanted, &outcome) { - 7342
(Some(_), _) => misread::Outcome::Escalated, - 7343
(None, TurnOutcome::Aborted { .. }) => misread::Outcome::Abandoned, - 7344
_ => misread::Outcome::Held, - 7345
}; - 7346
misread::MisreadLedger::new(&self.sessions_home()).record( - 7347
&resolved_intent.reading, - 7348
resolved_intent.provenance.tier, - 7349
resolved_intent.provenance.resolver_version, - 7350
outcome, - 7351
wanted, - 7352
self.reading_sliced(&resolved_intent.reading), - 7353
); - 7354
} - 7355
- 7356
// Close the episode with what it actually achieved. `Learned` and - 7357
// `Stalled` are deliberately different: a turn that answered - 7358
// substantively but moved no criterion reduced uncertainty and must - 7359
// not count against the stall breaker. - 7360
if !episodes.is_empty() { - 7361
let tool_calls = session - 7362
.chain_to_root() - 7363
.iter() - 7364
.filter(|entry| match &entry.payload { - 7365
vak_session::EntryPayload::Message(record) => record - 7366
.message - 7367
.content - 7368
.iter() - 7369
.any(|block| matches!(block, vak_llm::ContentBlock::ToolUse { .. })), - 7370
_ => false, - 7371
}) - 7372
.count(); - 7373
// Reuses the shared estimator rather than multiplying tokens by a - 7374
// rate here: it already handles the cache-creation and cache-read - 7375
// tiers, and a second cost formula would drift from the ledger's. - 7376
let prices = &self.inner.config.finops.price_overrides; - 7377
let spend = session - 7378
.receipts() - 7379
.iter() - 7380
.skip(receipts_before) - 7381
.flat_map(|receipt| { - 7382
let model = receipt.model.clone(); - 7383
receipt.attempts.iter().filter_map(move |attempt| { - 7384
attempt.usage.as_ref().map(|usage| (model.clone(), usage)) - 7385
}) - 7386
}) - 7387
.filter_map(|(model, usage)| { - 7388
vak_config::finops::estimate_cost_usd(&model, usage, prices) - 7389
}) - 7390
.sum::<f64>(); - 7391
// Spend is attributed to the primary strand's commitment; the - 7392
// others record the advancement at zero cost rather than - 7393
// double-counting one turn's dispatches. - 7394
for (index, episode) in episodes.iter().enumerate() { - 7395
commitments::end_episode( - 7396
&self.sessions_home(), - 7397
episode, - 7398
commitments::classify(&outcome, tool_calls, Vec::new()), - 7399
if index == 0 { spend } else { 0.0 }, - 7400
); - 7401
} - 7402
} - 7403
- 7404
// Phase B: fold this run's dispatches into the routing evidence - 7405
// ledger (success / failure / unknown by settlement). - 7406
let new_receipts: Vec<vak_llm::WorkReceipt> = session - 7407
.receipts() - 7408
.into_iter() - 7409
.skip(receipts_before) - 7410
.cloned() - 7411
.collect(); - 7412
if !new_receipts.is_empty() { - 7413
routing::EvidenceLedger::new(&self.sessions_home()).record_receipts(&new_receipts); - 7414
// Phase R: fold the same dispatches into session beliefs. - 7415
// Domain-weighted doubt accumulates per leg; one success - 7416
// clears it. Cancelled attempts say nothing. - 7417
for r in &new_receipts { - 7418
for a in &r.attempts { - 7419
let (provider, model) = r.attempt_leg(a); - 7420
if provider.is_empty() || model.is_empty() { - 7421
continue; - 7422
} - 7423
match a.settlement { - 7424
vak_llm::Settlement::Ok => { - 7425
self.inner - 7426
.beliefs - 7427
.record_outcome(provider, model, a.domain, true); - 7428
} - 7429
vak_llm::Settlement::Failed => { - 7430
self.inner - 7431
.beliefs - 7432
.record_outcome(provider, model, a.domain, false); - 7433
} - 7434
_ => {} - 7435
} - 7436
} - 7437
} - 7438
} - 7439
- 7440
// The model is idle now that this turn is done: this is the only - 7441
// place a horizon-ladder probe may start (docs/design/68 §1), and - 7442
// it never blocks the return below. - 7443
self.maybe_start_capacity_probe(&turn_primary_leg).await; - 7444
- 7445
Ok((outcome, session)) - 7446
} - 7447
- 7448
fn next_checkpoint_seq(&self, session_id: &str) -> u32 { - 7449
checkpoints::next_seq(&self.sessions_home(), session_id) - 7450
} - 7451
- 7452
/// User-invoked compaction (`/compact`): summarize older turns into a - 7453
/// compaction entry now, regardless of the automatic trigger threshold. - 7454
/// Append-only; a receipt entry audits the summarizer dispatch. The - 7455
/// session always returns; failures land in `CompactOutcome.error`. - 7456
/// - 7457
/// Uses the same incremental, card-based mechanism as the agent loop's - 7458
/// own compaction (docs/design/68-context-engine.md §4): plan the - 7459
/// working set with a metadata-only `CapacityProfile` (no live route - 7460
/// leg to probe here), and if the plan finds a packet range, summarize - 7461
/// its turn cards and append one `Compaction` entry covering it. - 7462
pub async fn compact_session_now( - 7463
&self, - 7464
mut session: SessionLog, - 7465
cancel: tokio_util::sync::CancellationToken, - 7466
) -> (SessionLog, CompactOutcome) { - 7467
let profile = vak_context::capacity::CapacityProfile::from_metadata_only( - 7468
self.inner.config.context_window, - 7469
u64::from(self.inner.config.max_tokens), - 7470
"compact-session-now".to_string(), - 7471
std::time::SystemTime::now(), - 7472
); - 7473
let system = self.system_prompt(); - 7474
let tool_defs = vak_tools::definitions(&self.agent_tools()); - 7475
let prefix_chars = (system.len() as u64) - 7476
+ tool_defs - 7477
.iter() - 7478
.map(|t| { - 7479
(t.name.len() + t.description.len()) as u64 - 7480
+ serde_json::to_string(&t.parameters) - 7481
.map(|s| s.len() as u64) - 7482
.unwrap_or(0) - 7483
}) - 7484
.sum::<u64>(); - 7485
let prefix_tokens = profile.estimate_tokens(prefix_chars); - 7486
let plan_now = |session: &SessionLog| -> vak_session::WorkingSetPlan { - 7487
vak_context::plan_for_session(session, &profile, prefix_tokens, 0) - 7488
}; - 7489
let plan = plan_now(&session); - 7490
let Some((first_turn_id, last_turn_id)) = plan.packet_range else { - 7491
return ( - 7492
session, - 7493
CompactOutcome { - 7494
report: None, - 7495
error: None, - 7496
}, - 7497
); - 7498
}; - 7499
let (transcript, transcript_chars) = - 7500
session.packet_transcript(&first_turn_id, &last_turn_id); - 7501
let before = profile.estimate_tokens(transcript_chars); - 7502
let provider = match self.provider() { - 7503
Ok(p) => p, - 7504
Err(e) => return (session, CompactOutcome::failed(e.to_string())), - 7505
}; - 7506
let model = self.effective_model(); - 7507
let req = vak_context::assemble::compaction_request(&model, &transcript); - 7508
- 7509
let started = std::time::Instant::now(); - 7510
let mut receipt = vak_llm::WorkReceipt::new( - 7511
vak_llm::WorkPurpose::Summarize, - 7512
self.effective_provider(), - 7513
&model, - 7514
); - 7515
let summary = match provider.stream(req, cancel).await { - 7516
Ok(stream) => match stream.result().await { - 7517
Ok(msg) => { - 7518
receipt.record( - 7519
vak_llm::AttemptReason::Initial, - 7520
vak_llm::FailureDomain::Unknown, - 7521
vak_llm::Settlement::Ok, - 7522
started.elapsed().as_millis() as u64, - 7523
Some(msg.usage.clone()), - 7524
None, - 7525
); - 7526
msg.text_content() - 7527
} - 7528
Err(e) => { - 7529
receipt.record( - 7530
vak_llm::AttemptReason::Initial, - 7531
vak_llm::FailureDomain::Unknown, - 7532
vak_llm::Settlement::Failed, - 7533
started.elapsed().as_millis() as u64, - 7534
None, - 7535
Some(e.to_string()), - 7536
); - 7537
let _ = session.append_receipt(receipt); - 7538
return (session, CompactOutcome::failed(e.to_string())); - 7539
} - 7540
}, - 7541
Err(e) => { - 7542
receipt.record( - 7543
vak_llm::AttemptReason::Initial, - 7544
vak_llm::FailureDomain::Unknown, - 7545
vak_llm::Settlement::Cancelled, - 7546
started.elapsed().as_millis() as u64, - 7547
None, - 7548
Some(e.to_string()), - 7549
); - 7550
let _ = session.append_receipt(receipt); - 7551
return (session, CompactOutcome::failed(e.to_string())); - 7552
} - 7553
}; - 7554
if summary.trim().is_empty() { - 7555
let _ = session.append_receipt(receipt); - 7556
return ( - 7557
session, - 7558
CompactOutcome::failed("compaction produced an empty summary".into()), - 7559
); - 7560
} - 7561
let summarized = { - 7562
let index = vak_session::TurnIndex::from_log(&session); - 7563
let ids: Vec<&str> = index.turns.iter().map(|t| t.id.as_str()).collect(); - 7564
match ( - 7565
ids.iter().position(|id| *id == first_turn_id), - 7566
ids.iter().position(|id| *id == last_turn_id), - 7567
) { - 7568
(Some(lo), Some(hi)) => hi.saturating_sub(lo) + 1, - 7569
_ => 0, - 7570
} - 7571
}; - 7572
if let Err(e) = session.append_incremental_compaction( - 7573
&first_turn_id, - 7574
&last_turn_id, - 7575
&model, - 7576
summary, - 7577
before, - 7578
) { - 7579
return ( - 7580
session, - 7581
CompactOutcome::failed(format!("compaction write failed: {e}")), - 7582
); - 7583
} - 7584
// Semantic memory & entity distillation: distill learned invariants, - 7585
// domain procedural rules, and semantic entities before older history fades. - 7586
let _ = self.consolidate_memory(); - 7587
let _ = session.append_receipt(receipt); - 7588
let after = plan_now(&session).spent; - 7589
( - 7590
session, - 7591
CompactOutcome { - 7592
report: Some(CompactReport { - 7593
before_tokens: before, - 7594
after_tokens: after, - 7595
summarized_messages: summarized, - 7596
}), - 7597
error: None, - 7598
}, - 7599
) - 7600
} - 7601
- 7602
fn build_sandbox(&self) -> Option<std::sync::Arc<dyn vak_tools::sandbox::Sandbox>> { - 7603
let mode = match self.effective_permission_mode() { - 7604
vak_config::PermissionMode::ReadOnly => SandboxMode::ReadOnly, - 7605
vak_config::PermissionMode::WorkspaceWrite => SandboxMode::WorkspaceWrite, - 7606
vak_config::PermissionMode::FullAccess if self.task_copy_boundary => { - 7607
SandboxMode::WorkspaceWrite - 7608
} - 7609
vak_config::PermissionMode::FullAccess => return None, - 7610
}; - 7611
if self.task_copy_boundary { - 7612
if self.effective_sandbox_backend() == "docker" { - 7613
return Some(std::sync::Arc::new(vak_tools::sandbox::DenySandbox::new( - 7614
"strict task-copy containment is unavailable for the Docker backend", - 7615
))); - 7616
} - 7617
#[cfg(target_os = "macos")] - 7618
return Some(std::sync::Arc::new( - 7619
vak_tools::sandbox::Seatbelt::task_copy(mode, self.inner.cwd.as_path()), - 7620
)); - 7621
#[cfg(target_os = "linux")] - 7622
return Some(std::sync::Arc::new( - 7623
vak_tools::landlock::Landlock::task_copy(mode, self.inner.cwd.as_path()), - 7624
)); - 7625
#[cfg(not(any(target_os = "macos", target_os = "linux")))] - 7626
return Some(std::sync::Arc::new(vak_tools::sandbox::DenySandbox::new( - 7627
"strict task-copy containment is unsupported on this platform", - 7628
))); - 7629
} - 7630
let backend = self.effective_sandbox_backend(); - 7631
let backend = backend.as_str(); - 7632
if backend == "docker" { - 7633
// Fail closed at call time if the daemon is unreachable: - 7634
// BashTool surfaces the wrapped command's error verbatim, and - 7635
// the probe keeps startup cheap. - 7636
return Some(std::sync::Arc::new(sandbox_docker::DockerSandbox::new( - 7637
mode,
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.