- 4345
/// A distributed-bus secret (`vak_config::BUS_NATS_JWT_VAR`, - 4346
/// `BUS_NATS_NKEY_SEED_VAR`) through the secrets chain. - 4347
pub fn bus_secret(&self, env_var: &str) -> Option<String> { - 4348
self.scoped_secret(env_var) - 4349
} - 4350
- 4351
fn provider_secret(&self, env_var: &str) -> Option<String> { - 4352
self.scoped_secret(env_var) - 4353
.or_else(|| vak_config::get_var(env_var)) - 4354
} - 4355
- 4356
fn scoped_secret(&self, env_var: &str) -> Option<String> { - 4357
if self.agent_identity.is_some() - 4358
&& let Some(val) = - 4359
vak_config::read_env_file_var(&self.sessions_home().join(".env"), env_var) - 4360
{ - 4361
return Some(val); - 4362
} - 4363
vak_config::read_env_file_var(&self.inner.cwd.join(".env"), env_var) - 4364
.or_else(|| vak_config::read_env_file_var(&self.user_env_file(), env_var)) - 4365
.or_else(|| std::env::var(env_var).ok()) - 4366
} - 4367
- 4368
/// The chat transports a bot can be created on. - 4369
/// - 4370
/// A surface is a **transport, not a credential slot** (AGENTS.md - 4371
/// invariant 23): several bots can share one, each with its own token, - 4372
/// policy, permission mode, and route. The fixed per-surface env vars - 4373
/// this replaces (`TELEGRAM_BOT_TOKEN` and friends) could only ever - 4374
/// describe one bot per transport, which is why they are gone. - 4375
/// - 4376
/// **Alphabetical, and this is the only list.** Every surface — API, - 4377
/// admin console, desktop — renders from here rather than carrying its - 4378
/// own copy, so adding a transport is one edit and no channel can - 4379
/// quietly become the default by being first or by being the one a UI - 4380
/// happens to hardcode. - 4381
pub const SURFACES: &'static [ChatSurface] = &[ - 4382
ChatSurface { - 4383
id: "discord", - 4384
label: "Discord", - 4385
}, - 4386
ChatSurface { - 4387
id: "slack", - 4388
label: "Slack", - 4389
}, - 4390
ChatSurface { - 4391
id: "telegram", - 4392
label: "Telegram", - 4393
}, - 4394
]; - 4395
- 4396
/// True when `surface` names a transport vak can bridge. - 4397
pub fn is_surface(surface: &str) -> bool { - 4398
Self::SURFACES.iter().any(|s| s.id == surface) - 4399
} - 4400
- 4401
/// The human label for a surface id, falling back to the id itself so - 4402
/// an unknown value renders as data rather than as an empty cell. - 4403
pub fn surface_label(id: &str) -> &str { - 4404
Self::SURFACES - 4405
.iter() - 4406
.find(|s| s.id == id) - 4407
.map(|s| s.label) - 4408
.unwrap_or(id) - 4409
} - 4410
- 4411
/// Persist a bot's token into the Shared secret scope and - 4412
/// register a runtime override so an in-process check is correct right - 4413
/// away. `env` is the bot's own `token_env` (`BOT_TOKEN__<ID>`); the - 4414
/// bridge unit re-reads the credential store itself on restart. The - 4415
/// token never re-enters any response. - 4416
pub fn set_bot_token(&self, env: &str, token: &str) -> Result<String, CoreError> { - 4417
let token = token.trim(); - 4418
if token.is_empty() { - 4419
return Err(CoreError::InvalidConfig(format!("empty token for {env}"))); - 4420
} - 4421
let path = self.user_env_file(); - 4422
vak_config::upsert_env_file(&path, env, token) - 4423
.map_err(|e| CoreError::InvalidConfig(format!("writing {path:?}: {e}")))?; - 4424
vak_config::set_override(env, token); - 4425
Ok(env.to_string()) - 4426
} - 4427
- 4428
/// Revoke a bot's stored token: strip it from the user secret scope - 4429
/// and drop the runtime override. A token exported in the real environment - 4430
/// cannot be unset from here — the caller is told so it can say as much. - 4431
pub fn remove_bot_token(&self, env: &str) -> Result<RemovedKey, CoreError> { - 4432
let path = self.user_env_file(); - 4433
vak_config::remove_env_file_key(&path, env) - 4434
.map_err(|e| CoreError::InvalidConfig(format!("writing {path:?}: {e}")))?; - 4435
vak_config::clear_override(env); - 4436
vak_config::forget_dotenv_var(env); - 4437
Ok(RemovedKey { - 4438
env_var: env.to_string(), - 4439
shadowed_by_env: vak_config::get_var(env).is_some(), - 4440
}) - 4441
} - 4442
- 4443
/// Ask `provider` which models its configured key can actually reach. - 4444
/// - 4445
/// There is no baked-in catalogue: an out-of-date table silently hides - 4446
/// models a provider shipped yesterday and offers ones the key cannot - 4447
/// use. Results are cached briefly because pickers poll this. - 4448
pub async fn discover_models(&self, provider: &str) -> Result<Vec<String>, CoreError> { - 4449
const TTL: std::time::Duration = std::time::Duration::from_secs(300); - 4450
let pool = self.provider_auth_pool_for(provider)?; - 4451
let mut all = Vec::new(); - 4452
let mut last_error = None; - 4453
for auth in pool { - 4454
let key = ( - 4455
provider.to_string(), - 4456
auth.credential_id.clone().unwrap_or_default(), - 4457
); - 4458
if let Ok(cache) = self.inner.models_cache.lock() - 4459
&& let Some((at, models)) = cache.get(&key) - 4460
&& at.elapsed() < TTL - 4461
{ - 4462
all.extend(models.iter().cloned()); - 4463
continue; - 4464
} - 4465
match vak_llm::models::list_models(provider, &auth).await { - 4466
Ok(models) => { - 4467
all.extend(models.iter().cloned()); - 4468
if let Ok(mut cache) = self.inner.models_cache.lock() { - 4469
cache.insert(key, (std::time::Instant::now(), models)); - 4470
} - 4471
} - 4472
Err(error) => last_error = Some(error), - 4473
} - 4474
} - 4475
all.sort(); - 4476
all.dedup(); - 4477
if all.is_empty() - 4478
&& let Some(error) = last_error - 4479
{ - 4480
return Err(error.into()); - 4481
} - 4482
Ok(all) - 4483
} - 4484
- 4485
pub async fn bedrock_model_availability( - 4486
&self, - 4487
model_ids: &[String], - 4488
) -> Result<Vec<vak_llm::models::BedrockModelAvailability>, CoreError> { - 4489
let auth = self.provider_auth_for("bedrock")?; - 4490
Ok(vak_llm::models::bedrock_model_availability(&auth, model_ids).await?) - 4491
} - 4492
- 4493
/// Read provider-published account metadata without exposing credentials. - 4494
pub async fn provider_status( - 4495
&self, - 4496
provider: &str, - 4497
) -> Result<vak_llm::provider_status::ProviderStatus, CoreError> { - 4498
let auth = self.provider_auth_for(provider)?; - 4499
Ok(vak_llm::provider_status::inspect(provider, &auth).await?) - 4500
} - 4501
- 4502
/// Long TTL for provider-reported model metadata (context window, - 4503
/// output max): it changes only when the model itself does, so once a - 4504
/// value is cached a turn should virtually never wait on it again. - 4505
/// `route_context_limits` and `capacity_profile_for` are the only two - 4506
/// callers and now share this one cache and this one TTL. - 4507
const MODEL_METADATA_TTL: std::time::Duration = std::time::Duration::from_secs(3600); - 4508
- 4509
/// Provider metadata for `(provider, model, credential_id)`, served - 4510
/// from `self.inner.model_context_cache`. A fresh hit returns with no - 4511
/// I/O. A stale hit still returns immediately — the stale value — and - 4512
/// kicks off a single-flighted background refresh so the *next* call - 4513
/// sees a fresh one; the caller that found it stale never waits on the - 4514
/// network. Only a cold key (nothing cached yet) is fetched inline, - 4515
/// since there is nothing else to serve. - 4516
async fn model_context_cached( - 4517
&self, - 4518
provider: &str, - 4519
model: &str, - 4520
credential_id: Option<&str>, - 4521
) -> Option<vak_llm::models::ModelContext> { - 4522
let cache_key = ( - 4523
provider.to_string(), - 4524
model.to_string(), - 4525
credential_id.unwrap_or_default().to_string(), - 4526
); - 4527
let cached = self - 4528
.inner - 4529
.model_context_cache - 4530
.lock() - 4531
.ok() - 4532
.and_then(|cache| cache.get(&cache_key).cloned()); - 4533
if let Some((fetched_at, value)) = cached { - 4534
if fetched_at.elapsed() < Self::MODEL_METADATA_TTL { - 4535
return value; - 4536
} - 4537
let should_spawn = self - 4538
.inner - 4539
.model_context_refreshing - 4540
.lock() - 4541
.is_ok_and(|mut refreshing| refreshing.insert(cache_key.clone())); - 4542
if should_spawn { - 4543
let core = self.clone(); - 4544
let refresh_key = cache_key.clone(); - 4545
tokio::spawn(async move { - 4546
let (provider, model, credential_id) = refresh_key.clone(); - 4547
let fresh = match core.provider_auth_for_leg( - 4548
&provider, - 4549
(!credential_id.is_empty()).then_some(credential_id.as_str()), - 4550
) { - 4551
Ok(auth) => vak_llm::models::model_context(&provider, &auth, &model) - 4552
.await - 4553
.ok() - 4554
.flatten(), - 4555
Err(_) => None, - 4556
}; - 4557
if let Ok(mut cache) = core.inner.model_context_cache.lock() { - 4558
cache.insert(refresh_key.clone(), (std::time::Instant::now(), fresh)); - 4559
} - 4560
if let Ok(mut refreshing) = core.inner.model_context_refreshing.lock() { - 4561
refreshing.remove(&refresh_key); - 4562
} - 4563
}); - 4564
} - 4565
return value; - 4566
} - 4567
- 4568
// Cold: nothing cached yet, so this call is the one that populates it. - 4569
let value = match self.provider_auth_for_leg(provider, credential_id) { - 4570
Ok(auth) => vak_llm::models::model_context(provider, &auth, model) - 4571
.await - 4572
.ok() - 4573
.flatten(), - 4574
Err(_) => None, - 4575
}; - 4576
if let Ok(mut cache) = self.inner.model_context_cache.lock() { - 4577
cache.insert(cache_key, (std::time::Instant::now(), value.clone())); - 4578
} - 4579
value - 4580
} - 4581
- 4582
async fn route_context_limits( - 4583
&self, - 4584
primary: &vak_llm::RouteLeg, - 4585
fallback: &[vak_llm::RouteLeg], - 4586
) -> (u64, u64) { - 4587
let mut legs = Vec::with_capacity(1 + fallback.len()); - 4588
legs.push(primary.clone()); - 4589
legs.extend(fallback.iter().cloned()); - 4590
let mut context_window = self.inner.config.context_window; - 4591
let mut max_output = u64::from(self.inner.config.max_tokens); - 4592
for leg in legs { - 4593
let metadata = self - 4594
.model_context_cached(&leg.provider, &leg.model, leg.credential_id.as_deref()) - 4595
.await; - 4596
if let Some(metadata) = metadata { - 4597
context_window = context_window.min(metadata.input_tokens); - 4598
if let Some(output) = metadata.output_tokens { - 4599
max_output = max_output.min(output); - 4600
} - 4601
} - 4602
} - 4603
(context_window, max_output.min(context_window).max(1)) - 4604
} - 4605
- 4606
/// Whether `leg` reaches a runner on this machine: named `ollama`, or - 4607
/// its resolved credential's base URL host is loopback. Local models - 4608
/// are always probed in full (docs/design/68-context-engine.md §1: - 4609
/// "the probe is free apart from time, and time is exactly what it - 4610
/// saves"); a hosted model only gets the full horizon ladder when the - 4611
/// operator opts in via `[probe] hosted = "full"`. - 4612
fn is_local_provider(&self, leg: &vak_llm::RouteLeg) -> bool { - 4613
if leg.provider == "ollama" { - 4614
return true; - 4615
} - 4616
self.provider_auth_for_leg(&leg.provider, leg.credential_id.as_deref()) - 4617
.ok() - 4618
.and_then(|auth| auth.base_url) - 4619
.as_deref() - 4620
.and_then(url_host) - 4621
.is_some_and(is_loopback_host) - 4622
} - 4623
- 4624
/// A background probe on a key repeatedly interrupted by real turns is - 4625
/// spaced out rather than respawned on every single turn completion. - 4626
const CAPACITY_PROBE_MIN_INTERVAL: std::time::Duration = std::time::Duration::from_secs(60); - 4627
/// Bounds any single probe rung request: a hung provider must not stall - 4628
/// the background probe indefinitely. - 4629
const PROBE_REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); - 4630
/// A rung whose expected prefill time (its token target divided by the - 4631
/// measured prefill throughput so far) exceeds this is skipped rather - 4632
/// than sent — every larger rung would be even slower, so the ladder - 4633
/// stops there and reports what it already confirmed. - 4634
const PROBE_MAX_RUNG_SECS: f64 = 45.0; - 4635
- 4636
/// The measured capacity profile for `leg` (docs/design/68-context-engine.md - 4637
/// §1), bound immediately from what is already known — the ledger's or - 4638
/// process cache's last profile for this key (even if stale: stale - 4639
/// beats blocking), or a metadata-only profile on a cold key. Never - 4640
/// runs the horizon-ladder probe itself: that never belongs on a - 4641
/// turn's critical path. When the bound profile is missing, stale, or - 4642
/// `needs_reprobe`, this cancels any in-flight background probe for - 4643
/// the same key (a real turn now wants the model) and leaves it to - 4644
/// `maybe_start_capacity_probe`, called once this turn is done, to - 4645
/// (re)start the ladder in the background. A freshly bound profile not - 4646
/// already in this session's ledger is recorded as a `CapacityProbe` - 4647
/// activity — this is also how a session catches up on a profile a - 4648
/// background probe delivered since it last bound this key. - 4649
async fn capacity_profile_for( - 4650
&self, - 4651
leg: &vak_llm::RouteLeg, - 4652
session: &mut SessionLog, - 4653
) -> vak_context::capacity::CapacityProfile { - 4654
let now = std::time::SystemTime::now(); - 4655
- 4656
// Metadata must be fetched before the key is built: Ollama's - 4657
// `/api/show` reports the running quantisation, and a requantised - 4658
// model must key its own profile rather than inheriting a stale one - 4659
// (docs/design/68-context-engine.md §1). Served from the shared, - 4660
// stale-while-revalidate cache (`model_context_cached`) rather than - 4661
// a fresh network call every turn. - 4662
let metadata = self - 4663
.model_context_cached(&leg.provider, &leg.model, leg.credential_id.as_deref()) - 4664
.await; - 4665
let key = vak_context::capacity::ProfileKey { - 4666
provider: leg.provider.clone(), - 4667
model: leg.model.clone(), - 4668
quantisation: metadata.as_ref().and_then(|m| m.quantisation.clone()), - 4669
}; - 4670
- 4671
// A real turn is about to use this model: a background probe for - 4672
// the same key would only compete with it for the same compute. - 4673
// The probe restarts its ladder on its next background attempt - 4674
// rather than resuming mid-rung — simpler and race-free, and cheap - 4675
// to redo ("the probe is free apart from time", §1). - 4676
if let Ok(mut probes) = self.inner.capacity_probes.lock() - 4677
&& let Some(token) = probes.remove(&key) - 4678
{ - 4679
token.cancel(); - 4680
} - 4681
- 4682
// The more recently probed of the ledger's latest record and the - 4683
// in-process cache wins: the ledger carries every feedback update - 4684
// the turn loop wrote (tightened horizons, calibrated tokens/char) - 4685
// and survives a restart, but a background probe (§1: "only while - 4686
// that model is idle") updates only the process cache, so it can - 4687
// now be strictly newer than what this session has ever logged. - 4688
let from_session = - 4689
session.latest_capacity_profile::<_, vak_context::capacity::CapacityProfile>(&key); - 4690
let from_cache = self - 4691
.inner - 4692
.capacity_cache - 4693
.lock() - 4694
.ok() - 4695
.and_then(|cache| cache.get(&key).cloned()); - 4696
let cached = match (from_session, from_cache) { - 4697
(Some(a), Some(b)) => Some(if b.provenance.probed_at > a.provenance.probed_at { - 4698
b - 4699
} else { - 4700
a - 4701
}), - 4702
(Some(a), None) => Some(a), - 4703
(None, Some(b)) => Some(b), - 4704
(None, None) => None, - 4705
}; - 4706
let declared_window = metadata - 4707
.as_ref() - 4708
.map(|m| m.input_tokens) - 4709
.unwrap_or(self.inner.config.context_window); - 4710
let output_reserve = metadata - 4711
.as_ref() - 4712
.and_then(|m| m.output_tokens) - 4713
.unwrap_or(u64::from(self.inner.config.max_tokens)); - 4714
let metadata_digest = format!("{declared_window}:{output_reserve}"); - 4715
- 4716
let mut profile = match cached { - 4717
Some(profile) => profile, - 4718
None => vak_context::capacity::CapacityProfile::from_metadata_only( - 4719
declared_window, - 4720
output_reserve, - 4721
metadata_digest, - 4722
now, - 4723
), - 4724
}; - 4725
profile.provenance.quantisation = key.quantisation.clone(); - 4726
- 4727
if let Ok(mut cache) = self.inner.capacity_cache.lock() { - 4728
cache.insert(key.clone(), profile.clone()); - 4729
} - 4730
- 4731
let already_logged = session - 4732
.latest_capacity_profile::<_, vak_context::capacity::CapacityProfile>(&key) - 4733
.is_some_and(|logged| logged.provenance.probed_at == profile.provenance.probed_at); - 4734
if !already_logged { - 4735
let mut data = std::collections::BTreeMap::new(); - 4736
if let Ok(key_json) = serde_json::to_string(&key) { - 4737
data.insert("key".into(), key_json); - 4738
} - 4739
if let Ok(profile_json) = serde_json::to_string(&profile) { - 4740
data.insert("profile".into(), profile_json); - 4741
} - 4742
let activity_id = now - 4743
.duration_since(std::time::UNIX_EPOCH) - 4744
.map(|d| format!("activity-{}", d.as_nanos())) - 4745
.unwrap_or_else(|_| format!("activity-{}", uuid_like())); - 4746
let _ = session.append_activity(vak_session::ActivityRecord { - 4747
activity_id, - 4748
turn: None, - 4749
kind: vak_session::ActivityKind::CapacityProbe, - 4750
status: vak_session::ActivityStatus::Succeeded, - 4751
label: format!("Capacity profile bound for {}/{}", leg.provider, leg.model), - 4752
detail: None, - 4753
data, - 4754
}); - 4755
} - 4756
- 4757
profile - 4758
} - 4759
- 4760
/// Starts a background horizon-ladder probe for `leg`'s profile key, - 4761
/// but only when warranted: this leg is eligible for a full probe - 4762
/// (local, or hosted with `[probe] hosted = "full"`), the cached - 4763
/// profile is missing/stale/`needs_reprobe`, no probe for this key is - 4764
/// already running (single-flight), and the key was not attempted - 4765
/// within `CAPACITY_PROBE_MIN_INTERVAL`. Called only once a turn is - 4766
/// done (docs/design/68 §1: "only while that model is idle") — never - 4767
/// from a turn's own critical path. The task updates the Core cache on - 4768
/// completion; the session ledger catches up at the *next* bind - 4769
/// (`capacity_profile_for`), since the ledger belongs to the running - 4770
/// turn, not to this detached task. - 4771
async fn maybe_start_capacity_probe(&self, leg: &vak_llm::RouteLeg) { - 4772
let local = self.is_local_provider(leg); - 4773
if !(local || self.inner.config.probe.hosted == "full") { - 4774
return; - 4775
} - 4776
let metadata = self - 4777
.model_context_cached(&leg.provider, &leg.model, leg.credential_id.as_deref()) - 4778
.await; - 4779
let key = vak_context::capacity::ProfileKey { - 4780
provider: leg.provider.clone(), - 4781
model: leg.model.clone(), - 4782
quantisation: metadata.as_ref().and_then(|m| m.quantisation.clone()), - 4783
}; - 4784
let now = std::time::SystemTime::now(); - 4785
let declared_window = metadata - 4786
.as_ref() - 4787
.map(|m| m.input_tokens) - 4788
.unwrap_or(self.inner.config.context_window); - 4789
let output_reserve = metadata - 4790
.as_ref() - 4791
.and_then(|m| m.output_tokens) - 4792
.unwrap_or(u64::from(self.inner.config.max_tokens)); - 4793
let metadata_digest = format!("{declared_window}:{output_reserve}"); - 4794
- 4795
let fresh = self - 4796
.inner - 4797
.capacity_cache - 4798
.lock() - 4799
.ok() - 4800
.and_then(|cache| cache.get(&key).cloned()) - 4801
.is_some_and(|profile| { - 4802
// A metadata-only profile (no rungs) is a placeholder, not - 4803
// a satisfied probe: this leg is eligible for the full - 4804
// ladder, so only an actual probe result counts as fresh. - 4805
!profile.provenance.rungs.is_empty() - 4806
&& !profile.needs_reprobe - 4807
&& profile.provenance.metadata_digest == metadata_digest - 4808
&& !profile.is_stale(now, local) - 4809
}); - 4810
if fresh { - 4811
return; - 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(),
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.