- 1
//! Routing evidence ledger + ladder admission (docs/design/15-reliability.md, - 2
//! Phase R upgrades ported from the vakrouter study). - 3
//! - 4
//! Evidence is append-only JSONL; reads apply a TTL and treat outcomes as - 5
//! success / failure / UNKNOWN. Billing without a verdict is UNKNOWN -- - 6
//! it shrinks confidence without punishing the model. - 7
//! - 8
//! Beliefs are the session-scoped complement to persisted evidence: - 9
//! domain-weighted doubt (a provider outage says more than one malformed - 10
//! stream) that demotes a leg below fully-trusted peers until ONE success - 11
//! on that leg clears it. Reality outranks priors. - 12
- 13
use std::collections::HashMap; - 14
use std::io::{BufRead, BufReader, Write}; - 15
use std::path::{Path, PathBuf}; - 16
use std::sync::Mutex; - 17
- 18
use serde::{Deserialize, Serialize}; - 19
use vak_llm::{EvidenceSnapshot, ModelEvidence, RouteLeg}; - 20
- 21
const EVIDENCE_TTL_DAYS: u64 = 30; - 22
- 23
#[derive(Debug, Clone, Serialize, Deserialize)] - 24
pub struct EvidenceRow { - 25
pub ts: chrono::DateTime<chrono::Utc>, - 26
pub provider: String, - 27
pub model: String, - 28
/// success | failure | unknown - 29
pub outcome: String, - 30
pub latency_ms: u64, - 31
} - 32
- 33
fn p50(samples: &mut [u64]) -> Option<u64> { - 34
if samples.is_empty() { - 35
return None; - 36
} - 37
samples.sort_unstable(); - 38
Some(samples[(samples.len() - 1) / 2]) - 39
} - 40
- 41
pub struct EvidenceLedger { - 42
path: PathBuf, - 43
} - 44
- 45
impl EvidenceLedger { - 46
pub fn new(sessions_home: &Path) -> Self { - 47
EvidenceLedger { - 48
path: sessions_home.join("routing-evidence.jsonl"), - 49
} - 50
} - 51
- 52
pub fn append(&self, row: &EvidenceRow) -> std::io::Result<()> { - 53
if let Some(parent) = self.path.parent() { - 54
std::fs::create_dir_all(parent)?; - 55
} - 56
let line = serde_json::to_string(row) - 57
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?; - 58
let mut f = std::fs::OpenOptions::new() - 59
.create(true) - 60
.append(true) - 61
.open(&self.path)?; - 62
writeln!(f, "{line}") - 63
} - 64
- 65
/// TTL-filtered snapshot for the ordering function. Corrupt lines are - 66
/// skipped rather than trusted -- never misranked. Latency is a true - 67
/// p50 over each leg's samples, not first-seen. - 68
pub fn snapshot(&self) -> EvidenceSnapshot { - 69
let cutoff = chrono::Utc::now() - chrono::Duration::days(EVIDENCE_TTL_DAYS as i64); - 70
let mut by_key: HashMap<(String, String), ModelEvidence> = HashMap::new(); - 71
let mut latencies: HashMap<(String, String), Vec<u64>> = HashMap::new(); - 72
let Ok(f) = std::fs::File::open(&self.path) else { - 73
return EvidenceSnapshot::default(); - 74
}; - 75
for line in BufReader::new(f).lines().map_while(Result::ok) { - 76
let Ok(row) = serde_json::from_str::<EvidenceRow>(&line) else { - 77
continue; - 78
}; - 79
if row.ts < cutoff { - 80
continue; - 81
} - 82
let key = (row.provider.clone(), row.model.clone()); - 83
let e = by_key.entry(key.clone()).or_default(); - 84
match row.outcome.as_str() { - 85
"success" => { - 86
e.success += 1; - 87
latencies.entry(key).or_default().push(row.latency_ms); - 88
} - 89
"failure" => e.failure += 1, - 90
_ => e.unknown += 1, - 91
} - 92
} - 93
for (key, e) in by_key.iter_mut() { - 94
e.p50_latency_ms = p50(latencies.entry(key.clone()).or_default()); - 95
} - 96
EvidenceSnapshot { by_key } - 97
} - 98
- 99
/// Fold a finished run's work receipts into evidence. Attribution is - 100
/// per attempt: fallback legs record against the leg that actually - 101
/// served (or failed) -- never against the receipt's final stamp. - 102
pub fn record_receipts(&self, receipts: &[vak_llm::WorkReceipt]) { - 103
let now = chrono::Utc::now(); - 104
for r in receipts { - 105
for a in &r.attempts { - 106
let outcome = match a.settlement { - 107
vak_llm::Settlement::Ok => Some("success"), - 108
vak_llm::Settlement::Failed => Some("failure"), - 109
vak_llm::Settlement::Unknown => Some("unknown"), - 110
vak_llm::Settlement::Cancelled => None, // contributes nothing - 111
}; - 112
if let Some(outcome) = outcome { - 113
let (provider, model) = r.attempt_leg(a); - 114
let _ = self.append(&EvidenceRow { - 115
ts: now, - 116
provider: provider.to_string(), - 117
model: model.to_string(), - 118
outcome: outcome.into(), - 119
latency_ms: a.latency_ms, - 120
}); - 121
} - 122
} - 123
} - 124
} - 125
} - 126
- 127
/// Doubt weight per failure domain. A dead provider (.6) is strong - 128
/// evidence against its legs; one malformed stream (.3) is weaker; an - 129
/// account-scoped rejection (.05) barely moves the needle. Governance and - 130
/// transport failures (request/network/deadline/unknown) say nothing about - 131
/// answer quality and are ignored. - 132
const DOMAIN_DOUBT: &[(vak_llm::FailureDomain, f64)] = &[ - 133
(vak_llm::FailureDomain::Provider, 0.6), - 134
(vak_llm::FailureDomain::Model, 0.3), - 135
(vak_llm::FailureDomain::Account, 0.05), - 136
]; - 137
- 138
/// Accumulated doubt never fully zeroes a candidate. - 139
pub const BELIEF_FLOOR: f64 = 0.1; - 140
- 141
/// Bounded doubt history per leg; oldest observations drop first. - 142
const MAX_BELIEF_OBSERVATIONS: usize = 128; - 143
- 144
#[derive(Debug, Clone)] - 145
struct Observation { - 146
weight: f64, - 147
} - 148
- 149
/// Session-scoped domain-weighted doubt per (provider, model) leg. - 150
#[derive(Default)] - 151
pub struct BeliefState { - 152
inner: Mutex<HashMap<(String, String), Vec<Observation>>>, - 153
} - 154
- 155
impl BeliefState { - 156
pub fn new() -> Self { - 157
Self::default() - 158
} - 159
- 160
/// Record one dispatch outcome. Success clears ALL doubt for the leg: - 161
/// reality outranks priors. Non-evidence domains contribute nothing. - 162
pub fn record_outcome( - 163
&self, - 164
provider: &str, - 165
model: &str, - 166
domain: vak_llm::FailureDomain, - 167
succeeded: bool, - 168
) { - 169
let mut inner = self - 170
.inner - 171
.lock() - 172
.unwrap_or_else(std::sync::PoisonError::into_inner); - 173
let key = (provider.to_string(), model.to_string()); - 174
if succeeded { - 175
inner.remove(&key); - 176
return; - 177
} - 178
let Some(weight) = DOMAIN_DOUBT - 179
.iter() - 180
.find(|(d, _)| *d == domain) - 181
.map(|(_, w)| *w) - 182
else { - 183
return; - 184
}; - 185
let obs = inner.entry(key).or_default(); - 186
obs.push(Observation { weight }); - 187
if obs.len() > MAX_BELIEF_OBSERVATIONS { - 188
obs.drain(0..obs.len() - MAX_BELIEF_OBSERVATIONS); - 189
} - 190
} - 191
- 192
/// Multiplier in [BELIEF_FLOOR, 1] plus the reasons behind it. - 193
pub fn multiplier(&self, provider: &str, model: &str) -> (f64, Vec<String>) { - 194
let inner = self - 195
.inner - 196
.lock() - 197
.unwrap_or_else(std::sync::PoisonError::into_inner); - 198
let mut mult = 1.0_f64; - 199
let mut reasons = Vec::new(); - 200
if let Some(obs) = inner.get(&(provider.to_string(), model.to_string())) { - 201
for o in obs { - 202
mult *= 1.0 - 0.25 * o.weight; - 203
reasons.push(format!("doubt w={:.2}", o.weight)); - 204
} - 205
} - 206
if mult < BELIEF_FLOOR { - 207
mult = BELIEF_FLOOR; - 208
} - 209
(mult, reasons) - 210
} - 211
- 212
/// Deterministic snapshot for the pure ordering function. - 213
pub fn snapshot(&self) -> BeliefSnapshot { - 214
let inner = self - 215
.inner - 216
.lock() - 217
.unwrap_or_else(std::sync::PoisonError::into_inner); - 218
BeliefSnapshot { - 219
multipliers: inner - 220
.iter() - 221
.map(|(k, obs)| { - 222
let mut m = 1.0_f64; - 223
for o in obs { - 224
m *= 1.0 - 0.25 * o.weight; - 225
} - 226
(k.clone(), m.max(BELIEF_FLOOR)) - 227
}) - 228
.collect(), - 229
} - 230
} - 231
} - 232
- 233
/// Immutable view of belief multipliers keyed by (provider, model). - 234
#[derive(Debug, Clone, Default)] - 235
pub struct BeliefSnapshot { - 236
pub multipliers: HashMap<(String, String), f64>, - 237
} - 238
- 239
impl BeliefSnapshot { - 240
pub fn get(&self, provider: &str, model: &str) -> f64 { - 241
self.multipliers - 242
.get(&(provider.to_string(), model.to_string())) - 243
.copied() - 244
.unwrap_or(1.0) - 245
} - 246
} - 247
- 248
/// The admission-time routing decision frozen into the contract header. - 249
#[derive(Debug, Clone)] - 250
pub struct RoutePlan { - 251
pub ladder: Vec<RouteLeg>, - 252
/// "utility" | "balanced" | "quality-critical". - 253
pub objective: String, - 254
/// Freeze-time warnings (thin chain, dominant failure domain, - 255
/// unreachable cross-model fallbacks). - 256
pub annotations: Vec<String>, - 257
} - 258
- 259
/// Admission-time ordering wrapper over the versioned pure function. - 260
pub fn order( - 261
candidates: Vec<RouteLeg>, - 262
snap: &EvidenceSnapshot, - 263
cost_of: impl Fn(&str) -> Option<f64>, - 264
) -> Vec<RouteLeg> { - 265
vak_llm::route::order_ladder_v1(candidates, snap, false, cost_of) - 266
} - 267
- 268
/// Assemble the frozen ladder from ranked candidates (Phase R diversity - 269
/// constraints, ported from the vakrouter chain study). - 270
/// - 271
/// The operator-selected primary NEVER loses its head position — v2 - 272
/// ordering decides the FALLBACK order, never whether the user's explicit - 273
/// choice serves first. Remaining seats cap per provider at - 274
/// ceil(max_total/3) so one failure domain cannot own every slot. - 275
/// Annotations are freeze-time warnings for traces/TUI, never model input. - 276
pub fn assemble_ladder( - 277
primary: &RouteLeg, - 278
ranked: Vec<RouteLeg>, - 279
max_total: usize, - 280
cross_model_requested: bool, - 281
) -> (Vec<RouteLeg>, Vec<String>) { - 282
let max_total = max_total.max(1); - 283
let seats_per_provider = max_total.div_ceil(3).max(1); - 284
let mut legs = vec![primary.clone()]; - 285
let mut counts: HashMap<String, usize> = HashMap::new(); - 286
counts.insert(primary.provider.clone(), 1); - 287
- 288
for leg in ranked { - 289
if legs.len() >= max_total { - 290
break; - 291
} - 292
if leg == *primary { - 293
continue; - 294
} - 295
let seats = counts.entry(leg.provider.clone()).or_insert(0); - 296
let independent_credential = leg.provider == primary.provider - 297
&& leg.credential_id.is_some() - 298
&& leg.credential_id != primary.credential_id; - 299
if *seats >= seats_per_provider && !independent_credential { - 300
continue; - 301
} - 302
*seats += 1; - 303
legs.push(leg); - 304
} - 305
- 306
let mut annotations = Vec::new(); - 307
if cross_model_requested && legs.iter().all(|l| l.model == primary.model) { - 308
annotations.push( - 309
"cross-model fallback configured but none of the allowed models are \ - 310
reachable by warm discovery" - 311
.to_string(), - 312
); - 313
} - 314
if legs.len() <= 1 { - 315
annotations.push("single point of failure: no executable fallback candidates".to_string()); - 316
} else { - 317
let len = legs.len(); - 318
if let Some((dominant, max_seats)) = counts - 319
.iter() - 320
.max_by_key(|(_, v)| **v) - 321
.filter(|(_, v)| **v > len.div_ceil(2) && len > 2) - 322
{ - 323
annotations.push(format!( - 324
"{dominant} holds {max_seats} of {len} ladder seats — one dominant failure domain" - 325
)); - 326
} - 327
} - 328
(legs, annotations) - 329
} - 330
- 331
#[cfg(test)] - 332
#[allow(clippy::unwrap_used, clippy::expect_used)] - 333
mod tests { - 334
use super::*; - 335
use tempfile::tempdir; - 336
- 337
#[test] - 338
fn snapshot_respects_ttl_and_skips_corrupt_lines() { - 339
let dir = tempdir().unwrap(); - 340
let ledger = EvidenceLedger::new(dir.path()); - 341
let old = chrono::Utc::now() - chrono::Duration::days(40); - 342
ledger - 343
.append(&EvidenceRow { - 344
ts: old, - 345
provider: "p".into(), - 346
model: "m".into(), - 347
outcome: "failure".into(), - 348
latency_ms: 5, - 349
}) - 350
.unwrap(); - 351
ledger - 352
.append(&EvidenceRow { - 353
ts: chrono::Utc::now(), - 354
provider: "p".into(), - 355
model: "m".into(), - 356
outcome: "success".into(), - 357
latency_ms: 900, - 358
}) - 359
.unwrap(); - 360
ledger - 361
.append(&EvidenceRow { - 362
ts: chrono::Utc::now(), - 363
provider: "p".into(), - 364
model: "m".into(), - 365
outcome: "success".into(), - 366
latency_ms: 42, - 367
}) - 368
.unwrap(); - 369
let mut file_path = dir.path().join("routing-evidence.jsonl"); - 370
let mut existing = std::fs::read_to_string(&file_path).unwrap(); - 371
existing.push_str("not json\n"); - 372
std::fs::write(&mut file_path, existing).unwrap(); - 373
- 374
let snap = ledger.snapshot(); - 375
let e = snap.get("p", "m"); - 376
assert_eq!(e.success, 2); - 377
assert_eq!(e.failure, 0, "40-day-old rows must be TTL-dropped"); - 378
assert_eq!( - 379
e.p50_latency_ms, - 380
Some(42), - 381
"p50 must be the median of samples, not first-seen" - 382
); - 383
} - 384
- 385
#[test] - 386
fn record_receipts_attributes_per_attempt_and_maps_settlements() { - 387
let dir = tempdir().unwrap(); - 388
let ledger = EvidenceLedger::new(dir.path()); - 389
let mut r = vak_llm::WorkReceipt::new(vak_llm::WorkPurpose::Execute, "openai", "gpt-x"); - 390
r.record( - 391
vak_llm::AttemptReason::Initial, - 392
vak_llm::FailureDomain::Network, - 393
vak_llm::Settlement::Unknown, - 394
10, - 395
None, - 396
None, - 397
); - 398
// Fall over to another leg; the winning attempt must attribute to - 399
// THAT leg, not the receipt's final stamp. - 400
r.stamp_leg("anthropic", "claude-x"); - 401
r.record( - 402
vak_llm::AttemptReason::RouteFallback, - 403
vak_llm::FailureDomain::Unknown, - 404
vak_llm::Settlement::Ok, - 405
20, - 406
None, - 407
None, - 408
); - 409
ledger.record_receipts(std::slice::from_ref(&r)); - 410
let snap = ledger.snapshot(); - 411
assert_eq!( - 412
snap.get("openai", "gpt-x").unknown, - 413
1, - 414
"failed leg keeps its own attribution" - 415
); - 416
assert_eq!( - 417
snap.get("anthropic", "claude-x").success, - 418
1, - 419
"winning leg attributes to itself" - 420
); - 421
assert_eq!(snap.get("", "gpt-x").unknown, 0, "no empty-provider rows"); - 422
} - 423
- 424
#[test] - 425
fn beliefs_demote_on_doubt_and_clear_on_success() { - 426
let b = BeliefState::new(); - 427
b.record_outcome("p", "flaky", vak_llm::FailureDomain::Provider, false); - 428
let (mult, _reasons) = b.multiplier("p", "flaky"); - 429
assert!(mult < 1.0, "provider failure must create doubt"); - 430
b.record_outcome("p", "flaky", vak_llm::FailureDomain::Network, false); - 431
let (unchanged, _) = b.multiplier("p", "flaky"); - 432
assert_eq!( - 433
unchanged, mult, - 434
"governance-domain failures must NOT add doubt" - 435
); - 436
let (clean, _) = b.multiplier("q", "other"); - 437
assert_eq!(clean, 1.0); - 438
b.record_outcome("p", "flaky", vak_llm::FailureDomain::Unknown, true); - 439
let (cleared, _) = b.multiplier("p", "flaky"); - 440
assert_eq!(cleared, 1.0, "one success clears all doubt"); - 441
} - 442
- 443
#[test] - 444
fn belief_multiplier_floors_but_never_zeroes() { - 445
let b = BeliefState::new(); - 446
for _ in 0..40 { - 447
b.record_outcome("p", "dead", vak_llm::FailureDomain::Provider, false); - 448
} - 449
let (mult, _) = b.multiplier("p", "dead"); - 450
assert_eq!(mult, BELIEF_FLOOR); - 451
let snap = b.snapshot(); - 452
assert_eq!(snap.get("p", "dead"), BELIEF_FLOOR); - 453
assert_eq!(snap.get("p", "absent"), 1.0); - 454
} - 455
- 456
#[test] - 457
fn assemble_pins_primary_head_and_caps_provider_seats() { - 458
let primary = RouteLeg { - 459
provider: "anthropic".into(), - 460
model: "claude-x".into(), - 461
dialect: vak_llm::EndpointDialect::AnthropicMessages, - 462
credential_id: None, - 463
}; - 464
let mk = |m: &str| RouteLeg { - 465
provider: "openai".into(), - 466
model: m.into(), - 467
dialect: vak_llm::EndpointDialect::Responses, - 468
credential_id: None, - 469
}; - 470
let ranked = vec![ - 471
mk("gpt-a"), - 472
mk("gpt-b"), - 473
RouteLeg { - 474
provider: "google".into(), - 475
model: "gemini-a".into(), - 476
dialect: vak_llm::EndpointDialect::GoogleGenerateContent, - 477
credential_id: None, - 478
}, - 479
]; - 480
// max_total 4 → seats/provider = 2; openai can hold at most two - 481
// fallback seats alongside the primary. - 482
let (legs, annotations) = assemble_ladder(&primary, ranked, 4, true); - 483
assert_eq!(legs[0], primary, "primary must stay at the head"); - 484
assert_eq!(legs.len(), 4); - 485
assert!( - 486
!annotations - 487
.iter() - 488
.any(|a| a.contains("dominant failure domain")), - 489
"two of four seats is exactly half — not dominant" - 490
); - 491
- 492
// A wide ladder lets one provider hold three fallback seats; that - 493
// IS a dominant failure domain worth surfacing. - 494
let (legs, annotations) = assemble_ladder( - 495
&primary, - 496
vec![mk("gpt-a"), mk("gpt-b"), mk("gpt-c")], - 497
7, - 498
false, - 499
); - 500
assert_eq!(legs.len(), 4); - 501
assert!( - 502
annotations - 503
.iter() - 504
.any(|a| a.contains("openai holds 3 of 4") && a.contains("dominant")), - 505
"three of four seats on one provider must warn, got {annotations:?}" - 506
); - 507
} - 508
- 509
#[test] - 510
fn assemble_annotates_thin_chain_and_unreachable_cross_model() { - 511
let primary = RouteLeg { - 512
provider: "anthropic".into(), - 513
model: "claude-x".into(), - 514
dialect: vak_llm::EndpointDialect::AnthropicMessages, - 515
credential_id: None, - 516
}; - 517
let (legs, annotations) = assemble_ladder(&primary, vec![], 4, false); - 518
assert_eq!(legs.len(), 1); - 519
assert!(annotations.iter().any(|a| a.contains("single point"))); - 520
- 521
let (_, annotations) = assemble_ladder( - 522
&primary, - 523
vec![RouteLeg { - 524
provider: "anthropic".into(), - 525
model: "claude-y".into(), - 526
dialect: vak_llm::EndpointDialect::AnthropicMessages, - 527
credential_id: None, - 528
}], - 529
4, - 530
true, - 531
); - 532
assert!( - 533
!annotations.iter().any(|a| a.contains("cross-model")), - 534
"a reachable alternate model means cross-model worked" - 535
); - 536
let (_, annotations) = assemble_ladder(&primary, vec![], 4, true); - 537
assert!( - 538
annotations.iter().any(|a| a.contains("cross-model")), - 539
"configured-but-unreachable fallbacks must be surfaced" - 540
); - 541
} - 542
- 543
#[test] - 544
fn assemble_drops_duplicate_primary() { - 545
let primary = RouteLeg { - 546
provider: "anthropic".into(), - 547
model: "claude-x".into(), - 548
dialect: vak_llm::EndpointDialect::AnthropicMessages, - 549
credential_id: None, - 550
}; - 551
let ranked = vec![ - 552
primary.clone(), - 553
RouteLeg { - 554
provider: "openai".into(), - 555
model: "gpt".into(), - 556
dialect: vak_llm::EndpointDialect::Responses, - 557
credential_id: None, - 558
}, - 559
]; - 560
let (legs, _) = assemble_ladder(&primary, ranked, 4, false); - 561
assert_eq!(legs.len(), 2); - 562
} - 563
- 564
#[test] - 565
fn assemble_keeps_same_model_on_a_distinct_credential() { - 566
let primary = RouteLeg { - 567
provider: "openrouter".into(), - 568
model: "shared-model".into(), - 569
dialect: vak_llm::EndpointDialect::Responses, - 570
credential_id: Some("key-a".into()), - 571
}; - 572
let alternate = RouteLeg { - 573
credential_id: Some("key-b".into()), - 574
..primary.clone() - 575
}; - 576
let (legs, _) = assemble_ladder(&primary, vec![alternate], 3, false); - 577
assert_eq!(legs.len(), 2); - 578
assert_eq!(legs[1].credential_id.as_deref(), Some("key-b")); - 579
} - 580
} - 581
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.