- 5456
- 5457
#[allow(clippy::too_many_arguments)] - 5458
pub async fn run_managed_turn_with( - 5459
&self, - 5460
session: SessionLog, - 5461
prompt: &str, - 5462
cancel: CancellationToken, - 5463
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5464
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5465
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5466
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5467
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5468
self.clone() - 5469
.run_turn_inner( - 5470
session, - 5471
vak_session::MessageRecord { - 5472
message: vak_llm::Message::user_text(prompt), - 5473
meta: None, - 5474
}, - 5475
cancel, - 5476
approver, - 5477
permission, - 5478
steering, - 5479
events, - 5480
None, - 5481
Some(WorkMode::Managed), - 5482
) - 5483
.await - 5484
} - 5485
- 5486
#[allow(clippy::too_many_arguments)] - 5487
pub async fn run_auto_turn_with( - 5488
&self, - 5489
session: SessionLog, - 5490
prompt: &str, - 5491
cancel: CancellationToken, - 5492
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5493
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5494
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5495
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5496
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5497
self.clone() - 5498
.run_turn_inner( - 5499
session, - 5500
vak_session::MessageRecord { - 5501
message: vak_llm::Message::user_text(prompt), - 5502
meta: None, - 5503
}, - 5504
cancel, - 5505
approver, - 5506
permission, - 5507
steering, - 5508
events, - 5509
None, - 5510
Some(WorkMode::Auto), - 5511
) - 5512
.await - 5513
} - 5514
- 5515
/// Goal-mode turn (docs/design/42-managed-work-contracts.md): the run may only end when - 5516
/// the objective's acceptance criteria pass an independent audit. - 5517
#[allow(clippy::too_many_arguments)] - 5518
pub async fn run_goal_turn_with( - 5519
&self, - 5520
session: SessionLog, - 5521
prompt: &str, - 5522
objective: &str, - 5523
criteria: Vec<String>, - 5524
cancel: CancellationToken, - 5525
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5526
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5527
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5528
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5529
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5530
self.clone() - 5531
.run_turn_inner( - 5532
session, - 5533
vak_session::MessageRecord { - 5534
message: vak_llm::Message::user_text(prompt), - 5535
meta: None, - 5536
}, - 5537
cancel, - 5538
approver, - 5539
permission, - 5540
steering, - 5541
events, - 5542
Some((objective.to_string(), criteria)), - 5543
None, - 5544
) - 5545
.await - 5546
} - 5547
- 5548
/// The standing authority for this workspace and surface. - 5549
/// - 5550
/// Autonomy is *granted* (configuration, and only from a trusted project); - 5551
/// attendance is *observed* (which surface this is, and whether the - 5552
/// installed approver can actually answer). Keeping them separate is what - 5553
/// stops vak nagging when you wanted autonomy and barrelling ahead when - 5554
/// nobody is watching. - 5555
pub fn turn_authority(&self) -> vak_intent::Authority { - 5556
self.turn_authority_for(self.surface()) - 5557
} - 5558
- 5559
/// The authority that would apply on a named surface. - 5560
/// - 5561
/// Separate from [`Core::turn_authority`] because `vak intent explain - 5562
/// --surface cron` asks "what would happen there", and deriving attendance - 5563
/// from *this* process's surface would answer a different question — it - 5564
/// reported a cron run as interactive. - 5565
pub fn turn_authority_for(&self, surface: &Surface) -> vak_intent::Authority { - 5566
let surface = intent::intent_surface(surface); - 5567
// An approver that cannot answer makes the surface unattended - 5568
// whatever it claims to be — the same fact `reach` already uses to - 5569
// stop advertising capabilities nobody can approve. - 5570
let attendance = if self.approver_answerable() { - 5571
surface.implied_attendance() - 5572
} else { - 5573
vak_intent::Attendance::Unattended - 5574
}; - 5575
// A channel may cap delegation below what the workspace granted and - 5576
// may never raise it, exactly like every other key on ChannelPolicy. - 5577
// A Telegram chat gets to say "propose only, in here"; it does not - 5578
// get to say "act freely" on a workspace whose operator did not. - 5579
let mut autonomy = intent::configured_autonomy(&self.inner.config); - 5580
if let Some(policy) = self.channel_policy() - 5581
&& let Some(ceiling) = policy - 5582
.autonomy_ceiling - 5583
.as_deref() - 5584
.and_then(vak_intent::Autonomy::parse) - 5585
{ - 5586
autonomy = autonomy.capped_by(ceiling); - 5587
} - 5588
vak_intent::Authority { - 5589
autonomy, - 5590
attendance, - 5591
} - 5592
} - 5593
- 5594
/// Resolve this turn's intent from the prompt and the session so far. - 5595
/// - 5596
/// Free tiers only; the surfaces that preview a reading (`vak intent - 5597
/// explain`, `GET /intent/explain`, the composer strip) use this. The - 5598
/// turn path uses [`Core::resolve_turn_intent_with_escalation`], which - 5599
/// may spend a classification dispatch on a weak reading. - 5600
pub fn resolve_turn_intent( - 5601
&self, - 5602
session: &SessionLog, - 5603
prompt: &vak_llm::Message, - 5604
) -> vak_intent::Intent { - 5605
self.resolve_turn_free_tiers(session, prompt, "").intent() - 5606
} - 5607
- 5608
/// Whether `reading` decided which admitted tools a turn loaded — the - 5609
/// kernel slices, and the reading's act was confident enough to — as - 5610
/// opposed to the orientation floor standing in for a reading too weak - 5611
/// to decide. Only the former can be contradicted by a deferred tool the - 5612
/// model then used (`misread::MisreadRow::sliced`). - 5613
fn reading_sliced(&self, reading: &vak_intent::Reading) -> bool { - 5614
let intent = &self.inner.config.intent; - 5615
intent.enabled - 5616
&& intent.slice_capabilities - 5617
&& reading.may_slice_capabilities(intent.accept_confidence) - 5618
} - 5619
- 5620
/// The free tiers, under the turn's base authority. A grant on a - 5621
/// commitment is applied afterwards, once the turn knows which - 5622
/// commitments its strands work on (`vak_intent::apply_envelopes`). - 5623
fn resolve_turn_free_tiers( - 5624
&self, - 5625
session: &SessionLog, - 5626
prompt: &vak_llm::Message, - 5627
turn_id: &str, - 5628
) -> vak_intent::Resolution { - 5629
let text = prompt.text_content(); - 5630
let attachments: Vec<vak_intent::Attachment> = prompt - 5631
.content - 5632
.iter() - 5633
.filter_map(|block| match block { - 5634
vak_llm::ContentBlock::Image { .. } => Some(vak_intent::Attachment { - 5635
modality: vak_intent::Modality::Image, - 5636
name: "attachment".into(), - 5637
}), - 5638
_ => None, - 5639
}) - 5640
.collect(); - 5641
intent::resolve_turn( - 5642
&text, - 5643
turn_id, - 5644
self.surface(), - 5645
&attachments, - 5646
intent::workspace_facts(&self.inner.cwd), - 5647
intent::history_facts(session), - 5648
&vak_intent::Declared::default(), - 5649
&self.turn_authority(), - 5650
&intent::resolver_config(&self.inner.config), - 5651
) - 5652
} - 5653
- 5654
/// Resolve this turn's intent, spending a classification dispatch when - 5655
/// the free tiers were not confident and `[intent] escalate` allows it. - 5656
/// - 5657
/// The dispatch is a dispatch like any other: a `Classify` work receipt - 5658
/// on the session, spend-gate admission under `max_classify_usd`, a - 5659
/// short watchdog, the run's cancellation token — and **fail-open** to - 5660
/// the free-tier reading on any failure. A classifier outage must never - 5661
/// block work, and the partial is never worse than the orienting - 5662
/// engagement. - 5663
/// - 5664
/// `escalate = "local"` means a model on this machine: only the keyless - 5665
/// `ollama` provider serves it, whatever `classify_model` names, so a - 5666
/// local setting can never send the request text off the machine. - 5667
pub async fn resolve_turn_intent_with_escalation( - 5668
&self, - 5669
session: &mut SessionLog, - 5670
prompt: &vak_llm::Message, - 5671
turn_id: &str, - 5672
cancel: &CancellationToken, - 5673
) -> vak_intent::Intent { - 5674
let resolution = self.resolve_turn_free_tiers(session, prompt, turn_id); - 5675
let (partial, reason) = match resolution { - 5676
vak_intent::Resolution::Settled(intent) => return intent, - 5677
vak_intent::Resolution::Escalate { partial, reason } => (partial, reason), - 5678
}; - 5679
let escalate = self.inner.config.intent.escalate.as_str(); - 5680
if escalate == "none" { - 5681
return partial; - 5682
} - 5683
let cloud = escalate == "cloud"; - 5684
let mut partial = partial; - 5685
let give_up = |partial: &mut vak_intent::Intent, why: String| { - 5686
partial.provenance.escalation_note = Some(format!( - 5687
"{}; escalation ({reason}) skipped: {why}", - 5688
partial - 5689
.provenance - 5690
.escalation_note - 5691
.clone() - 5692
.unwrap_or_default() - 5693
)); - 5694
}; - 5695
- 5696
// --- which leg ------------------------------------------------------ - 5697
// `classify_model` may be `provider/model` or a bare model name. Local - 5698
// escalation runs on the keyless `ollama` provider and nothing else — - 5699
// an Ollama model name may itself contain `/` (`hf.co/org/model`), so - 5700
// only a prefix naming another configured provider is refused rather - 5701
// than read as a model; cloud escalation runs on the effective - 5702
// provider unless a provider was named. - 5703
let configured = self.inner.config.intent.classify_model.clone(); - 5704
let known_providers = self.provider_names(); - 5705
let (provider_name, model) = match (cloud, configured) { - 5706
(true, Some(spec)) if spec.contains('/') => { - 5707
let (p, m) = spec.split_once('/').unwrap_or(("", "")); - 5708
(p.to_string(), m.to_string()) - 5709
} - 5710
(true, Some(model)) => (self.effective_provider(), model), - 5711
(true, None) => (self.effective_provider(), self.effective_model()), - 5712
(false, Some(spec)) => match spec.split_once('/') { - 5713
Some(("ollama", model)) => ("ollama".to_string(), model.to_string()), - 5714
Some((prefix, _)) if known_providers.iter().any(|p| p == prefix) => { - 5715
give_up( - 5716
&mut partial, - 5717
format!( - 5718
"escalate = \"local\" runs only on ollama, but classify_model names {prefix}" - 5719
), - 5720
); - 5721
return partial; - 5722
} - 5723
_ => ("ollama".to_string(), spec), - 5724
}, - 5725
(false, None) => { - 5726
if self.effective_provider() == "ollama" { - 5727
("ollama".to_string(), self.effective_model()) - 5728
} else { - 5729
give_up( - 5730
&mut partial, - 5731
"escalate = \"local\" needs [intent] classify_model or an ollama route" - 5732
.into(), - 5733
); - 5734
return partial; - 5735
} - 5736
} - 5737
}; - 5738
let provider = match self - 5739
.provider_auth_for_leg(&provider_name, None) - 5740
.and_then(|auth| { - 5741
self.inner - 5742
.registry - 5743
.get(&provider_name, &auth) - 5744
.map_err(CoreError::from) - 5745
}) { - 5746
Ok(provider) => provider, - 5747
Err(error) => { - 5748
give_up( - 5749
&mut partial, - 5750
format!("no usable {provider_name} provider: {error}"), - 5751
); - 5752
return partial; - 5753
} - 5754
}; - 5755
- 5756
// --- the request ------------------------------------------------------ - 5757
let user_prompt = vak_intent::classification_prompt(&partial); - 5758
let digest = vak_intent::prompt_digest(&user_prompt); - 5759
// One object per part: a fixed budget truncated the answer for a - 5760
// request with more than a few parts, and a truncated array parses - 5761
// as nothing. - 5762
let output_budget = vak_intent::classification_budget(partial.strands.len()); - 5763
let mut request = vak_llm::ChatRequest::new(&model); - 5764
request.system = - 5765
Some("You classify requests for an agent runtime. Answer with JSON only.".to_string()); - 5766
request.messages = vec![vak_llm::Message::user_text(user_prompt.clone())]; - 5767
request.max_tokens = output_budget; - 5768
// A strict-JSON answer, not a deliberation: measured live on a - 5769
// thinking model, the default spent the whole budget in its thinking - 5770
// channel and returned nothing. - 5771
request.think = Some(false); - 5772
- 5773
// --- admission ---------------------------------------------------------- - 5774
use vak_agent::SpendGate as _; - 5775
let session_id = session - 5776
.header() - 5777
.map(|header| header.session_id.clone()) - 5778
.unwrap_or_default(); - 5779
let gate = self.spend_gate_for(&session_id); - 5780
let planned = vak_llm::Usage { - 5781
input_tokens: (user_prompt.len() / 4) as u64 + 64, - 5782
output_tokens: u64::from(output_budget), - 5783
..Default::default() - 5784
}; - 5785
let cap = self.inner.config.intent.max_classify_usd; - 5786
// Every paid dispatch needs a price it can be held to; only the - 5787
// keyless local provider may run unpriced. A non-finite estimate is - 5788
// no estimate: `NaN > cap` is false and would wave anything through. - 5789
match gate.estimate_usd(&model, &planned) { - 5790
Some(est) if !est.is_finite() || est > cap => { - 5791
give_up( - 5792
&mut partial, - 5793
format!("estimated ${est:.4} exceeds max_classify_usd ${cap:.4}"), - 5794
); - 5795
return partial; - 5796
} - 5797
None if provider_name != "ollama" => { - 5798
give_up(&mut partial, format!("no price known for {model}")); - 5799
return partial; - 5800
} - 5801
_ => {} - 5802
} - 5803
if let Err(denied) = gate - 5804
.authorize(&vak_agent::SpendCheck { - 5805
model: &model, - 5806
provider: &provider_name, - 5807
session_id: &session_id, - 5808
est_input_tokens: planned.input_tokens, - 5809
planned_output_tokens: planned.output_tokens, - 5810
}) - 5811
.await - 5812
{ - 5813
give_up(&mut partial, format!("spend gate refused: {denied}")); - 5814
return partial; - 5815
} - 5816
- 5817
// --- dispatch ----------------------------------------------------------- - 5818
let watchdog = - 5819
std::time::Duration::from_secs(self.inner.config.intent.classify_timeout_secs); - 5820
let started = std::time::Instant::now(); - 5821
let mut receipt = - 5822
vak_llm::WorkReceipt::new(vak_llm::WorkPurpose::Classify, &provider_name, &model); - 5823
let child = cancel.child_token(); - 5824
let outcome = tokio::time::timeout(watchdog, async { - 5825
provider.stream(request, child).await?.result().await - 5826
}) - 5827
.await; - 5828
let answer = match outcome { - 5829
Ok(Ok(message)) => { - 5830
receipt.record( - 5831
vak_llm::AttemptReason::Initial, - 5832
vak_llm::FailureDomain::Unknown, - 5833
vak_llm::Settlement::Ok, - 5834
started.elapsed().as_millis() as u64, - 5835
Some(message.usage.clone()), - 5836
None, - 5837
); - 5838
gate.record_settled_with_latency( - 5839
&provider_name, - 5840
&model, - 5841
&session_id, - 5842
&message.usage, - 5843
started.elapsed().as_millis() as u64, - 5844
); - 5845
let _ = session.append_receipt(receipt); - 5846
message.text_content() - 5847
} - 5848
Ok(Err(error)) => { - 5849
receipt.record( - 5850
vak_llm::AttemptReason::Initial, - 5851
vak_llm::FailureDomain::Unknown, - 5852
vak_llm::Settlement::Failed, - 5853
started.elapsed().as_millis() as u64, - 5854
None, - 5855
Some(error.to_string()), - 5856
); - 5857
let _ = session.append_receipt(receipt); - 5858
give_up(&mut partial, format!("classifier failed: {error}")); - 5859
return partial; - 5860
} - 5861
Err(_) => { - 5862
receipt.record( - 5863
vak_llm::AttemptReason::Initial, - 5864
vak_llm::FailureDomain::Unknown, - 5865
vak_llm::Settlement::Cancelled, - 5866
started.elapsed().as_millis() as u64, - 5867
None, - 5868
Some(format!("watchdog {}s", watchdog.as_secs())), - 5869
); - 5870
let _ = session.append_receipt(receipt); - 5871
give_up(&mut partial, "classifier exceeded its watchdog".into()); - 5872
return partial; - 5873
} - 5874
}; - 5875
- 5876
// --- fold in ------------------------------------------------------------ - 5877
let classifications = match vak_intent::parse_classifications(&answer) { - 5878
Ok(parsed) => parsed, - 5879
Err(error) => { - 5880
give_up( - 5881
&mut partial, - 5882
format!("unparseable classifier answer: {error}"), - 5883
); - 5884
return partial; - 5885
} - 5886
}; - 5887
// The tier names where the model actually ran, not which setting - 5888
// asked for it: `escalate = "cloud"` on an Ollama route is local. - 5889
let ran_off_machine = provider_name != "ollama"; - 5890
let mut applied = vak_intent::apply_classification( - 5891
partial, - 5892
&classifications, - 5893
&format!("{provider_name}/{model}"), - 5894
&digest, - 5895
&self.turn_authority(), - 5896
&intent::resolver_config(&self.inner.config), - 5897
ran_off_machine, - 5898
); - 5899
if !matches!( - 5900
applied.provenance.tier, - 5901
vak_intent::Tier::LocalModel | vak_intent::Tier::CloudModel - 5902
) { - 5903
// The answer parsed but set nothing: keep the start of it so the - 5904
// ledger says what the classifier actually said. - 5905
const NOTED_ANSWER_CHARS: usize = 400; - 5906
let shown: String = answer.chars().take(NOTED_ANSWER_CHARS).collect(); - 5907
applied.provenance.escalation_note = Some(format!( - 5908
"{}; answer: {shown:?}", - 5909
applied - 5910
.provenance - 5911
.escalation_note - 5912
.clone() - 5913
.unwrap_or_default() - 5914
)); - 5915
} - 5916
applied - 5917
} - 5918
- 5919
/// Takes `self` by value (an `Arc` bump plus a few small per-turn - 5920
/// fields) so the single reconciliation below can correct this turn's - 5921
/// answerability before anything reads it. Every public `run_*_with` - 5922
/// funnels through here, which is what makes that one place enough. - 5923
#[allow(clippy::too_many_arguments)] - 5924
async fn run_turn_inner( - 5925
mut self, - 5926
mut session: SessionLog, - 5927
prompt: vak_session::MessageRecord, - 5928
cancel: CancellationToken, - 5929
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5930
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5931
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5932
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5933
goal: Option<(String, Vec<String>)>, - 5934
work_mode: Option<WorkMode>, - 5935
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5936
let vak_session::MessageRecord { - 5937
message: prompt, - 5938
meta: prompt_meta, - 5939
} = prompt; - 5940
let prompt_text = prompt.text_content(); - 5941
self.agent_identity = session.header().and_then(|header| header.agent.clone()); - 5942
// The approver that will actually serve this run is the authority on - 5943
// whether its gates reach anyone. Whatever the host stamped earlier - 5944
// loses to it, and a disagreement is recorded rather than believed. - 5945
self.reconcile_answerability(approver.as_ref()); - 5946
let permission_lease = self.permission_lease(); - 5947
let live_child_sessions = session - 5948
.header() - 5949
.map(|header| { - 5950
self.inner - 5951
.workers - 5952
.active_for(&header.session_id) - 5953
.into_iter() - 5954
.map(|child| child.id) - 5955
.collect::<std::collections::HashSet<_>>() - 5956
}) - 5957
.unwrap_or_default(); - 5958
session.reconcile_running_work_with_child_ledgers( - 5959
&live_child_sessions, - 5960
Some(&self.inner.sessions_home), - 5961
)?; - 5962
let session_contract = session.header().map(|header| header.contract.clone()); - 5963
// Per-turn routing: provider/model are always resolved from the live - 5964
// effective_route(), never from the session's FrozenContract. The - 5965
// contract is authority only for capabilities and permission_mode - 5966
// (security/audit boundaries). This means: - 5967
// - A user changing provider in Settings takes effect on the next turn - 5968
// of any open session, not just new sessions. - 5969
// - The auto-routing algorithm (evidence, beliefs, v2 ordering) is - 5970
// re-evaluated every turn, not frozen at admission. - 5971
// - Workers inherit the Core's current effective route, not the - 5972
// parent session's admission snapshot. - 5973
// - Per-turn dispatch is recorded in WorkReceipt; audit is preserved. - 5974
let (provider, model) = (self.provider()?, self.effective_model()); - 5975
let registry = self.capability_registry(); - 5976
if registry.current().await.epoch == 0 || registry.has_pending_changes().await { - 5977
registry.reconcile().await; - 5978
} - 5979
// Rendered from the re-bound packet, not from the frozen string. - 5980
// - 5981
// The contract's admitted set is still the authority; only the - 5982
// rendering of what it already admits is refreshed. See - 5983
// `rebound_capabilities`. - 5984
let frozen_system_prompt = match session_contract.as_ref() { - 5985
Some(contract) => { - 5986
let rebound = self.rebound_capabilities(contract).await; - 5987
self.resolve_prompt(&rebound).text - 5988
} - 5989
None => { - 5990
let admitted = self.admitted_capabilities().await; - 5991
self.resolve_prompt(&admitted).text - 5992
} - 5993
}; - 5994
let mut cfg = AgentConfig::new(frozen_system_prompt.clone()); - 5995
- 5996
// ---- intent resolution (docs/design/47-commitment-kernel.md) ---- - 5997
// Runs before anything reads a knob it governs. Everything derived - 5998
// from it narrows: the projections in `crate::intent` take a baseline - 5999
// and return something no wider, so a misread can make this turn more - 6000
// cautious and never less. - 6001
// - 6002
// The turn id is minted here, once: every strand and thread id this - 6003
// turn records derives from it, so threads — and the commitments - 6004
// keyed by them — are unique across turns and sessions. - 6005
let turn_id = uuid_like(); - 6006
let resolved_intent = self - 6007
.resolve_turn_intent_with_escalation(&mut session, &prompt, &turn_id, &cancel) - 6008
.await; - 6009
// Which commitments this turn works on is decided now, before any - 6010
// knob is read, because a grant on one of them narrows the strands - 6011
// that serve it. Nothing is written until the turn is about to run. - 6012
let turn_authority = self.turn_authority(); - 6013
let episode_plan = commitments::plan_episodes( - 6014
&self.sessions_home(), - 6015
&self.inner.config, - 6016
&resolved_intent, - 6017
chrono::Utc::now(), - 6018
); - 6019
let envelopes = episode_plan.envelopes(); - 6020
let resolved_intent = if envelopes.is_empty() { - 6021
resolved_intent - 6022
} else { - 6023
vak_intent::apply_envelopes( - 6024
resolved_intent, - 6025
&envelopes, - 6026
turn_authority.autonomy, - 6027
chrono::Utc::now(), - 6028
) - 6029
}; - 6030
let engagement = resolved_intent.engagement.clone(); - 6031
debug_assert!( - 6032
intent::projection_is_narrowing(&engagement.limits), - 6033
"a derived engagement widened the baseline" - 6034
); - 6035
let mut admitted_outcome = - 6036
vak_intent::OutcomeSpec::from_intent(prompt.text_content(), &resolved_intent); - 6037
admitted_outcome.evidence_max_age_secs = Some(self.effective_evidence_max_age_secs()); - 6038
cfg.continued_saved_file = - 6039
continued_saved_file(&session, &resolved_intent, &self.inner.cwd); - 6040
cfg.outcome = Some(admitted_outcome.clone()); - 6041
- 6042
// ---- turn-capability assembly (docs/design/41-capability-registry.md § Turn) ---- - 6043
// Admission is policy only — channel, reach, revocation — and is the - 6044
// same for every kind. What the reading predicts the turn will need - 6045
// decides only which admitted tools are *loaded* (the surface, below); - 6046
// a misread therefore costs one `find_tools` call, never a capability. - 6047
let cap_set = registry.current().await; - 6048
let revoked_ids = registry.revoked_ids().await; - 6049
let reach_standings = self.capability_standings(); - 6050
let channel_policy = self.channel_policy().unwrap_or_default(); - 6051
let mcp_inventory = cap_set.mcp_inventory(); - 6052
let builtin_names: std::collections::BTreeSet<String> = - 6053
self.tool_names().into_iter().collect(); - 6054
let mut turn_capabilities = capability::TurnCapabilities::build(&capability::TurnProbe { - 6055
capabilities: cap_set.as_ref(), - 6056
revoked_ids, - 6057
channel_policy: &channel_policy, - 6058
reach_standings: &reach_standings, - 6059
mcp_inventory: &mcp_inventory, - 6060
builtin_names: &builtin_names, - 6061
}); - 6062
let revoke_registry = registry.clone(); - 6063
{ - 6064
// Grounding: which calls reach outside information is decided from - 6065
// what each capability declares it serves (`serves`), never from - 6066
// a tool's name or its output. - 6067
let servers: std::collections::BTreeMap<String, Vec<String>> = self - 6068
.effective_mcp() - 6069
.servers - 6070
.iter() - 6071
.map(|(name, server)| (name.clone(), server.serves.clone())) - 6072
.collect(); - 6073
let index = cfg.mcp_tool_index.clone(); - 6074
let tool_serves: std::collections::BTreeMap<String, Vec<String>> = self - 6075
.tool_declarations() - 6076
.into_iter() - 6077
.map(|(name, serves)| (name, serves.iter().map(|d| d.to_string()).collect())) - 6078
.collect(); - 6079
let observing: std::collections::BTreeSet<String> = tool_serves - 6080
.iter() - 6081
.filter(|(_, serves)| capability::provider::serves_observation(serves)) - 6082
.map(|(name, _)| name.clone()) - 6083
.collect(); - 6084
cfg.observation_check = Some(Arc::new(move |name, _input| observing.contains(name))); - 6085
cfg.retrieval_check = Some(Arc::new(move |name, input| { - 6086
let server = index - 6087
.lock() - 6088
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6089
.get(name) - 6090
.cloned(); - 6091
capability::provider::call_retrieves_external( - 6092
name, - 6093
input, - 6094
server.as_deref(), - 6095
&|tool| tool_serves.get(tool).cloned().unwrap_or_default(), - 6096
&|server| servers.get(server).cloned().unwrap_or_default(), - 6097
) - 6098
})); - 6099
} - 6100
{ - 6101
let recipes = vak_delivery::built_in_recipes(); - 6102
cfg.presentation_check = Some(Arc::new(move |text, offered| { - 6103
presentation_tools::presentation_check_nudge(text, offered, &recipes) - 6104
})); - 6105
} - 6106
{ - 6107
// Presentations are ledger entries (docs/design/68-context-engine.md - 6108
// §10): the agent loop has no card-shape or skill-registry - 6109
// knowledge (AGENTS.md invariant 14 — the tool itself executes - 6110
// across the worker/broker boundary and has no session-log - 6111
// access either), so `Core` supplies the rebuild as a closure, - 6112
// the same pattern `presentation_check`/`retrieval_check` use. - 6113
let skills = vak_delivery::built_in_skill_registry(); - 6114
cfg.presentation_rebuild = Some(Arc::new(move |name, input| { - 6115
// A tool call already passes its own tool name; a fence - 6116
// (docs/design/68-context-engine.md §10's inline-fence - 6117
// fallback) has no tool name at all — vak-agent has no - 6118
// card-shape knowledge (invariant 14) and cannot resolve - 6119
// `semantic_type` to the matching `emit_*_card` tool - 6120
// itself, so this hook does it here, preferring the - 6121
// payload's own declared type and falling back to - 6122
// whatever name the caller passed. - 6123
let resolved_name = input - 6124
.get("semantic_type") - 6125
.and_then(serde_json::Value::as_str) - 6126
.and_then(presentation_tools::emit_tool_for) - 6127
.unwrap_or(name); - 6128
presentation_tools::presentation_info(resolved_name, input, &skills).map(|info| { - 6129
vak_tools::PresentationCard { - 6130
semantic_type: info.semantic_type, - 6131
skill_id: info.skill_id, - 6132
skill_version: info.skill_version, - 6133
schema_version: info.schema_version, - 6134
payload: info.payload, - 6135
title: info.title, - 6136
identity_digest: info.identity_digest, - 6137
} - 6138
}) - 6139
})); - 6140
} - 6141
cfg.revocation_check = Some(Arc::new(move |name, input| { - 6142
let id = if name == "mcp" { - 6143
input - 6144
.get("server") - 6145
.and_then(serde_json::Value::as_str) - 6146
.map(|server| capability::CapabilityId::new(CapabilityKind::McpServer, server)) - 6147
} else if name == "skill" { - 6148
input - 6149
.get("name") - 6150
.and_then(serde_json::Value::as_str) - 6151
.map(|skill| capability::CapabilityId::new(CapabilityKind::Skill, skill)) - 6152
} else { - 6153
Some(capability::CapabilityId::new(CapabilityKind::Tool, name)) - 6154
}; - 6155
id.is_some_and(|id| revoke_registry.revoked_now(&id)) - 6156
})); - 6157
// Prompt text is a projection of the same selected descriptor set as - 6158
// the schemas. This is deliberately after intent/policy filtering so - 6159
// a removed or withheld capability cannot remain in prose. - 6160
// - 6161
// The prefix (identity, contract, guardrails, tool surface) is - 6162
// byte-stable and carried as `system_prefix`; the clock instant and - 6163
// epistemic stance are per-turn and carried separately as `tail` so - 6164
// the request assembler can render them into the moving tail instead - 6165
// of the cached prefix (docs/design/68-context-engine.md §4/§6). - 6166
// Progressive disclosure: with it on, the reading decides which admitted - 6167
// tools are loaded and the rest are listed in the stable catalogue; - 6168
// with it off, everything is loaded and there is nothing to list. - 6169
let progressive = - 6170
self.inner.config.intent.enabled && self.inner.config.intent.slice_capabilities; - 6171
let loaded_domains = if progressive { - 6172
engagement.limits.required_domains.clone() - 6173
} else { - 6174
vak_intent::DomainSet::All - 6175
}; - 6176
let tool_catalogue = if progressive { - 6177
self.tool_catalogue_for(&turn_capabilities.tool_names) - 6178
} else { - 6179
String::new() - 6180
}; - 6181
let (turn_resolution, turn_temporal, turn_stance) = self.resolve_prompt_with_stance_parts( - 6182
&turn_capabilities.descriptors, - 6183
Some(engagement.posture.epistemic_stance), - 6184
&tool_catalogue, - 6185
); - 6186
cfg.system_prefix = turn_resolution.text; - 6187
cfg.tail = vak_agent::TailInput { - 6188
temporal: turn_temporal.trim().to_string(), - 6189
stance: prompts::stance_with_card_clarifier(&turn_stance), - 6190
}; - 6191
- 6192
let work_config = self.effective_work(); - 6193
// Managed-ness follows from the reading's horizon rather than from a - 6194
// keyword scan. The old `is_managed_work_request` fired on any two of - 6195
// `and`/`then`/`first`, so "explain what this and that mean" read as - 6196
// durable multi-step work; `Horizon` is derived from recurrence and - 6197
// enumeration instead. An explicit run-scoped mode still wins. - 6198
let default_work_mode = match work_config.default_mode.as_str() { - 6199
"managed" => WorkMode::Managed, - 6200
"auto" if engagement.posture.managed => WorkMode::Managed, - 6201
_ => WorkMode::Direct, - 6202
}; - 6203
cfg.work_mode = work_mode.unwrap_or(default_work_mode); - 6204
if cfg.work_mode == WorkMode::Auto { - 6205
cfg.work_mode = if engagement.posture.managed { - 6206
WorkMode::Managed - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.