- 10347
core = core.with_agent_identity(Some(agent)); - 10348
- 10349
assert_eq!( - 10350
core.provider_secret("ANTHROPIC_API_KEY"), - 10351
Some("shared-anthropic-key".into()) - 10352
); - 10353
- 10354
let agent_home = core.sessions_home(); - 10355
vak_config::upsert_env_file( - 10356
&agent_home.join(".env"), - 10357
"ANTHROPIC_API_KEY", - 10358
"agent-private-key", - 10359
) - 10360
.unwrap(); - 10361
- 10362
assert_eq!( - 10363
core.provider_secret("ANTHROPIC_API_KEY"), - 10364
Some("agent-private-key".into()) - 10365
); - 10366
} - 10367
} - 10368
- 10369
/// docs/design/68-context-engine.md §1: a loopback ("ollama") leg is always - 10370
/// probed in full, a hosted leg is not unless `[probe] hosted = "full"`. - 10371
/// Uses the "Scripted"/`Fn`-provider pattern already used throughout - 10372
/// `vak-agent`'s tests (e.g. `crates/vak-agent/tests/doom_loop.rs`): a fake - 10373
/// `Provider` registered directly into the registry, no real network. The - 10374
/// hosted comparison leg names a provider `provider_auth_for` does not - 10375
/// recognize, so its auth resolution fails deterministically regardless of - 10376
/// what real credentials happen to be set in the developer's environment — - 10377
/// the point being tested is "no probe without opt-in", not credential - 10378
/// plumbing. - 10379
#[cfg(test)] - 10380
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 10381
mod capacity_probe_tests { - 10382
use super::Core; - 10383
use std::sync::Arc; - 10384
use tokio_util::sync::CancellationToken; - 10385
use vak_llm::{ - 10386
AssistantMessage, ChatRequest, ContentBlock, LlmError, Provider, StopReason, Usage, - 10387
}; - 10388
use vak_session::SessionLog; - 10389
use vak_session::types::{FrozenContract, SessionHeader}; - 10390
- 10391
/// Always answers a probe rung by calling `probe_ack` — enough to drive - 10392
/// the horizon ladder to convergence without any network I/O. - 10393
struct AlwaysFollowsProbe; - 10394
- 10395
#[async_trait::async_trait] - 10396
impl Provider for AlwaysFollowsProbe { - 10397
fn name(&self) -> &str { - 10398
"ollama" - 10399
} - 10400
- 10401
async fn stream( - 10402
&self, - 10403
_request: ChatRequest, - 10404
_cancel: CancellationToken, - 10405
) -> Result<vak_llm::EventStream, LlmError> { - 10406
let (mut sink, rx) = vak_llm::stream::channel(4); - 10407
sink.close_message(AssistantMessage { - 10408
content: vec![ContentBlock::ToolUse { - 10409
id: "probe-1".into(), - 10410
name: "probe_ack".into(), - 10411
input: serde_json::json!({"ok": true}), - 10412
}], - 10413
stop_reason: StopReason::ToolUse, - 10414
usage: Usage { - 10415
input_tokens: 4_000, - 10416
output_tokens: 5, - 10417
prefill_ms: Some(50), - 10418
..Default::default() - 10419
}, - 10420
model: "fake-ollama-model".into(), - 10421
response_id: None, - 10422
}) - 10423
.await; - 10424
Ok(rx) - 10425
} - 10426
} - 10427
- 10428
/// Polls the process capacity cache until `maybe_start_capacity_probe`'s - 10429
/// detached task has recorded a converged profile (at least one rung), - 10430
/// or panics after a generous bound. Fake providers in this module - 10431
/// answer in-process with no real I/O, so convergence is normally a - 10432
/// handful of scheduler ticks. - 10433
async fn wait_for_capacity_probe( - 10434
core: &Core, - 10435
key: &vak_context::capacity::ProfileKey, - 10436
) -> vak_context::capacity::CapacityProfile { - 10437
for _ in 0..200 { - 10438
if let Some(profile) = core - 10439
.inner - 10440
.capacity_cache - 10441
.lock() - 10442
.ok() - 10443
.and_then(|cache| cache.get(key).cloned()) - 10444
&& !profile.provenance.rungs.is_empty() - 10445
{ - 10446
return profile; - 10447
} - 10448
tokio::time::sleep(std::time::Duration::from_millis(10)).await; - 10449
} - 10450
panic!("background capacity probe did not converge in time"); - 10451
} - 10452
- 10453
fn header() -> SessionHeader { - 10454
SessionHeader { - 10455
agent: None, - 10456
session_id: "s-capacity-probe".into(), - 10457
created_at: chrono::Utc::now(), - 10458
cwd: std::path::PathBuf::from("/tmp/proj"), - 10459
parent_session_id: None, - 10460
contract_id: None, - 10461
work_item_id: None, - 10462
conversation: None, - 10463
contract: FrozenContract { - 10464
app_version: "0.1.0".into(), - 10465
provider: "ollama".into(), - 10466
model: "fake-ollama-model".into(), - 10467
route_ladder: Vec::new(), - 10468
route_objective: String::new(), - 10469
route_annotations: Vec::new(), - 10470
system_prompt: String::new(), - 10471
permission_mode: "workspace-write".into(), - 10472
capabilities: Vec::new(), - 10473
prompt_layers: Vec::new(), - 10474
}, - 10475
} - 10476
} - 10477
- 10478
#[tokio::test] - 10479
async fn loopback_leg_is_probed_and_hosted_leg_is_not_by_default() { - 10480
super::isolate_global_config(); - 10481
let dir = tempfile::tempdir().unwrap(); - 10482
let core = Core::new_with_trust(dir.path().to_path_buf(), true).unwrap(); - 10483
// No real Ollama server needs to exist: metadata discovery - 10484
// (`model_context`) is a plain reqwest call this test does not - 10485
// control, so it is pointed at a bound-then-dropped loopback port — - 10486
// guaranteed connection-refused, which `model_context`'s own - 10487
// fallback turns into a fixed 8192/4096 declared window/output - 10488
// reserve, deterministically and without depending on what may or - 10489
// may not be listening on the default Ollama port. - 10490
let unused_port = { - 10491
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); - 10492
listener.local_addr().unwrap().port() - 10493
}; - 10494
vak_config::set_override( - 10495
"VAK_OLLAMA_BASE_URL", - 10496
format!("http://127.0.0.1:{unused_port}/v1"), - 10497
); - 10498
- 10499
// The probe ladder's provider client comes from the registry, not - 10500
// from `model_context`'s direct reqwest call — registering a fake - 10501
// factory here intercepts exactly that seam. - 10502
core.inner.registry.register("ollama", |_auth| { - 10503
Ok(Arc::new(AlwaysFollowsProbe) as Arc<dyn Provider>) - 10504
}); - 10505
- 10506
let loopback_leg = vak_llm::RouteLeg { - 10507
provider: "ollama".into(), - 10508
model: "fake-ollama-model".into(), - 10509
dialect: vak_llm::EndpointDialect::default(), - 10510
credential_id: None, - 10511
}; - 10512
assert!(core.is_local_provider(&loopback_leg)); - 10513
- 10514
let session_path = dir.path().join("probe-session.jsonl"); - 10515
let mut session = SessionLog::create(session_path, header()).unwrap(); - 10516
// `capacity_profile_for` never blocks on a probe: the very first - 10517
// bind returns a metadata-only profile immediately. - 10518
let bound = core.capacity_profile_for(&loopback_leg, &mut session).await; - 10519
assert_eq!(bound.declared_window, 8_192, "the ollama metadata fallback"); - 10520
assert_eq!(bound.output_reserve, 4_096); - 10521
assert!( - 10522
bound.provenance.rungs.is_empty(), - 10523
"no ladder rung may run on the turn's own critical path" - 10524
); - 10525
- 10526
// The ladder only runs once `maybe_start_capacity_probe` is called - 10527
// (docs/design/68 §1: "only while that model is idle: start after - 10528
// a turn completes") — never inside `capacity_profile_for` itself. - 10529
let key = vak_context::capacity::ProfileKey { - 10530
provider: "ollama".into(), - 10531
model: "fake-ollama-model".into(), - 10532
quantisation: None, - 10533
}; - 10534
core.maybe_start_capacity_probe(&loopback_leg).await; - 10535
let probed = wait_for_capacity_probe(&core, &key).await; - 10536
assert_eq!( - 10537
probed.instruction_horizon.tokens, 4_000, - 10538
"the single rung under declared_window * 0.9 that the fake provider followed" - 10539
); - 10540
assert_eq!(probed.instruction_horizon.confidence, 0.9); - 10541
assert_eq!(probed.provenance.rungs.len(), 1); - 10542
assert!(probed.provenance.rungs[0].accepted); - 10543
assert_eq!(probed.provenance.rungs[0].followed_instruction, Some(true)); - 10544
- 10545
// The next bind catches the session's ledger up on what the - 10546
// background probe delivered (the ledger belongs to the turn, the - 10547
// detached task has none). - 10548
let caught_up = core.capacity_profile_for(&loopback_leg, &mut session).await; - 10549
assert_eq!(caught_up.instruction_horizon.tokens, 4_000); - 10550
let recorded: vak_context::capacity::CapacityProfile = session - 10551
.latest_capacity_profile(&key) - 10552
.expect("the probe must be recorded as a ledger activity"); - 10553
assert_eq!(recorded.instruction_horizon.tokens, 4_000); - 10554
- 10555
// A hosted provider `provider_auth_for` has no wiring for at all - 10556
// fails auth resolution deterministically, regardless of any real - 10557
// credential the environment happens to carry for a KNOWN - 10558
// provider name — exactly what "not probed by default" needs to be - 10559
// hermetic. - 10560
let hosted_leg = vak_llm::RouteLeg { - 10561
provider: "hosted-test-provider-not-wired".into(), - 10562
model: "some-frontier-model".into(), - 10563
dialect: vak_llm::EndpointDialect::default(), - 10564
credential_id: None, - 10565
}; - 10566
assert!(!core.is_local_provider(&hosted_leg)); - 10567
assert_eq!(core.inner.config.probe.hosted, "none"); - 10568
- 10569
let hosted_session_path = dir.path().join("hosted-session.jsonl"); - 10570
let mut hosted_session = SessionLog::create(hosted_session_path, header()).unwrap(); - 10571
let hosted_profile = core - 10572
.capacity_profile_for(&hosted_leg, &mut hosted_session) - 10573
.await; - 10574
core.maybe_start_capacity_probe(&hosted_leg).await; - 10575
- 10576
assert!( - 10577
hosted_profile.provenance.rungs.is_empty(), - 10578
"no ladder rung should run for a hosted leg without opt-in" - 10579
); - 10580
assert_eq!(hosted_profile.instruction_horizon.confidence, 0.3); - 10581
assert_eq!( - 10582
hosted_profile.instruction_horizon.tokens, hosted_profile.declared_window, - 10583
"an unprobed hosted profile starts the horizon at the declared window" - 10584
); - 10585
- 10586
vak_config::clear_override("VAK_OLLAMA_BASE_URL"); - 10587
} - 10588
- 10589
/// Always follows the probe instruction; a request identical to the - 10590
/// previous one (a repeated sample inside a rung, or the cache rung's - 10591
/// second request) reports a cache hit, so the test can assert - 10592
/// `CacheBehaviour::ProviderReported` end to end through - 10593
/// `capacity_profile_for` without any network I/O. - 10594
struct FollowsProbeAndReportsCacheOnThirdCall { - 10595
calls: std::sync::atomic::AtomicU32, - 10596
last_fingerprint: std::sync::Mutex<Option<String>>, - 10597
} - 10598
- 10599
#[async_trait::async_trait] - 10600
impl Provider for FollowsProbeAndReportsCacheOnThirdCall { - 10601
fn name(&self) -> &str { - 10602
"ollama" - 10603
} - 10604
- 10605
async fn stream( - 10606
&self, - 10607
_request: ChatRequest, - 10608
_cancel: CancellationToken, - 10609
) -> Result<vak_llm::EventStream, LlmError> { - 10610
let call = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - 10611
let fingerprint = serde_json::to_string(&_request.messages).unwrap_or_default(); - 10612
let repeated = self - 10613
.last_fingerprint - 10614
.lock() - 10615
.map(|mut last| { - 10616
let same = last.as_deref() == Some(fingerprint.as_str()); - 10617
*last = Some(fingerprint); - 10618
same - 10619
}) - 10620
.unwrap_or(false); - 10621
let (mut sink, rx) = vak_llm::stream::channel(4); - 10622
sink.close_message(AssistantMessage { - 10623
content: vec![ContentBlock::ToolUse { - 10624
id: format!("probe-{call}"), - 10625
name: "probe_ack".into(), - 10626
input: serde_json::json!({"ok": true}), - 10627
}], - 10628
stop_reason: StopReason::ToolUse, - 10629
usage: Usage { - 10630
input_tokens: 4_000, - 10631
output_tokens: 5, - 10632
cache_read_input_tokens: repeated.then_some(3_500), - 10633
..Default::default() - 10634
}, - 10635
model: "fake-ollama-model".into(), - 10636
response_id: None, - 10637
}) - 10638
.await; - 10639
Ok(rx) - 10640
} - 10641
} - 10642
- 10643
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 10644
async fn cache_rung_detects_a_provider_reported_hit_and_records_both_rungs_in_signals() { - 10645
super::isolate_global_config(); - 10646
let dir = tempfile::tempdir().unwrap(); - 10647
let core = Core::new_with_trust(dir.path().to_path_buf(), true).unwrap(); - 10648
let unused_port = { - 10649
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); - 10650
listener.local_addr().unwrap().port() - 10651
}; - 10652
vak_config::set_override( - 10653
"VAK_OLLAMA_BASE_URL", - 10654
format!("http://127.0.0.1:{unused_port}/v1"), - 10655
); - 10656
core.inner.registry.register("ollama", |_auth| { - 10657
Ok(Arc::new(FollowsProbeAndReportsCacheOnThirdCall { - 10658
calls: std::sync::atomic::AtomicU32::new(0), - 10659
last_fingerprint: std::sync::Mutex::new(None), - 10660
}) as Arc<dyn Provider>) - 10661
}); - 10662
- 10663
let loopback_leg = vak_llm::RouteLeg { - 10664
provider: "ollama".into(), - 10665
model: "fake-ollama-model".into(), - 10666
dialect: vak_llm::EndpointDialect::default(), - 10667
credential_id: None, - 10668
}; - 10669
let session_path = dir.path().join("cache-rung-session.jsonl"); - 10670
let mut session = SessionLog::create(session_path, header()).unwrap(); - 10671
let _ = core.capacity_profile_for(&loopback_leg, &mut session).await; - 10672
core.maybe_start_capacity_probe(&loopback_leg).await; - 10673
let key = vak_context::capacity::ProfileKey { - 10674
provider: "ollama".into(), - 10675
model: "fake-ollama-model".into(), - 10676
quantisation: None, - 10677
}; - 10678
let probed = wait_for_capacity_probe(&core, &key).await; - 10679
- 10680
assert_eq!( - 10681
probed.cache, - 10682
vak_context::capacity::CacheBehaviour::ProviderReported - 10683
); - 10684
assert!( - 10685
probed - 10686
.provenance - 10687
.signals - 10688
.iter() - 10689
.any(|s| s.contains("cache rung") && s.contains("ProviderReported")), - 10690
"both cache-rung latencies and the outcome must be recorded as \ - 10691
provenance signals: {:?}", - 10692
probed.provenance.signals - 10693
); - 10694
- 10695
vak_config::clear_override("VAK_OLLAMA_BASE_URL"); - 10696
} - 10697
} - 10698
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.