- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
//! Work receipts + dispatch ceiling (docs/design/42-managed-work-contracts.md): every provider - 4
//! dispatch lands as a typed ledger entry on every exit path; the ceiling - 5
//! fails closed without another paid call. - 6
- 7
use std::path::Path; - 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::{ - 19
AttemptReason, EventStream, FailureDomain, LlmError, Provider, Settlement, WorkPurpose, - 20
}; - 21
use vak_permission::{Mode, PermissionEngine}; - 22
use vak_session::types::{EntryPayload, FrozenContract, SessionHeader}; - 23
use vak_session::{SessionLog, SessionPath}; - 24
- 25
fn text_msg(t: &str) -> AssistantMessage { - 26
AssistantMessage { - 27
content: vec![ContentBlock::text(t)], - 28
stop_reason: StopReason::EndTurn, - 29
usage: Usage { - 30
input_tokens: 5, - 31
output_tokens: 3, - 32
..Default::default() - 33
}, - 34
model: "test-model".into(), - 35
response_id: None, - 36
} - 37
} - 38
- 39
/// Always fails with a rate-limit error: informed transience feeding - 40
/// retries/endurance but never tripping the breaker. - 41
struct AlwaysRateLimited; - 42
- 43
#[async_trait::async_trait] - 44
impl Provider for AlwaysRateLimited { - 45
fn name(&self) -> &str { - 46
"always-429" - 47
} - 48
- 49
async fn stream( - 50
&self, - 51
_request: ChatRequest, - 52
_cancel: CancellationToken, - 53
) -> Result<EventStream, LlmError> { - 54
let (mut sink, rx) = stream::channel(64); - 55
sink.close_error(LlmError::RateLimit { - 56
message: "slow down".into(), - 57
retry_after_secs: None, - 58
}) - 59
.await; - 60
Ok(rx) - 61
} - 62
} - 63
- 64
/// Succeeds after N failed dispatches (network resets). - 65
struct FailThenSucceed { - 66
remaining_failures: AtomicU32, - 67
} - 68
- 69
#[async_trait::async_trait] - 70
impl Provider for FailThenSucceed { - 71
fn name(&self) -> &str { - 72
"flaky" - 73
} - 74
- 75
async fn stream( - 76
&self, - 77
_request: ChatRequest, - 78
_cancel: CancellationToken, - 79
) -> Result<EventStream, LlmError> { - 80
let (mut sink, rx) = stream::channel(64); - 81
if self.remaining_failures.fetch_sub(1, Ordering::SeqCst) > 0 { - 82
sink.close_error(LlmError::Network("conn reset".into())) - 83
.await; - 84
} else { - 85
let m = text_msg("done"); - 86
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 87
sink.close_message(m).await; - 88
} - 89
Ok(rx) - 90
} - 91
} - 92
- 93
/// Aborts mid-stream with partial output. - 94
struct AbortMidStream; - 95
- 96
#[async_trait::async_trait] - 97
impl Provider for AbortMidStream { - 98
fn name(&self) -> &str { - 99
"aborting" - 100
} - 101
- 102
async fn stream( - 103
&self, - 104
_request: ChatRequest, - 105
_cancel: CancellationToken, - 106
) -> Result<EventStream, LlmError> { - 107
let (mut sink, rx) = stream::channel(64); - 108
sink.close_error(LlmError::Aborted { - 109
partial: Some(Box::new(text_msg("partial answer"))), - 110
}) - 111
.await; - 112
Ok(rx) - 113
} - 114
} - 115
- 116
fn session_paths(dir: &Path) -> (std::path::PathBuf, std::path::PathBuf) { - 117
let cwd = dir.to_path_buf(); - 118
let home = cwd.join(".vak-home"); - 119
( - 120
home.clone(), - 121
SessionPath::new_session_file(&home, &cwd, "receipts"), - 122
) - 123
} - 124
- 125
fn setup_with( - 126
provider: Arc<dyn Provider>, - 127
customize: impl FnOnce(&mut AgentConfig), - 128
) -> (Agent, tempfile::TempDir, std::path::PathBuf) { - 129
let dir = tempdir().unwrap(); - 130
let (home, path) = session_paths(dir.path()); - 131
std::fs::create_dir_all(&home).unwrap(); - 132
let header = SessionHeader { - 133
agent: None, - 134
session_id: "receipts".into(), - 135
created_at: chrono::Utc::now(), - 136
cwd: dir.path().to_path_buf(), - 137
parent_session_id: None, - 138
contract_id: None, - 139
work_item_id: None, - 140
conversation: None, - 141
contract: FrozenContract { - 142
app_version: "0".into(), - 143
provider: provider.name().to_string(), - 144
model: "test-model".into(), - 145
route_ladder: Vec::new(), - 146
route_objective: String::new(), - 147
route_annotations: Vec::new(), - 148
system_prompt: "sys".into(), - 149
permission_mode: "workspace-write".into(), - 150
capabilities: Vec::new(), - 151
prompt_layers: Vec::new(), - 152
}, - 153
}; - 154
let log = SessionLog::create(path, header).unwrap(); - 155
let mut cfg = AgentConfig::new("sys"); - 156
cfg.model = "test-model".into(); - 157
cfg.mode = Mode::FullAccess; - 158
cfg.permission = Some(Arc::new(PermissionEngine::default())); - 159
cfg.approver = Some(Arc::new(AutoApprove)); - 160
// Keep tests fast regardless of backoff math. - 161
cfg.retry_base_backoff_ms = 1; - 162
cfg.run_retry_base_backoff_ms = 1; - 163
customize(&mut cfg); - 164
(Agent::new(provider, log, cfg), dir, home) - 165
} - 166
- 167
async fn run(agent: &mut Agent) -> TurnOutcome { - 168
let (ev_tx, mut ev_rx) = mpsc::channel(256); - 169
tokio::spawn(async move { while ev_rx.recv().await.is_some() {} }); - 170
let cancel = CancellationToken::new(); - 171
let steering = SteeringQueues::new(); - 172
agent.run("hello", &steering, cancel, ev_tx).await - 173
} - 174
- 175
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 176
async fn success_step_writes_execute_receipt_and_projection_ignores_it() { - 177
let (mut agent, _dir, _home) = setup_with( - 178
Arc::new(FailThenSucceed { - 179
remaining_failures: AtomicU32::new(1), - 180
}), - 181
|_| {}, - 182
); - 183
let outcome = run(&mut agent).await; - 184
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 185
- 186
let session = agent.into_session().await; - 187
- 188
// Model-visible means logged — and ONLY that: receipts never enter the - 189
// projection. - 190
assert_eq!(session.derive_messages().len(), 2, "prompt + reply only"); - 191
- 192
let receipts = session.receipts(); - 193
assert_eq!(receipts.len(), 1, "one work unit => one receipt"); - 194
let r = receipts[0]; - 195
assert_eq!(r.purpose, WorkPurpose::Execute); - 196
assert_eq!(r.attempts.len(), 2, "network failure, then success"); - 197
assert_eq!(r.attempts[0].reason, AttemptReason::Initial); - 198
assert_eq!(r.attempts[0].domain, FailureDomain::Network); - 199
assert_eq!(r.attempts[0].settlement, Settlement::Unknown); - 200
assert_eq!(r.attempts[1].reason, AttemptReason::Retry); - 201
assert_eq!(r.attempts[1].settlement, Settlement::Ok); - 202
assert_eq!(r.winning_attempt, Some(1)); - 203
assert_eq!( - 204
r.attempts[1].usage.as_ref().map(|u| u.input_tokens), - 205
Some(5) - 206
); - 207
} - 208
- 209
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 210
async fn ceiling_exhaustion_fails_closed_with_failure_receipt() { - 211
let (mut agent, _dir, _home) = setup_with(Arc::new(AlwaysRateLimited), |cfg| { - 212
cfg.max_retries = 1; - 213
cfg.run_retry_attempts = 0; - 214
cfg.dispatch_ceiling = 2; - 215
}); - 216
let outcome = run(&mut agent).await; - 217
match outcome { - 218
TurnOutcome::Failed { error } => { - 219
assert!( - 220
error.to_string().contains("dispatch ceiling of 2"), - 221
"unexpected error: {error}" - 222
); - 223
} - 224
other => panic!("expected Failed at ceiling, got {other:?}"), - 225
} - 226
- 227
let session = agent.into_session().await; - 228
let receipts = session.receipts(); - 229
assert_eq!(receipts.len(), 1); - 230
let r = receipts[0]; - 231
assert_eq!(r.attempts.len(), 2, "initial + one retry, then closed"); - 232
assert!(r.winning_attempt.is_none()); - 233
assert!( - 234
r.attempts - 235
.iter() - 236
.all(|a| a.domain == FailureDomain::Account && a.settlement == Settlement::Failed) - 237
); - 238
} - 239
- 240
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 241
async fn mid_stream_abort_records_cancelled_settlement_and_keeps_partial() { - 242
let (mut agent, _dir, _home) = setup_with(Arc::new(AbortMidStream), |_| {}); - 243
let outcome = run(&mut agent).await; - 244
assert!(matches!(outcome, TurnOutcome::Aborted { .. })); - 245
- 246
let session = agent.into_session().await; - 247
let messages = session.derive_messages(); - 248
assert!( - 249
messages - 250
.iter() - 251
.any(|m| m.text_content().contains("partial answer")), - 252
"partial output survives abort (invariant 5)" - 253
); - 254
- 255
let receipts = session.receipts(); - 256
assert_eq!(receipts.len(), 1); - 257
let attempts = &receipts[0].attempts; - 258
assert_eq!(attempts.len(), 1); - 259
assert_eq!(attempts[0].settlement, Settlement::Cancelled); - 260
assert!(receipts[0].winning_attempt.is_none()); - 261
} - 262
- 263
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 264
async fn receipts_round_trip_through_disk() { - 265
let (mut agent, dir, _home) = setup_with( - 266
Arc::new(FailThenSucceed { - 267
remaining_failures: AtomicU32::new(0), - 268
}), - 269
|_| {}, - 270
); - 271
let outcome = run(&mut agent).await; - 272
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 273
- 274
let (_, path) = session_paths(dir.path()); - 275
let session = agent.into_session().await; - 276
let kinds: Vec<&str> = session - 277
.chain_to_root() - 278
.iter() - 279
.map(|e| match e.payload { - 280
EntryPayload::Header(_) => "header", - 281
EntryPayload::Message(_) => "message", - 282
EntryPayload::Compaction(_) => "compaction", - 283
EntryPayload::Receipt(_) => "receipt", - 284
EntryPayload::Goal(_) => "goal", - 285
EntryPayload::GoalUpdate(_) => "goal-update", - 286
EntryPayload::Activity(_) => "activity", - 287
EntryPayload::Work(_) => "work", - 288
EntryPayload::Intent(_) => "intent", - 289
EntryPayload::TurnCapabilitiesBound(_) | EntryPayload::TurnCapabilitiesRef(_) => { - 290
"capabilities" - 291
} - 292
EntryPayload::ChildRun { .. } => "child-run", - 293
EntryPayload::Presentation(_) => "presentation", - 294
EntryPayload::TurnCard(_) => "turn-card", - 295
EntryPayload::EvidenceBody(_) => "evidence-body", - 296
}) - 297
.collect(); - 298
assert_eq!( - 299
kinds, - 300
vec!["header", "message", "receipt", "message", "turn-card"], - 301
"receipt sits between the prompt and the reply it produced; the turn-close hook \ - 302
appends the TurnCard once the reply lands" - 303
); - 304
drop(session); - 305
- 306
let reopened = SessionLog::open(path).unwrap(); - 307
assert_eq!(reopened.receipts().len(), 1, "audit survives restart"); - 308
assert_eq!(reopened.derive_messages().len(), 2, "projection unchanged"); - 309
assert!(reopened.warnings().is_empty()); - 310
} - 311
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.