- 6207
} else { - 6208
WorkMode::Direct - 6209
}; - 6210
} - 6211
cfg.work_enabled = work_config.enabled; - 6212
cfg.max_work_items = work_config.max_items; - 6213
cfg.max_work_revisions = work_config.max_revisions; - 6214
let capabilities = turn_capabilities.descriptors.clone(); - 6215
cfg.input_normalizer = Some(Arc::new(move |message| { - 6216
normalize_capability_message(message, &capabilities) - 6217
})); - 6218
cfg.model = model.clone(); - 6219
cfg.tools = self.agent_tools(); - 6220
// The intent cap is enforced by the admission/commitment contract; - 6221
// the agent loop's counter includes tool round-trips and is therefore - 6222
// not a faithful model-turn budget for general-purpose turns. - 6223
cfg.max_turns = self.effective_max_turns(); - 6224
cfg.parallel_tools = true; - 6225
cfg.max_retries = self.inner.config.max_retries; - 6226
cfg.retry_base_backoff_ms = self.inner.config.retry_base_backoff_ms; - 6227
cfg.run_retry_attempts = self.inner.config.run_retry_attempts; - 6228
cfg.run_retry_base_backoff_ms = self.inner.config.run_retry_base_backoff_ms; - 6229
cfg.request_timeout = if self.inner.config.request_timeout_secs == 0 { - 6230
None - 6231
} else { - 6232
Some(std::time::Duration::from_secs( - 6233
self.inner.config.request_timeout_secs, - 6234
)) - 6235
}; - 6236
// Per-turn route planning: assemble a fresh ladder using the current - 6237
// evidence ledger, belief state, warm discovery cache, and demand facts - 6238
// from this turn's intent resolution. This replaces the admission-frozen - 6239
// ladder; the session header's route_ladder is now an initial snapshot. - 6240
let needs_tools = !self.tool_names().is_empty(); - 6241
let turn_primary_credential_id = self - 6242
.provider_auth_for_leg(&self.effective_provider(), None) - 6243
.ok() - 6244
.and_then(|auth| auth.credential_id); - 6245
// An injected provider instance serves the turn itself, so it is the - 6246
// identity routing and receipts record; otherwise the configured one. - 6247
let turn_primary_provider = if self.provider_instance_override().is_some() { - 6248
provider.name().to_string() - 6249
} else { - 6250
self.effective_provider() - 6251
}; - 6252
let turn_primary_leg = vak_llm::RouteLeg { - 6253
provider: turn_primary_provider.clone(), - 6254
model: model.clone(), - 6255
// The dialect names the wire the turn is actually served on, - 6256
// which is the adapter's, not the configured provider's. - 6257
dialect: vak_llm::EndpointDialect::for_provider( - 6258
provider.name(), - 6259
needs_tools || engagement.posture.demand.reasoning_required, - 6260
), - 6261
credential_id: turn_primary_credential_id, - 6262
}; - 6263
cfg.provider_name = Some(turn_primary_provider.clone()); - 6264
let mut turn_plan = - 6265
self.plan_route_ladder(turn_primary_leg.clone(), Some(engagement.posture.demand)); - 6266
// The engagement's modality constraint: a leg that cannot see is not - 6267
// a valid fallback for a vision turn. With no operator hints every - 6268
// leg is assumed capable; with hints and no capable leg, the turn - 6269
// fails typed rather than quietly dropping the image (invariant 10). - 6270
// The reading never shortens the ladder: a fallback is resilience, - 6271
// and a misread greeting must not cost a turn its recovery. - 6272
let modality_hints = self.inner.config.route.modality_hints.clone(); - 6273
if !engagement.limits.required_modalities.is_empty() && !modality_hints.is_empty() { - 6274
let supports = |model: &str| { - 6275
intent::leg_supports_modalities( - 6276
model, - 6277
&engagement.limits.required_modalities, - 6278
&modality_hints, - 6279
) - 6280
}; - 6281
// The primary leg is dispatched first whatever the ladder says, - 6282
// so it has to be capable itself; fallbacks are then filtered. - 6283
if !supports(&turn_primary_leg.model) { - 6284
let wanted: Vec<&str> = engagement - 6285
.limits - 6286
.required_modalities - 6287
.iter() - 6288
.map(|m| m.as_str()) - 6289
.collect(); - 6290
return Err(CoreError::UnsupportedModality { - 6291
modalities: wanted.join(", "), - 6292
model: turn_primary_leg.model.clone(), - 6293
}); - 6294
} - 6295
turn_plan.ladder = turn_plan - 6296
.ladder - 6297
.iter() - 6298
.enumerate() - 6299
.filter(|(index, leg)| *index == 0 || supports(&leg.model)) - 6300
.map(|(_, leg)| leg.clone()) - 6301
.collect(); - 6302
} - 6303
let (context_window, max_output) = self - 6304
.route_context_limits(&turn_primary_leg, &turn_plan.ladder) - 6305
.await; - 6306
cfg.declared_window = context_window; - 6307
cfg.max_output = max_output; - 6308
// Measured capacity (docs/design/68-context-engine.md §1): bound - 6309
// immediately from what is already known (cache, ledger, or a - 6310
// metadata-only profile) — never from a live probe. The ladder - 6311
// itself, when this key needs one, runs in the background after - 6312
// this turn completes (see the `maybe_start_capacity_probe` call - 6313
// near this function's return). `session` is the same ledger this - 6314
// turn is about to append to, so the bind's Activity lands before - 6315
// the turn's own messages. - 6316
let capacity = self - 6317
.capacity_profile_for(&turn_primary_leg, &mut session) - 6318
.await; - 6319
cfg.capacity_key = Some(vak_context::capacity::ProfileKey { - 6320
provider: turn_primary_leg.provider.clone(), - 6321
model: turn_primary_leg.model.clone(), - 6322
quantisation: capacity.provenance.quantisation.clone(), - 6323
}); - 6324
cfg.capacity = Some(capacity); - 6325
let sp = &self.inner.config.stop_policy; - 6326
cfg.stop_policy = if sp.enabled { - 6327
Some(vak_agent::StopPolicy { - 6328
marker_gate: sp.marker_gate, - 6329
verify_gate: sp.verify_gate, - 6330
max_blocks: sp.max_blocks, - 6331
}) - 6332
} else { - 6333
None - 6334
}; - 6335
cfg.circuit_breaker = Some(self.inner.breaker.clone()); - 6336
cfg.handoff_reset = self.inner.config.goal.handoff_reset; - 6337
cfg.max_audit_blocks = self.inner.config.goal.max_audit_blocks; - 6338
cfg.approver = approver.clone(); - 6339
// The engagement supplies a CEILING on approval permissiveness, never - 6340
// a floor: `approval_mode` takes the stricter of it and configuration, - 6341
// so an irreversible turn reaches a human even under auto-approve, and - 6342
// nothing here can skip a gate the operator asked for. - 6343
cfg.approval_mode = match intent::approval_mode( - 6344
self.effective_approval_mode(), - 6345
engagement.limits.approval_ceiling, - 6346
self.inner.config.intent.posture, - 6347
) { - 6348
vak_config::ApprovalMode::Ask => vak_agent::ApprovalMode::Ask, - 6349
vak_config::ApprovalMode::ApproveSafe => vak_agent::ApprovalMode::ApproveSafe, - 6350
vak_config::ApprovalMode::AutoApprove => vak_agent::ApprovalMode::AutoApprove, - 6351
}; - 6352
// Every provider dispatch is settled into the FinOps ledger. Budget - 6353
// admission remains a no-op when no cap is configured, while unknown - 6354
// prices are retained as explicit unpriced rows for auditability. - 6355
// Reuse one gate per session so run/day accounting spans turns. - 6356
let sid = session - 6357
.header() - 6358
.map(|h| h.session_id.clone()) - 6359
.unwrap_or_default(); - 6360
let turn_gate = self.spend_gate_for(&sid); - 6361
// An envelope's lifetime spend limit meets the configured run cap; - 6362
// the smaller governs. - 6363
if let Some(ceiling) = engagement.limits.spend_ceiling_usd { - 6364
turn_gate.narrow_run_cap(ceiling); - 6365
} - 6366
cfg.spend_gate = Some(turn_gate); - 6367
- 6368
// MEA substrate (Phase H): auditor sees the workspace delta between - 6369
// this run's start checkpoint and the live tree. - 6370
{ - 6371
let home = self.sessions_home(); - 6372
let seq = self.next_checkpoint_seq(&sid); - 6373
let cwd = self.inner.cwd.clone(); - 6374
cfg.workspace_delta = Some(Arc::new(CheckpointDelta { - 6375
home: home.clone(), - 6376
sid: sid.clone(), - 6377
seq, - 6378
cwd, - 6379
})); - 6380
} - 6381
// An envelope's permission ceiling narrows the mode through the same - 6382
// door a gateway channel override uses; it can never raise it. - 6383
cfg.mode = match intent::permission_mode( - 6384
self.effective_permission_mode(), - 6385
engagement.limits.permission_ceiling, - 6386
) { - 6387
vak_config::PermissionMode::ReadOnly => vak_permission::Mode::ReadOnly, - 6388
vak_config::PermissionMode::WorkspaceWrite => vak_permission::Mode::WorkspaceWrite, - 6389
vak_config::PermissionMode::FullAccess => vak_permission::Mode::FullAccess, - 6390
}; - 6391
if self.task_copy_boundary && cfg.mode == vak_permission::Mode::FullAccess { - 6392
cfg.mode = vak_permission::Mode::WorkspaceWrite; - 6393
} - 6394
let permission_rules = self.channel_permission_rules(); - 6395
cfg.permission = Some(match permission { - 6396
Some(p) => p, - 6397
None => std::sync::Arc::new(self.build_permission_engine(&permission_rules)?), - 6398
}); - 6399
if self.task_copy_boundary { - 6400
let Some(engine) = cfg.permission.take() else { - 6401
return Err(CoreError::MissingEngine); - 6402
}; - 6403
cfg.permission = Some(std::sync::Arc::new( - 6404
engine.as_ref().clone().restrict_tools(TASK_COPY_TOOLS), - 6405
)); - 6406
} - 6407
let Some(engine) = cfg.permission.clone() else { - 6408
return Err(CoreError::MissingEngine); - 6409
}; - 6410
let session_id = session - 6411
.header() - 6412
.map(|header| header.session_id.clone()) - 6413
.unwrap_or_default(); - 6414
cfg.sandbox = self.session_sandbox(&session_id); - 6415
- 6416
// Per-turn fallback ladder: built from the freshly planned turn_plan, - 6417
// not the admission-frozen session contract. This reflects the current - 6418
// evidence ledger and belief state, so a provider that failed earlier - 6419
// this session or was demoted by the auto-routing algorithm is correctly - 6420
// ranked. Unresolvable legs (missing key/registry) skip silently. - 6421
for leg in turn_plan.ladder.iter().skip(1) { - 6422
if let Ok(auth) = - 6423
self.provider_auth_for_leg(&leg.provider, leg.credential_id.as_deref()) - 6424
&& let Ok(p) = self.inner.registry.get(&adapter_name_for_leg(leg), &auth) - 6425
{ - 6426
cfg.ladder.push((p, leg.model.clone())); - 6427
cfg.ladder_provider_names.push(leg.provider.clone()); - 6428
} - 6429
} - 6430
- 6431
let mut tools = self.scoped_tools(&ToolScope { - 6432
session_id: session - 6433
.header() - 6434
.map(|h| h.session_id.clone()) - 6435
.unwrap_or_default(), - 6436
agent_id: session - 6437
.header() - 6438
.and_then(|header| header.agent.as_ref().map(|agent| agent.id.clone())) - 6439
.or_else(|| self.agent_identity.as_ref().map(|agent| agent.id.clone())), - 6440
audience_id: session.header().and_then(|header| { - 6441
header - 6442
.conversation - 6443
.as_ref() - 6444
.map(|context| context.audience_id.clone()) - 6445
}), - 6446
}); - 6447
let frozen_skills = std::mem::take(&mut turn_capabilities.frozen_skills); - 6448
let skill_tool = (!frozen_skills.is_empty()) - 6449
.then(|| Arc::new(skills::SkillTool::new(frozen_skills)) as Arc<dyn vak_tools::Tool>); - 6450
if let Some(skill_tool) = &skill_tool { - 6451
tools.push(skill_tool.clone()); - 6452
} - 6453
if let Some(manager) = self.mcp_manager() { - 6454
// Turn admission never blocks on an optional integration: the - 6455
// manager is reused across turns and this lazy meta-tool only - 6456
// connects when the model calls it. - 6457
let policy = self.channel_policy().unwrap_or_default(); - 6458
let context = self.plugin_mcp_invocation_context(); - 6459
let activity_ledger = finops::ActivityLedger::new(&self.sessions_home()); - 6460
let activity_session = session.header().map(|h| h.session_id.clone()); - 6461
let recorder = Arc::new( - 6462
move |server: &str, tool: &str, success: bool, duration_ms: u64| { - 6463
let plugin = context - 6464
.iter() - 6465
.find(|(_, plugin, _)| server.starts_with(&format!("plugin.{plugin}."))) - 6466
.map(|(_, plugin, _)| plugin.clone()); - 6467
let _ = activity_ledger.append(&finops::ActivityRow { - 6468
ts: chrono::Utc::now(), - 6469
kind: "mcp".into(), - 6470
name: format!("{server}/{tool}"), - 6471
success, - 6472
duration_ms: Some(duration_ms), - 6473
session_id: activity_session.clone(), - 6474
plugin, - 6475
}); - 6476
for (store, plugin, trace_id) in &context { - 6477
if server.starts_with(&format!("plugin.{plugin}.")) { - 6478
let _ = store.record_invocation( - 6479
trace_id, - 6480
plugin, - 6481
&format!("mcp:{server}/{tool}"), - 6482
success, - 6483
); - 6484
break; - 6485
} - 6486
} - 6487
}, - 6488
); - 6489
let admitted_servers = turn_capabilities.mcp_server_names.clone(); - 6490
// One index for the config's whole life, filled in place so the - 6491
// retrieval check that captured it earlier sees live entries. - 6492
let index = cfg.mcp_tool_index.clone(); - 6493
*index - 6494
.lock() - 6495
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 6496
turn_capabilities.mcp_tool_index.clone(); - 6497
let admitted_for_catalog: std::collections::BTreeSet<String> = - 6498
admitted_servers.iter().cloned().collect(); - 6499
let builtins_for_catalog: std::collections::BTreeSet<String> = - 6500
self.tool_names().into_iter().collect(); - 6501
let mcp_tool = vak_mcp::McpTool::with_policy_and_recorder( - 6502
manager, - 6503
Some(policy.mcp_allow.clone().unwrap_or_else(|| { - 6504
admitted_servers.iter().map(|s| format!("{s}/*")).collect() - 6505
})), - 6506
policy.mcp_deny.clone(), - 6507
recorder, - 6508
) - 6509
.with_catalog_observer(Arc::new(move |catalog| { - 6510
*index - 6511
.lock() - 6512
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 6513
capability::turn::mcp_tool_index( - 6514
catalog, - 6515
&admitted_for_catalog, - 6516
&builtins_for_catalog, - 6517
); - 6518
})); - 6519
tools.push(Arc::new(mcp_tool)); - 6520
} - 6521
// Admission owns all filtering: factories may construct a tool the - 6522
// turn did not admit, and it is dropped here, before a child could - 6523
// inherit it. - 6524
tools.retain(|tool| turn_capabilities.tool_names.contains(tool.name())); - 6525
if turn_capabilities.tool_names.contains("task") - 6526
&& let Some(parent_id) = session.header().map(|h| h.session_id.clone()) - 6527
{ - 6528
let managed_projection = session.work_projection().ok().flatten(); - 6529
let managed_contract_id = managed_projection - 6530
.as_ref() - 6531
.map(|projection| projection.contract.contract_id.clone()); - 6532
let managed_work_item_ids = managed_projection - 6533
.map(|projection| projection.items.into_keys().collect()) - 6534
.unwrap_or_default(); - 6535
let mut read_only_tools = self.agent_read_only_tools(); - 6536
if let Some(skill_tool) = &skill_tool { - 6537
read_only_tools.push(skill_tool.clone()); - 6538
} - 6539
// A child's reader is this agent, not a person, so it gets the - 6540
// `Worker` surface rather than inheriting a human-facing one — - 6541
// a research child spawned from a phone chat is not on a phone. - 6542
let child_core = self.clone().with_surface(Surface::Worker); - 6543
let child_capability_set = turn_capabilities.descriptors.clone(); - 6544
let child_default_prompt = - 6545
child_core.system_prompt_for_capabilities(&child_capability_set); - 6546
let role_prompts = child_core.role_prompts(&child_capability_set); - 6547
tools.push(Arc::new(vak_agent::TaskTool::new(vak_agent::TaskDeps { - 6548
parent_agent_identity: self.agent_identity().cloned(), - 6549
outcome_objective: Some(prompt_text.to_string()), - 6550
outcome: cfg.outcome.clone(), - 6551
provider: provider.clone(), - 6552
system_prompt: child_default_prompt, - 6553
tail: cfg.tail.clone(), - 6554
role_prompts, - 6555
model: model.clone(), - 6556
tools: tools.clone(), - 6557
capabilities: turn_capabilities.descriptors.clone(), - 6558
hooks: cfg.hooks.clone(), - 6559
revocation_check: cfg.revocation_check.clone(), - 6560
presentation_rebuild: cfg.presentation_rebuild.clone(), - 6561
mcp_tool_index: Some(cfg.mcp_tool_index.clone()), - 6562
input_normalizer: cfg.input_normalizer.clone(), - 6563
read_only_tools, - 6564
max_turns: self.effective_max_turns(), - 6565
max_retries: cfg.max_retries, - 6566
retry_base_backoff_ms: cfg.retry_base_backoff_ms, - 6567
request_timeout: cfg.request_timeout, - 6568
circuit_breaker: cfg.circuit_breaker.clone(), - 6569
run_retry_attempts: cfg.run_retry_attempts, - 6570
run_retry_base_backoff_ms: cfg.run_retry_base_backoff_ms, - 6571
dispatch_ceiling: cfg.dispatch_ceiling, - 6572
spend_gate: cfg.spend_gate.clone(), - 6573
permission: Some(engine.clone()), - 6574
mode: cfg.mode, - 6575
approval_mode: cfg.approval_mode, - 6576
approver: approver.clone(), - 6577
sandbox: cfg.sandbox.clone(), - 6578
cwd: self.inner.cwd.clone(), - 6579
sessions_home: self.inner.sessions_home.clone(), - 6580
parent_session_id: parent_id, - 6581
contract_id: managed_contract_id, - 6582
work_item_id: None, - 6583
work_item_ids: managed_work_item_ids, - 6584
events: Some(events.clone()), - 6585
registry: Some(self.inner.workers.clone()), - 6586
}))); - 6587
} - 6588
let turn_standings = self.capability_standings(); - 6589
for detail in reach::audit_details(&turn_standings) { - 6590
security_events::record( - 6591
&self.sessions_home(), - 6592
security_events::EventKind::CapabilityUnreachable, - 6593
"capability_unreachable", - 6594
&detail, - 6595
None, - 6596
); - 6597
} - 6598
let predicted_cards = presentation_tools::predicted_card_tools( - 6599
&prompt_text, - 6600
&vak_delivery::built_in_recipes(), - 6601
); - 6602
let surface = capability::build_tool_surface(&tools, &loaded_domains, &predicted_cards); - 6603
// `find_tools` is synthetic and never itself deferred. Search spans - 6604
// every admitted tool, not only the deferred subset: weaker/local - 6605
// models often use discovery to relocate an already-loaded broker - 6606
// such as `mcp`, and an empty result must not be mistaken for absence. - 6607
// Domain labels stay search-only rather than mutating tool schemas. - 6608
let searchable_tools = surface - 6609
.core - 6610
.iter() - 6611
.chain(surface.deferred.iter()) - 6612
.cloned() - 6613
.collect(); - 6614
let search_keywords = tools - 6615
.iter() - 6616
.map(|tool| { - 6617
( - 6618
tool.name().to_string(), - 6619
tool.serves() - 6620
.iter() - 6621
.map(|domain| (*domain).to_string()) - 6622
.collect(), - 6623
) - 6624
}) - 6625
.collect(); - 6626
let find_tools_tool = Arc::new( - 6627
vak_tools::FindToolsTool::new(searchable_tools) - 6628
.with_keywords(search_keywords) - 6629
.with_discovered_sink(cfg.discovered_tools.clone()), - 6630
); - 6631
let find_tools_def = vak_llm::ToolDefinition::new( - 6632
vak_tools::Tool::name(find_tools_tool.as_ref()), - 6633
vak_tools::Tool::description(find_tools_tool.as_ref()), - 6634
vak_tools::Tool::schema(find_tools_tool.as_ref()), - 6635
); - 6636
tools.push(find_tools_tool); - 6637
let core_tool_names: Vec<String> = - 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,
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.