- 5270
!route_cfg.fallback_models.is_empty(), - 5271
); - 5272
routing::RoutePlan { - 5273
ladder, - 5274
objective: objective.as_str().to_string(), - 5275
annotations, - 5276
} - 5277
} - 5278
- 5279
pub async fn start_session(&self) -> Result<SessionLog, CoreError> { - 5280
let route = self.refresh_persisted_route()?; - 5281
self.start_session_with_route(route.provider, route.model) - 5282
.await - 5283
} - 5284
- 5285
/// Create a frozen session from an explicit provider/model pair without - 5286
/// mutating the Core default. Gateway/channel overrides use this path so - 5287
/// concurrent surfaces cannot overwrite one another's admission route. - 5288
pub async fn start_session_with_route( - 5289
&self, - 5290
provider: String, - 5291
model: String, - 5292
) -> Result<SessionLog, CoreError> { - 5293
self.start_session_with_route_for(provider, model, None) - 5294
.await - 5295
} - 5296
- 5297
/// Open a session whose route ladder is ordered for a known first request. - 5298
/// - 5299
/// Surfaces that have the opening prompt in hand should use this: the - 5300
/// ladder is frozen once at admission, so the demand facts available at - 5301
/// that moment are the only ones that can ever influence its ordering. - 5302
pub async fn start_session_for_prompt( - 5303
&self, - 5304
provider: String, - 5305
model: String, - 5306
prompt: &str, - 5307
) -> Result<SessionLog, CoreError> { - 5308
let resolution = intent::resolve_turn( - 5309
prompt, - 5310
"", - 5311
self.surface(), - 5312
&[], - 5313
intent::workspace_facts(&self.inner.cwd), - 5314
vak_intent::HistoryFacts::default(), - 5315
&vak_intent::Declared::default(), - 5316
&self.turn_authority(), - 5317
&intent::resolver_config(&self.inner.config), - 5318
); - 5319
let demand = resolution.peek().engagement.posture.demand; - 5320
self.start_session_with_route_for(provider, model, Some(demand)) - 5321
.await - 5322
} - 5323
- 5324
async fn start_session_with_route_for( - 5325
&self, - 5326
provider: String, - 5327
model: String, - 5328
demand: Option<vak_intent::DemandHint>, - 5329
) -> Result<SessionLog, CoreError> { - 5330
let session_id = uuid_like(); - 5331
let path = vak_session::SessionPath::new_session_file( - 5332
&self.sessions_home(), - 5333
&self.inner.cwd, - 5334
&session_id, - 5335
); - 5336
let primary_credential_id = self - 5337
.provider_auth_for_leg(&provider, None) - 5338
.ok() - 5339
.and_then(|auth| auth.credential_id); - 5340
let needs_tools_or_reasoning = - 5341
demand.is_some_and(|hint| hint.reasoning_required) || !self.tool_names().is_empty(); - 5342
let plan = self.plan_route_ladder( - 5343
vak_llm::RouteLeg { - 5344
provider: provider.clone(), - 5345
model: model.clone(), - 5346
dialect: vak_llm::EndpointDialect::for_provider( - 5347
&provider, - 5348
needs_tools_or_reasoning, - 5349
), - 5350
credential_id: primary_credential_id, - 5351
}, - 5352
demand, - 5353
); - 5354
let capabilities = self.admitted_capabilities().await; - 5355
let resolution = self.resolve_prompt(&capabilities); - 5356
let system_prompt = resolution.text; - 5357
let conversation = self.conversation_context.clone().or_else(|| { - 5358
Some(vak_session::ConversationContext::local( - 5359
&session_id, - 5360
self.surface.slug(), - 5361
)) - 5362
}); - 5363
let header = SessionHeader { - 5364
agent: self.agent_identity.clone(), - 5365
session_id, - 5366
created_at: chrono::Utc::now(), - 5367
cwd: self.inner.cwd.clone(), - 5368
parent_session_id: None, - 5369
contract_id: None, - 5370
work_item_id: None, - 5371
conversation, - 5372
contract: FrozenContract { - 5373
app_version: APP_VERSION.into(), - 5374
provider, - 5375
model, - 5376
// Frozen-ladder admission (docs/design/15-reliability.md + - 5377
// Phase R): primary leg always first; additional legs - 5378
// ONLY from warm discovery caches -- the same model on - 5379
// other keyed providers, plus explicit `[route]` - 5380
// fallback_models when warm discovery reaches them. - 5381
// No invented ids, no network at admission. Ordered by - 5382
// demand-scored v2 over TTL-filtered evidence and - 5383
// session beliefs, diversity-capped, then frozen. - 5384
route_ladder: plan.ladder, - 5385
route_objective: plan.objective, - 5386
route_annotations: plan.annotations, - 5387
system_prompt, - 5388
permission_mode: format!("{:?}", self.effective_permission_mode()) - 5389
.to_kebab_lowercase(), - 5390
capabilities, - 5391
prompt_layers: resolution.descriptors, - 5392
}, - 5393
}; - 5394
Ok(SessionLog::create(path, header)?) - 5395
} - 5396
- 5397
pub async fn run_turn( - 5398
&self, - 5399
session: SessionLog, - 5400
prompt: &str, - 5401
cancel: CancellationToken, - 5402
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5403
) -> Result<TurnOutcome, CoreError> { - 5404
let (outcome, _session) = self - 5405
.run_turn_with(session, prompt, cancel, None, None, None, events) - 5406
.await?; - 5407
Ok(outcome) - 5408
} - 5409
- 5410
#[allow(clippy::too_many_arguments)] - 5411
pub async fn run_turn_with_message( - 5412
&self, - 5413
session: SessionLog, - 5414
prompt: vak_session::MessageRecord, - 5415
cancel: CancellationToken, - 5416
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5417
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5418
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5419
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5420
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5421
self.clone() - 5422
.run_turn_inner( - 5423
session, prompt, cancel, approver, permission, steering, events, None, None, - 5424
) - 5425
.await - 5426
} - 5427
- 5428
#[allow(clippy::too_many_arguments)] - 5429
pub async fn run_turn_with( - 5430
&self, - 5431
session: SessionLog, - 5432
prompt: &str, - 5433
cancel: CancellationToken, - 5434
approver: Option<std::sync::Arc<dyn vak_agent::Approver>>, - 5435
permission: Option<std::sync::Arc<vak_permission::PermissionEngine>>, - 5436
steering: Option<std::sync::Arc<vak_agent::SteeringQueues>>, - 5437
events: tokio::sync::mpsc::Sender<AgentEvent>, - 5438
) -> Result<(TurnOutcome, SessionLog), CoreError> { - 5439
self.clone() - 5440
.run_turn_inner( - 5441
session, - 5442
vak_session::MessageRecord { - 5443
message: vak_llm::Message::user_text(prompt), - 5444
meta: None, - 5445
}, - 5446
cancel, - 5447
approver, - 5448
permission, - 5449
steering, - 5450
events, - 5451
None, - 5452
None, - 5453
) - 5454
.await - 5455
} - 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).
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.