- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
//! Frozen-ladder routing (docs/design/15-reliability.md): dispatch walks the - 4
//! frozen candidate legs on typed failures; ceiling/receipts/endurance - 5
//! are shared across legs; walking the ladder is contract execution. - 6
- 7
use std::collections::VecDeque; - 8
use std::sync::Arc; - 9
use std::sync::atomic::{AtomicU32, Ordering}; - 10
- 11
use tempfile::tempdir; - 12
use tokio::sync::mpsc; - 13
use tokio_util::sync::CancellationToken; - 14
- 15
use vak_agent::{Agent, AgentConfig, AutoApprove, SteeringQueues, TurnOutcome}; - 16
use vak_llm::stream; - 17
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; - 18
use vak_llm::{AttemptReason, EventStream, LlmError, Provider}; - 19
use vak_permission::{Mode, PermissionEngine}; - 20
use vak_session::types::{FrozenContract, SessionHeader}; - 21
use vak_session::{SessionLog, SessionPath}; - 22
- 23
fn text_msg(t: &str) -> AssistantMessage { - 24
AssistantMessage { - 25
content: vec![ContentBlock::text(t)], - 26
stop_reason: StopReason::EndTurn, - 27
usage: Usage { - 28
input_tokens: 5, - 29
output_tokens: 3, - 30
..Default::default() - 31
}, - 32
model: "fallback-model".into(), - 33
response_id: None, - 34
} - 35
} - 36
- 37
/// Always fails with a network error (blind failure domain). - 38
struct AlwaysNetwork { - 39
calls: AtomicU32, - 40
} - 41
- 42
struct PanickingProvider; - 43
- 44
struct TerminalQuota; - 45
- 46
#[async_trait::async_trait] - 47
impl Provider for TerminalQuota { - 48
fn name(&self) -> &str { - 49
"quota-limited" - 50
} - 51
- 52
async fn stream( - 53
&self, - 54
_request: ChatRequest, - 55
_cancel: CancellationToken, - 56
) -> Result<EventStream, LlmError> { - 57
let (mut sink, rx) = stream::channel(64); - 58
sink.close_error(LlmError::RateLimit { - 59
message: "free-models-per-day exhausted".into(), - 60
retry_after_secs: None, - 61
}) - 62
.await; - 63
Ok(rx) - 64
} - 65
} - 66
- 67
#[async_trait::async_trait] - 68
impl Provider for PanickingProvider { - 69
fn name(&self) -> &str { - 70
"primary-panics" - 71
} - 72
- 73
async fn stream( - 74
&self, - 75
_request: ChatRequest, - 76
_cancel: CancellationToken, - 77
) -> Result<EventStream, LlmError> { - 78
panic!("provider adapter panic"); - 79
} - 80
} - 81
- 82
#[async_trait::async_trait] - 83
impl Provider for AlwaysNetwork { - 84
fn name(&self) -> &str { - 85
"primary-net-dead" - 86
} - 87
- 88
async fn stream( - 89
&self, - 90
_request: ChatRequest, - 91
_cancel: CancellationToken, - 92
) -> Result<EventStream, LlmError> { - 93
self.calls.fetch_add(1, Ordering::SeqCst); - 94
let (mut sink, rx) = stream::channel(64); - 95
sink.close_error(LlmError::Network("conn reset".into())) - 96
.await; - 97
Ok(rx) - 98
} - 99
} - 100
- 101
struct Scripted { - 102
responses: std::sync::Mutex<VecDeque<AssistantMessage>>, - 103
} - 104
- 105
#[async_trait::async_trait] - 106
impl Provider for Scripted { - 107
fn name(&self) -> &str { - 108
"fallback-ok" - 109
} - 110
- 111
async fn stream( - 112
&self, - 113
_request: ChatRequest, - 114
_cancel: CancellationToken, - 115
) -> Result<EventStream, LlmError> { - 116
let next = self.responses.lock().unwrap().pop_front(); - 117
let (mut sink, rx) = stream::channel(64); - 118
match next { - 119
Some(m) => { - 120
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 121
sink.close_message(m).await; - 122
} - 123
None => { - 124
sink.close_error(LlmError::Parse("script exhausted".into())) - 125
.await - 126
} - 127
} - 128
Ok(rx) - 129
} - 130
} - 131
- 132
fn setup_with_primary( - 133
primary: Arc<dyn Provider>, - 134
ladder: Vec<(Arc<dyn Provider>, String)>, - 135
) -> (Agent, tempfile::TempDir) { - 136
let dir = tempdir().unwrap(); - 137
let cwd = dir.path().to_path_buf(); - 138
let header = SessionHeader { - 139
agent: None, - 140
session_id: "ladder".into(), - 141
created_at: chrono::Utc::now(), - 142
cwd: cwd.clone(), - 143
parent_session_id: None, - 144
contract_id: None, - 145
work_item_id: None, - 146
conversation: None, - 147
contract: FrozenContract { - 148
app_version: "0".into(), - 149
provider: "primary-net-dead".into(), - 150
model: "primary-model".into(), - 151
route_ladder: Vec::new(), - 152
route_objective: String::new(), - 153
route_annotations: Vec::new(), - 154
system_prompt: "sys".into(), - 155
permission_mode: "workspace-write".into(), - 156
capabilities: Vec::new(), - 157
prompt_layers: Vec::new(), - 158
}, - 159
}; - 160
let home = cwd.join(".vak-home"); - 161
std::fs::create_dir_all(&home).unwrap(); - 162
let log = - 163
SessionLog::create(SessionPath::new_session_file(&home, &cwd, "ladder"), header).unwrap(); - 164
let mut cfg = AgentConfig::new("sys"); - 165
cfg.model = "primary-model".into(); - 166
cfg.mode = Mode::FullAccess; - 167
cfg.permission = Some(Arc::new(PermissionEngine::default())); - 168
cfg.approver = Some(Arc::new(AutoApprove)); - 169
cfg.retry_base_backoff_ms = 1; - 170
cfg.run_retry_base_backoff_ms = 1; - 171
cfg.max_retries = 0; // one dispatch per leg keeps counting exact - 172
cfg.run_retry_attempts = 0; - 173
cfg.dispatch_ceiling = 3; - 174
cfg.ladder = ladder; - 175
cfg.ladder_provider_names = vec!["openrouter".into(); cfg.ladder.len()]; - 176
(Agent::new(primary, log, cfg), dir) - 177
} - 178
- 179
fn setup(ladder: Vec<(Arc<dyn Provider>, String)>) -> (Agent, tempfile::TempDir) { - 180
setup_with_primary( - 181
Arc::new(AlwaysNetwork { - 182
calls: AtomicU32::new(0), - 183
}), - 184
ladder, - 185
) - 186
} - 187
- 188
async fn run(agent: &mut Agent) -> TurnOutcome { - 189
let (ev_tx, mut ev_rx) = mpsc::channel(512); - 190
tokio::spawn(async move { while ev_rx.recv().await.is_some() {} }); - 191
let cancel = CancellationToken::new(); - 192
let steering = SteeringQueues::new(); - 193
agent.run("hello", &steering, cancel, ev_tx).await - 194
} - 195
- 196
fn fallback_provider() -> Arc<Scripted> { - 197
Arc::new(Scripted { - 198
responses: std::sync::Mutex::new(VecDeque::from(vec![text_msg("answered via fallback")])), - 199
}) - 200
} - 201
- 202
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 203
async fn primary_failure_walks_to_fallback_and_receipt_records_it() { - 204
let fb = fallback_provider(); - 205
let model_fb = "fallback-model".to_string(); - 206
let (mut agent, _dir) = setup(vec![(fb.clone() as Arc<dyn Provider>, model_fb.clone())]); - 207
let outcome = run(&mut agent).await; - 208
assert!( - 209
matches!(outcome, TurnOutcome::Completed { .. }), - 210
"fallback must rescue the step" - 211
); - 212
- 213
let session = agent.into_session().await; - 214
assert_eq!(session.receipts().len(), 1); - 215
let r = &session.receipts()[0]; - 216
assert_eq!(r.model, "fallback-model", "receipt names the winning leg"); - 217
assert_eq!( - 218
r.provider, "openrouter", - 219
"receipt preserves configured route identity" - 220
); - 221
assert_eq!(r.winning_attempt, Some(1)); - 222
assert_eq!(r.attempts[0].reason, AttemptReason::Initial); - 223
assert_eq!(r.attempts[0].domain, vak_llm::FailureDomain::Network); - 224
assert_eq!( - 225
r.attempts[1].reason, - 226
AttemptReason::RouteFallback, - 227
"first dispatch of the next leg is the fallback reason" - 228
); - 229
assert_eq!(r.attempts[1].settlement, vak_llm::Settlement::Ok); - 230
- 231
// Shared ceiling: exactly two paid dispatches consumed. - 232
// (asserted indirectly: a third dispatch would have exceeded ceiling=3? no) - 233
// Projection stays clean. - 234
assert_eq!(session.derive_messages().len(), 2); - 235
} - 236
- 237
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 238
async fn provider_panic_is_contained_and_fallback_completes() { - 239
let fb = fallback_provider(); - 240
let (mut agent, _dir) = setup_with_primary( - 241
Arc::new(PanickingProvider), - 242
vec![(fb as Arc<dyn Provider>, "fallback-model".to_string())], - 243
); - 244
- 245
assert!(matches!( - 246
run(&mut agent).await, - 247
TurnOutcome::Completed { .. } - 248
)); - 249
let session = agent.into_session().await; - 250
let receipt = &session.receipts()[0]; - 251
assert_eq!(receipt.attempts.len(), 2); - 252
assert_eq!(receipt.attempts[0].domain, vak_llm::FailureDomain::Network); - 253
assert_eq!(receipt.attempts[1].reason, AttemptReason::RouteFallback); - 254
assert_eq!(receipt.winning_attempt, Some(1)); - 255
} - 256
- 257
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 258
async fn terminal_quota_walks_to_fallback_without_retry_storm() { - 259
let fb = fallback_provider(); - 260
let (mut agent, _dir) = setup_with_primary( - 261
Arc::new(TerminalQuota), - 262
vec![(fb as Arc<dyn Provider>, "fallback-model".to_string())], - 263
); - 264
- 265
assert!(matches!( - 266
run(&mut agent).await, - 267
TurnOutcome::Completed { .. } - 268
)); - 269
let session = agent.into_session().await; - 270
let receipt = &session.receipts()[0]; - 271
assert_eq!(receipt.attempts.len(), 2); - 272
assert_eq!(receipt.attempts[0].reason, AttemptReason::Initial); - 273
assert_eq!(receipt.attempts[1].reason, AttemptReason::RouteFallback); - 274
assert_eq!(receipt.winning_attempt, Some(1)); - 275
} - 276
- 277
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 278
async fn all_legs_exhausted_fails_closed_within_ceiling() { - 279
// Fallback also fails (script empty => Parse error, non-retryable). - 280
let dead_fb = Arc::new(Scripted { - 281
responses: std::sync::Mutex::new(VecDeque::new()), - 282
}); - 283
let (mut agent, _dir) = setup(vec![( - 284
dead_fb as Arc<dyn Provider>, - 285
"fallback-model".to_string(), - 286
)]); - 287
match run(&mut agent).await { - 288
TurnOutcome::Failed { error } => { - 289
assert!(!error.to_string().is_empty()); - 290
} - 291
other => panic!("expected Failed after all legs, got {other:?}"), - 292
} - 293
let session = agent.into_session().await; - 294
let r = &session.receipts()[0]; - 295
assert_eq!(r.attempts.len(), 2, "one dispatch per leg at max_retries=0"); - 296
assert_eq!(r.attempts[1].reason, AttemptReason::RouteFallback); - 297
assert!(r.winning_attempt.is_none()); - 298
assert!( - 299
r.attempts - 300
.iter() - 301
.all(|a| a.settlement == vak_llm::Settlement::Unknown), - 302
"transport ambiguity stays UNKNOWN" - 303
); - 304
} - 305
- 306
/// Phase R: per-attempt leg attribution survives a cross-model walk, and - 307
/// consumers see one RouteFallback event per leg change. - 308
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 309
async fn fallback_legs_carry_provider_attribution_and_events() { - 310
let fb = fallback_provider(); - 311
let (mut agent, _dir) = setup(vec![( - 312
fb as Arc<dyn Provider>, - 313
"fallback-model".to_string(), - 314
)]); - 315
- 316
let (ev_tx, mut ev_rx) = mpsc::channel(512); - 317
let events = tokio::spawn(async move { - 318
let mut fallbacks = Vec::new(); - 319
while let Some(ev) = ev_rx.recv().await { - 320
if let vak_agent::AgentEvent::RouteFallback { - 321
to_provider, - 322
to_model, - 323
} = ev - 324
{ - 325
fallbacks.push((to_provider, to_model)); - 326
} - 327
} - 328
fallbacks - 329
}); - 330
- 331
let cancel = CancellationToken::new(); - 332
let steering = SteeringQueues::new(); - 333
let outcome = agent.run("hello", &steering, cancel, ev_tx).await; - 334
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 335
- 336
let session = agent.into_session().await; - 337
let r = &session.receipts()[0]; - 338
// Receipt-level stamp names the WINNING leg. - 339
assert_eq!( - 340
(r.provider.as_str(), r.model.as_str()), - 341
("openrouter", "fallback-model") - 342
); - 343
// Attempt 0 keeps its own (failed) leg attribution. - 344
assert_eq!( - 345
r.attempt_leg(&r.attempts[0]), - 346
("primary-net-dead", "primary-model") - 347
); - 348
// Attempt 1 attributes to the serving fallback leg. - 349
assert_eq!( - 350
r.attempt_leg(&r.attempts[1]), - 351
("openrouter", "fallback-model") - 352
); - 353
- 354
let fallbacks = events.await.unwrap(); - 355
assert_eq!( - 356
fallbacks, - 357
vec![("openrouter".into(), "fallback-model".into())], - 358
"exactly one RouteFallback event for the single leg change" - 359
); - 360
} - 361
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.