- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
use std::sync::{Arc, Mutex}; - 4
- 5
use tokio::sync::mpsc; - 6
use tokio_util::sync::CancellationToken; - 7
- 8
use tempfile::tempdir; - 9
- 10
use vak_agent::{Agent, AgentConfig, AgentEvent, TurnOutcome}; - 11
use vak_llm::stream; - 12
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; - 13
use vak_llm::{EventStream, LlmError, Provider}; - 14
use vak_session::SessionLog; - 15
use vak_session::types::{FrozenContract, SessionHeader}; - 16
- 17
/// Fails the first `failures` calls with `error`, then succeeds. - 18
struct FlakyThenGood { - 19
calls: Arc<Mutex<u32>>, - 20
failures: u32, - 21
error: LlmError, - 22
} - 23
- 24
#[async_trait::async_trait] - 25
impl Provider for FlakyThenGood { - 26
fn name(&self) -> &str { - 27
"flaky-then-good" - 28
} - 29
- 30
async fn stream( - 31
&self, - 32
_request: ChatRequest, - 33
_cancel: CancellationToken, - 34
) -> Result<EventStream, LlmError> { - 35
let n = { - 36
let mut c = self.calls.lock().unwrap(); - 37
*c += 1; - 38
*c - 39
}; - 40
let (mut sink, rx) = stream::channel(8); - 41
if n <= self.failures { - 42
sink.close_error(self.error.clone()).await; - 43
} else { - 44
let msg = AssistantMessage { - 45
content: vec![ContentBlock::text("recovered")], - 46
stop_reason: StopReason::EndTurn, - 47
usage: Usage::default(), - 48
model: "test-model".into(), - 49
response_id: None, - 50
}; - 51
sink.push(stream::StreamEvent::Start { - 52
partial: msg.clone(), - 53
}); - 54
sink.close_message(msg).await; - 55
} - 56
Ok(rx) - 57
} - 58
} - 59
- 60
fn build_agent(provider: Arc<dyn Provider>, session_id: &str, attempts: u32) -> Agent { - 61
let dir = tempdir().unwrap(); - 62
let header = SessionHeader { - 63
agent: None, - 64
session_id: session_id.into(), - 65
created_at: chrono::Utc::now(), - 66
cwd: dir.path().to_path_buf(), - 67
parent_session_id: None, - 68
contract_id: None, - 69
work_item_id: None, - 70
conversation: None, - 71
contract: FrozenContract { - 72
app_version: "0".into(), - 73
provider: "scripted".into(), - 74
model: "test-model".into(), - 75
route_ladder: Vec::new(), - 76
route_objective: String::new(), - 77
route_annotations: Vec::new(), - 78
system_prompt: "sys".into(), - 79
permission_mode: "full-access".into(), - 80
capabilities: Vec::new(), - 81
prompt_layers: Vec::new(), - 82
}, - 83
}; - 84
let log = SessionLog::create(dir.path().join("s.jsonl"), header).unwrap(); - 85
let mut cfg = AgentConfig::new("sys"); - 86
cfg.max_retries = 0; - 87
cfg.retry_base_backoff_ms = 1; - 88
cfg.run_retry_attempts = attempts; - 89
cfg.run_retry_base_backoff_ms = 1; - 90
std::mem::forget(dir); - 91
Agent::new(provider, log, cfg) - 92
} - 93
- 94
async fn run_agent(agent: &mut Agent) -> (TurnOutcome, Vec<AgentEvent>) { - 95
let (tx, mut rx) = mpsc::channel(256); - 96
let cancel = CancellationToken::new(); - 97
let outcome = agent - 98
.run("go", &Default::default(), cancel, tx.clone()) - 99
.await; - 100
drop(tx); - 101
let mut events = Vec::new(); - 102
while let Ok(ev) = rx.try_recv() { - 103
events.push(ev); - 104
} - 105
(outcome, events) - 106
} - 107
- 108
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] - 109
async fn endurance_bridges_transient_windows_then_completes() { - 110
for err in [ - 111
LlmError::Overloaded("storm".into()), - 112
LlmError::RateLimit { - 113
message: "429".into(), - 114
retry_after_secs: None, - 115
}, - 116
LlmError::Network("connection reset".into()), - 117
LlmError::Parse("stream closed before finish_reason".into()), - 118
] { - 119
let calls = Arc::new(Mutex::new(0u32)); - 120
let provider = Arc::new(FlakyThenGood { - 121
calls: calls.clone(), - 122
failures: 3, - 123
error: err.clone(), - 124
}); - 125
let mut agent = build_agent(provider, "endure", 5); - 126
let (outcome, events) = run_agent(&mut agent).await; - 127
assert!( - 128
matches!(outcome, TurnOutcome::Completed { .. }), - 129
"expected completion after endurance, got {outcome:?} ({err:?})" - 130
); - 131
assert_eq!( - 132
*calls.lock().unwrap(), - 133
4, - 134
"3 failures + 1 success ({err:?})" - 135
); - 136
let retries = events - 137
.iter() - 138
.filter(|e| matches!(e, AgentEvent::RetryScheduled { .. })) - 139
.count(); - 140
assert_eq!( - 141
retries, 3, - 142
"run-level re-attempts surfaced as events ({err:?})" - 143
); - 144
} - 145
} - 146
- 147
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] - 148
async fn endurance_budget_exhaustion_fails_cleanly() { - 149
let calls = Arc::new(Mutex::new(0u32)); - 150
let provider = Arc::new(FlakyThenGood { - 151
calls: calls.clone(), - 152
failures: 10, - 153
error: LlmError::Overloaded("long storm".into()), - 154
}); - 155
let mut agent = build_agent(provider, "exhaust", 2); - 156
let (outcome, _) = run_agent(&mut agent).await; - 157
match outcome { - 158
TurnOutcome::Failed { error } => { - 159
assert!(error.to_string().contains("long storm")); - 160
} - 161
other => panic!("expected failed after budget exhausted, got {other:?}"), - 162
} - 163
assert_eq!(*calls.lock().unwrap(), 3, "initial + 2 run-level attempts"); - 164
} - 165
- 166
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] - 167
async fn permanent_errors_are_not_endured() { - 168
for err in [ - 169
LlmError::Auth("bad key".into()), - 170
LlmError::InvalidRequest("bad body".into()), - 171
LlmError::Api { - 172
status: 500, - 173
message: "server".into(), - 174
}, - 175
] { - 176
let calls = Arc::new(Mutex::new(0u32)); - 177
let provider = Arc::new(FlakyThenGood { - 178
calls: calls.clone(), - 179
failures: 99, - 180
error: err.clone(), - 181
}); - 182
let mut agent = build_agent(provider, "permanent", 5); - 183
let (outcome, events) = run_agent(&mut agent).await; - 184
assert!(matches!(outcome, TurnOutcome::Failed { .. }), "{err:?}"); - 185
assert_eq!(*calls.lock().unwrap(), 1, "no run-level retry for {err:?}"); - 186
assert!( - 187
!events - 188
.iter() - 189
.any(|e| matches!(e, AgentEvent::RetryScheduled { .. })), - 190
"no RetryScheduled for {err:?}" - 191
); - 192
} - 193
} - 194
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.