- 4812
} - 4813
if self - 4814
.inner - 4815
.capacity_probes - 4816
.lock() - 4817
.is_ok_and(|probes| probes.contains_key(&key)) - 4818
{ - 4819
return; - 4820
} - 4821
if self - 4822
.inner - 4823
.capacity_probe_attempted - 4824
.lock() - 4825
.ok() - 4826
.and_then(|attempts| attempts.get(&key).copied()) - 4827
.is_some_and(|at| at.elapsed() < Self::CAPACITY_PROBE_MIN_INTERVAL) - 4828
{ - 4829
return; - 4830
} - 4831
let Ok(auth) = self.provider_auth_for_leg(&leg.provider, leg.credential_id.as_deref()) - 4832
else { - 4833
return; - 4834
}; - 4835
let Ok(provider_client) = self.inner.registry.get(&adapter_name_for_leg(leg), &auth) else { - 4836
return; - 4837
}; - 4838
- 4839
let token = CancellationToken::new(); - 4840
if let Ok(mut probes) = self.inner.capacity_probes.lock() { - 4841
probes.insert(key.clone(), token.clone()); - 4842
} - 4843
if let Ok(mut attempts) = self.inner.capacity_probe_attempted.lock() { - 4844
attempts.insert(key.clone(), std::time::Instant::now()); - 4845
} - 4846
- 4847
let core = self.clone(); - 4848
let leg = leg.clone(); - 4849
let probe_key = key.clone(); - 4850
tokio::spawn(async move { - 4851
let mut profile = core - 4852
.run_capacity_probe( - 4853
&leg, - 4854
provider_client, - 4855
ProbeMetadata { - 4856
declared_window, - 4857
output_reserve, - 4858
metadata_digest, - 4859
probed_at: now, - 4860
}, - 4861
&token, - 4862
) - 4863
.await; - 4864
profile.provenance.quantisation = probe_key.quantisation.clone(); - 4865
if let Ok(mut cache) = core.inner.capacity_cache.lock() { - 4866
cache.insert(probe_key.clone(), profile); - 4867
} - 4868
if let Ok(mut probes) = core.inner.capacity_probes.lock() { - 4869
probes.remove(&probe_key); - 4870
} - 4871
}); - 4872
} - 4873
- 4874
/// Runs the horizon-ladder probe (docs/design/68 §1 "Horizon ladder") - 4875
/// against a live provider client, one rung per request, honoring - 4876
/// `cancel`. A rung the provider rejects with a context-length error - 4877
/// (`LlmError::Context`) counts as a failed rung; any other transport - 4878
/// error stops the ladder early rather than fabricating more rungs, and - 4879
/// the ladder's result-so-far becomes the profile. - 4880
async fn run_capacity_probe( - 4881
&self, - 4882
leg: &vak_llm::RouteLeg, - 4883
provider_client: Arc<dyn Provider>, - 4884
metadata: ProbeMetadata, - 4885
cancel: &CancellationToken, - 4886
) -> vak_context::capacity::CapacityProfile { - 4887
let ProbeMetadata { - 4888
declared_window, - 4889
output_reserve, - 4890
metadata_digest, - 4891
probed_at, - 4892
} = metadata; - 4893
let mut ladder = vak_context::capacity::Ladder::new(declared_window); - 4894
let mut rungs: Vec<vak_context::capacity::Rung> = Vec::new(); - 4895
let mut signals: Vec<String> = Vec::new(); - 4896
// Refined as soon as one rung reports real usage, so later rungs - 4897
// land closer to their token target on a real tokenizer. - 4898
let mut tokens_per_char_hint = 0.25_f64; - 4899
// Refined from each rung's measured prefill time; 0 means unknown - 4900
// (no rung has reported one yet), in which case no rung is skipped. - 4901
let mut prefill_tps_hint = 0.0_f64; - 4902
- 4903
while let Some(target) = ladder.next_rung() { - 4904
if cancel.is_cancelled() { - 4905
signals.push("probe cancelled before convergence".into()); - 4906
break; - 4907
} - 4908
if prefill_tps_hint > 0.0 - 4909
&& target as f64 / prefill_tps_hint > Self::PROBE_MAX_RUNG_SECS - 4910
{ - 4911
signals.push(format!( - 4912
"rung {target} skipped: expected prefill ~{:.0}s exceeds the {:.0}s bound", - 4913
target as f64 / prefill_tps_hint, - 4914
Self::PROBE_MAX_RUNG_SECS - 4915
)); - 4916
break; - 4917
} - 4918
let request = - 4919
vak_context::capacity::probe_request(target, tokens_per_char_hint, &leg.model); - 4920
let sent_chars: u64 = request - 4921
.messages - 4922
.iter() - 4923
.map(|m| m.text_content().len() as u64) - 4924
.sum(); - 4925
// One completion from a sampling model is one coin flip; a rung - 4926
// is decided by the majority of up to PROBE_SAMPLES_PER_RUNG - 4927
// identical requests (identical on purpose: the prefix is - 4928
// cached after the first, so the extra samples cost decode - 4929
// time only), stopping as soon as the majority is settled. - 4930
let mut passes = 0u32; - 4931
let mut fails = 0u32; - 4932
let mut rejected: Option<String> = None; - 4933
let mut transport_error: Option<vak_llm::LlmError> = None; - 4934
let mut first_prefill_ms: Option<u64> = None; - 4935
let needed = vak_context::capacity::PROBE_SAMPLES_PER_RUNG / 2 + 1; - 4936
while passes < needed && fails < needed && rejected.is_none() { - 4937
if cancel.is_cancelled() { - 4938
break; - 4939
} - 4940
let outcome = match tokio::time::timeout( - 4941
Self::PROBE_REQUEST_TIMEOUT, - 4942
provider_client.stream(request.clone(), cancel.clone()), - 4943
) - 4944
.await - 4945
{ - 4946
Ok(Ok(stream)) => stream.result().await, - 4947
Ok(Err(e)) => Err(e), - 4948
Err(_) => Err(vak_llm::LlmError::Network(format!( - 4949
"probe rung {target} exceeded its {:?} timeout", - 4950
Self::PROBE_REQUEST_TIMEOUT - 4951
))), - 4952
}; - 4953
match outcome { - 4954
Ok(message) => { - 4955
if sent_chars > 0 && message.usage.input_tokens > 0 { - 4956
tokens_per_char_hint = - 4957
message.usage.prompt_tokens() as f64 / sent_chars as f64; - 4958
} - 4959
if first_prefill_ms.is_none() { - 4960
first_prefill_ms = message.usage.prefill_ms; - 4961
} - 4962
if let Some(ms) = message.usage.prefill_ms - 4963
&& ms > 0 - 4964
{ - 4965
prefill_tps_hint = - 4966
message.usage.input_tokens as f64 / (ms as f64 / 1000.0); - 4967
} - 4968
if vak_context::capacity::followed(&message) { - 4969
passes += 1; - 4970
} else { - 4971
fails += 1; - 4972
} - 4973
} - 4974
Err(vak_llm::LlmError::Context(msg)) => rejected = Some(msg), - 4975
Err(e) => { - 4976
transport_error = Some(e); - 4977
break; - 4978
} - 4979
} - 4980
} - 4981
if let Some(e) = transport_error { - 4982
signals.push(format!("rung {target} probe failed: {e}")); - 4983
break; - 4984
} - 4985
if let Some(msg) = rejected { - 4986
signals.push(format!("rung {target} rejected: {msg}")); - 4987
rungs.push(vak_context::capacity::Rung { - 4988
tokens: target, - 4989
accepted: false, - 4990
followed_instruction: None, - 4991
prefill_ms: None, - 4992
}); - 4993
ladder.report(target, false, false); - 4994
continue; - 4995
} - 4996
if passes + fails == 0 { - 4997
break; - 4998
} - 4999
let followed = passes >= needed; - 5000
signals.push(format!( - 5001
"rung {target}: {passes} followed / {fails} did not" - 5002
)); - 5003
rungs.push(vak_context::capacity::Rung { - 5004
tokens: target, - 5005
accepted: true, - 5006
followed_instruction: Some(followed), - 5007
prefill_ms: first_prefill_ms, - 5008
}); - 5009
ladder.report(target, true, followed); - 5010
} - 5011
- 5012
let verified_window = ladder.verified_window(); - 5013
let horizon = ladder.result().unwrap_or_else(|| { - 5014
let largest_followed = rungs - 5015
.iter() - 5016
.filter(|r| r.followed_instruction == Some(true)) - 5017
.map(|r| r.tokens) - 5018
.max(); - 5019
vak_context::capacity::Horizon { - 5020
tokens: largest_followed.unwrap_or_else(|| declared_window.min(4_000)), - 5021
confidence: if largest_followed.is_some() { 0.5 } else { 0.3 }, - 5022
last_confirmed: probed_at, - 5023
} - 5024
}); - 5025
- 5026
// Cache rung (docs/design/68-context-engine.md §1 point 3): two - 5027
// identical requests sent back to back at a fixed, modest size — - 5028
// separate from the horizon ladder, which varies size to find the - 5029
// instruction-following boundary rather than to probe caching. - 5030
let cache = if cancel.is_cancelled() { - 5031
vak_context::capacity::CacheBehaviour::Unknown - 5032
} else { - 5033
let cache_request = - 5034
vak_context::capacity::probe_request(4_000, tokens_per_char_hint, &leg.model); - 5035
let (first_outcome, first_latency_ms) = Self::stream_with_first_token_latency( - 5036
&provider_client, - 5037
cache_request.clone(), - 5038
cancel, - 5039
) - 5040
.await; - 5041
match first_outcome { - 5042
Ok(_) => { - 5043
let (second_outcome, second_latency_ms) = - 5044
Self::stream_with_first_token_latency( - 5045
&provider_client, - 5046
cache_request, - 5047
cancel, - 5048
) - 5049
.await; - 5050
match second_outcome { - 5051
Ok(second_message) => { - 5052
let first_ms = first_latency_ms.unwrap_or(0); - 5053
let second_ms = second_latency_ms.unwrap_or(0); - 5054
let behaviour = vak_context::capacity::classify_cache_rung( - 5055
first_ms, - 5056
second_ms, - 5057
&second_message.usage, - 5058
); - 5059
signals.push(format!( - 5060
"cache rung: first={first_ms}ms second={second_ms}ms \ - 5061
cache_read_input_tokens={:?} -> {behaviour:?}", - 5062
second_message.usage.cache_read_input_tokens - 5063
)); - 5064
behaviour - 5065
} - 5066
Err(e) => { - 5067
signals.push(format!("cache rung second request failed: {e}")); - 5068
vak_context::capacity::CacheBehaviour::Unknown - 5069
} - 5070
} - 5071
} - 5072
Err(e) => { - 5073
signals.push(format!("cache rung first request failed: {e}")); - 5074
vak_context::capacity::CacheBehaviour::Unknown - 5075
} - 5076
} - 5077
}; - 5078
- 5079
vak_context::capacity::CapacityProfile::from_probe( - 5080
declared_window, - 5081
verified_window, - 5082
horizon, - 5083
cache, - 5084
output_reserve, - 5085
vak_context::capacity::ProbeProvenance { - 5086
probed_at, - 5087
rungs, - 5088
signals, - 5089
metadata_digest, - 5090
quantisation: None, - 5091
}, - 5092
) - 5093
} - 5094
- 5095
/// Streams `request` and reports the wall-clock time from just before - 5096
/// the request is sent to the first `StreamEvent` off the wire, - 5097
/// alongside the final outcome. Used by the horizon ladder's cache rung - 5098
/// (docs/design/68-context-engine.md §1) to measure a provider's - 5099
/// prefix-cache behaviour purely from timing when it reports nothing. - 5100
async fn stream_with_first_token_latency( - 5101
provider_client: &Arc<dyn Provider>, - 5102
request: vak_llm::ChatRequest, - 5103
cancel: &CancellationToken, - 5104
) -> ( - 5105
Result<vak_llm::AssistantMessage, vak_llm::LlmError>, - 5106
Option<u64>, - 5107
) { - 5108
let started = std::time::Instant::now(); - 5109
match provider_client.stream(request, cancel.clone()).await { - 5110
Ok(mut stream) => { - 5111
let mut first_ms = None; - 5112
while let Some(_event) = futures::StreamExt::next(&mut stream).await { - 5113
if first_ms.is_none() { - 5114
first_ms = Some(started.elapsed().as_millis() as u64); - 5115
} - 5116
} - 5117
(stream.result().await, first_ms) - 5118
} - 5119
Err(e) => (Err(e), None), - 5120
} - 5121
} - 5122
- 5123
/// Drop memoised discovery for `provider` (or all of it) so the next - 5124
/// read reflects a key that just changed. - 5125
pub fn invalidate_models_cache(&self, provider: Option<&str>) { - 5126
if let Ok(mut cache) = self.inner.models_cache.lock() { - 5127
match provider { - 5128
Some(p) => cache.retain(|(cached_provider, _), _| cached_provider != p), - 5129
None => cache.clear(), - 5130
} - 5131
} - 5132
if let Ok(mut cache) = self.inner.model_context_cache.lock() { - 5133
match provider { - 5134
Some(p) => cache.retain(|(provider, _, _), _| provider != p), - 5135
None => cache.clear(), - 5136
} - 5137
} - 5138
} - 5139
- 5140
/// Frozen-ladder admission (docs/design/15-reliability.md + Phase R). - 5141
/// - 5142
/// Pure with respect to its inputs: warm discovery caches, the - 5143
/// evidence ledger, session beliefs, config, and tool count. No - 5144
/// network, no invented model ids. The operator-selected primary is - 5145
/// pinned to the head; v2 ordering decides only the FALLBACK order. - 5146
/// Order the frozen ladder for a session. - 5147
/// - 5148
/// `demand` is what the first turn's reading concluded about this work. - 5149
/// It is optional because a session can be opened before anyone has said - 5150
/// what it is for; when absent the demand facts fall back to the - 5151
/// conservative defaults below rather than being fabricated. - 5152
fn plan_route_ladder( - 5153
&self, - 5154
primary: vak_llm::RouteLeg, - 5155
demand: Option<vak_intent::DemandHint>, - 5156
) -> routing::RoutePlan { - 5157
const DISCOVERY_TTL: std::time::Duration = std::time::Duration::from_secs(300); - 5158
let mut candidates = vec![primary.clone()]; - 5159
- 5160
// Same-model legs on other keyed providers (legacy Phase B set). - 5161
if let Ok(cache) = self.inner.models_cache.lock() { - 5162
for ((p, credential_id), (fetched_at, models)) in cache.iter() { - 5163
if fetched_at.elapsed() >= DISCOVERY_TTL { - 5164
continue; - 5165
} - 5166
if models.contains(&primary.model) - 5167
&& !candidates.iter().any(|c| { - 5168
c.provider == *p && c.credential_id.as_deref() == Some(credential_id) - 5169
}) - 5170
&& self.provider_auth_for_leg(p, Some(credential_id)).is_ok() - 5171
{ - 5172
candidates.push(vak_llm::RouteLeg { - 5173
provider: p.clone(), - 5174
model: primary.model.clone(), - 5175
dialect: vak_llm::EndpointDialect::for_provider( - 5176
p, - 5177
!self.tool_names().is_empty(), - 5178
), - 5179
credential_id: Some(credential_id.clone()), - 5180
}); - 5181
} - 5182
} - 5183
} - 5184
- 5185
// Phase R cross-model legs: ONLY exact ids from the explicit - 5186
// `[route].fallback_models` allowlist, admitted when warm - 5187
// discovery shows a configured key reaches them. - 5188
let route_cfg = &self.inner.config.route; - 5189
if !route_cfg.fallback_models.is_empty() - 5190
&& let Ok(cache) = self.inner.models_cache.lock() - 5191
{ - 5192
for ((p, credential_id), (fetched_at, models)) in cache.iter() { - 5193
if fetched_at.elapsed() >= DISCOVERY_TTL { - 5194
continue; - 5195
} - 5196
if self.provider_auth_for_leg(p, Some(credential_id)).is_err() { - 5197
continue; - 5198
} - 5199
for m in models { - 5200
if route_cfg.fallback_models.contains(m) - 5201
&& m != &primary.model - 5202
&& !candidates.iter().any(|c| { - 5203
c.provider == *p - 5204
&& c.model == *m - 5205
&& c.credential_id.as_deref() == Some(credential_id) - 5206
}) - 5207
{ - 5208
candidates.push(vak_llm::RouteLeg { - 5209
provider: p.clone(), - 5210
model: m.clone(), - 5211
dialect: vak_llm::EndpointDialect::for_provider( - 5212
p, - 5213
!self.tool_names().is_empty(), - 5214
), - 5215
credential_id: Some(credential_id.clone()), - 5216
}); - 5217
} - 5218
} - 5219
} - 5220
} - 5221
candidates.sort(); - 5222
candidates.dedup(); - 5223
- 5224
// Demand scoring from facts available at admission. Unknown context - 5225
// still reads as moderate -- never zero, never fabricated -- but the - 5226
// three behavioural facts now come from the turn's reading instead of - 5227
// being hardcoded `false`. Passing constants here is why every session - 5228
// scored identical demand and the objective was effectively fixed, - 5229
// leaving the ordering function inert. - 5230
let hint = demand.unwrap_or(vak_intent::DemandHint { - 5231
reasoning_required: false, - 5232
evidence_required: false, - 5233
structured_output: false, - 5234
}); - 5235
let demand = vak_llm::score_demand(vak_llm::DemandInput { - 5236
estimated_input_tokens: 0, - 5237
output_budget_tokens: u64::from(self.inner.config.max_tokens), - 5238
tool_count: self.tool_names().len(), - 5239
structured_output: hint.structured_output, - 5240
reasoning_required: hint.reasoning_required, - 5241
evidence_required: hint.evidence_required, - 5242
}); - 5243
let objective = vak_llm::QualityObjective::resolve( - 5244
(route_cfg.objective != "auto").then_some(route_cfg.objective.as_str()), - 5245
demand.band, - 5246
); - 5247
- 5248
let belief_map = vak_llm::BeliefMap { - 5249
multipliers: self.inner.beliefs.snapshot().multipliers, - 5250
}; - 5251
let finops_cfg = self.inner.config.finops.clone(); - 5252
let home = self.sessions_home(); - 5253
let hints = route_cfg.quality_hints.clone(); - 5254
let ranked = vak_llm::order_ladder_v2( - 5255
candidates, - 5256
&routing::EvidenceLedger::new(&home).snapshot(), - 5257
&belief_map, - 5258
objective, - 5259
&hints, - 5260
move |m: &str| { - 5261
vak_config::finops::resolve_usd_per_mtok(m, &finops_cfg.price_overrides) - 5262
.map(|(_, out)| out) - 5263
}, - 5264
); - 5265
- 5266
let (ladder, annotations) = routing::assemble_ladder( - 5267
&primary, - 5268
ranked, - 5269
route_cfg.max_fallbacks, - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.