- 1
//! Work receipts: typed audit records for provider dispatches (docs/design/42-managed-work-contracts.md//! Phase A). Receipts are ledger data, never model-visible input. - 2
- 3
use crate::{LlmError, Usage}; - 4
use serde::{Deserialize, Serialize}; - 5
use std::time::Instant; - 6
- 7
/// Why this work exists. Extends as call sites adopt receipts; every - 8
/// variant must map to a caller-visible purpose, never a hidden retry. - 9
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] - 10
#[serde(rename_all = "snake_case")] - 11
pub enum WorkPurpose { - 12
/// Structured managed-work contract authoring. - 13
Plan, - 14
/// A main agent-loop model step. - 15
Execute, - 16
/// The compaction summarizer call. - 17
Summarize, - 18
/// Completion-audit judge call (Phase H). - 19
Verify, - 20
/// An intent-classification call (docs/design/47-commitment-kernel.md, - 21
/// resolver tiers 2 and 3). - 22
Classify, - 23
/// A Gemini Live API text-to-speech dispatch (voice/personality). - 24
VoiceSynthesis, - 25
/// Streaming or batch speech-to-text dispatch. - 26
SpeechRecognition, - 27
} - 28
- 29
/// Why a dispatch was made. `Retry` covers same-candidate transient - 30
/// retries; route-level reasons arrive with the frozen ladder (Phase B). - 31
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] - 32
#[serde(rename_all = "snake_case")] - 33
pub enum AttemptReason { - 34
Initial, - 35
Retry, - 36
/// First dispatch of the NEXT frozen-ladder candidate after typed - 37
// failure of the previous one (Phase B). - 38
RouteFallback, - 39
EnduranceRetry, - 40
} - 41
- 42
/// Which slice of the world failed. Consumed by breaker/endurance - 43
/// classification instead of error-string matching. - 44
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] - 45
#[serde(rename_all = "snake_case")] - 46
pub enum FailureDomain { - 47
/// Credentials, quota, rate limits: account-scoped facts. - 48
Account, - 49
/// The provider service itself (overload, outage). - 50
Provider, - 51
/// The model produced unusable output (malformed stream). - 52
Model, - 53
/// Our request was rejected before generation (bad request, auth). - 54
Request, - 55
/// Transport-level failure. - 56
Network, - 57
/// Watchdog deadline. - 58
Deadline, - 59
/// Local context admission/compaction failure. - 60
Context, - 61
/// Unclassified. - 62
Unknown, - 63
} - 64
- 65
/// What actually happened to a dispatch. Billing without a rated outcome - 66
/// is `Unknown`, never `Ok` — absent is UNKNOWN, never zero. - 67
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] - 68
#[serde(rename_all = "snake_case")] - 69
pub enum Settlement { - 70
Ok, - 71
Failed, - 72
Cancelled, - 73
Unknown, - 74
} - 75
- 76
#[derive(Debug, Clone, Serialize, Deserialize)] - 77
pub struct DispatchAttempt { - 78
pub ordinal: u32, - 79
pub reason: AttemptReason, - 80
pub domain: FailureDomain, - 81
pub settlement: Settlement, - 82
pub latency_ms: u64, - 83
#[serde(default, skip_serializing_if = "Option::is_none")] - 84
pub usage: Option<Usage>, - 85
#[serde(default, skip_serializing_if = "Option::is_none")] - 86
pub error: Option<String>, - 87
/// Per-attempt leg attribution override. A receipt walks multiple - 88
/// frozen-ladder legs, so the receipt-level provider/model names only - 89
/// the FINAL leg; fallback legs must be attributed to what actually - 90
/// failed over FROM. None ⇒ attribute to the receipt-level fields - 91
/// (single-leg receipts, and receipts written before per-leg - 92
/// attribution existed — which stay readable permanently, per the - 93
/// additive-only contract in docs/design/46 VII.3). - 94
#[serde(default, skip_serializing_if = "Option::is_none")] - 95
pub provider: Option<String>, - 96
#[serde(default, skip_serializing_if = "Option::is_none")] - 97
pub model: Option<String>, - 98
} - 99
- 100
#[derive(Debug, Clone, Serialize, Deserialize)] - 101
pub struct WorkReceipt { - 102
pub purpose: WorkPurpose, - 103
#[serde(default)] - 104
pub provider: String, - 105
pub model: String, - 106
/// Ordinal of the attempt that produced committed output; None when no - 107
/// attempt succeeded. - 108
#[serde(default, skip_serializing_if = "Option::is_none")] - 109
pub winning_attempt: Option<u32>, - 110
pub attempts: Vec<DispatchAttempt>, - 111
/// SHA-256 hex digest of the stable prefix (system prompt + tool - 112
/// schemas) this dispatch sent, per `vak_context::assemble::prefix_digest` - 113
/// (docs/design/68-context-engine.md §4). Empty for a receipt that - 114
/// predates prefix tracking or a work purpose with no request prefix - 115
/// (e.g. voice synthesis). - 116
#[serde(default)] - 117
pub prefix_digest: String, - 118
/// Measured cost of the prefix in tokens: the provider's reported - 119
/// `usage.input_tokens` on the first request seen with this digest, - 120
/// minus an estimate of the messages alone — approximate, since the - 121
/// provider does not itemize prefix vs. messages in its usage report. - 122
/// `None` when not yet measured for this digest. - 123
#[serde(default, skip_serializing_if = "Option::is_none")] - 124
pub prefix_tokens: Option<u64>, - 125
} - 126
- 127
impl WorkReceipt { - 128
pub fn new( - 129
purpose: WorkPurpose, - 130
provider: impl Into<String>, - 131
model: impl Into<String>, - 132
) -> Self { - 133
WorkReceipt { - 134
purpose, - 135
provider: provider.into(), - 136
model: model.into(), - 137
winning_attempt: None, - 138
attempts: Vec::new(), - 139
prefix_digest: String::new(), - 140
prefix_tokens: None, - 141
} - 142
} - 143
- 144
/// Restamp the receipt for a frozen-ladder leg about to dispatch. The - 145
/// previous leg's attempts keep their stamped attribution via the - 146
/// per-attempt overrides recorded alongside them. - 147
pub fn stamp_leg(&mut self, provider: &str, model: &str) { - 148
for a in &mut self.attempts { - 149
if a.provider.is_none() { - 150
let prev = std::mem::take(&mut self.provider); - 151
a.provider = Some(prev); - 152
let prev_model = std::mem::take(&mut self.model); - 153
a.model = Some(prev_model); - 154
} - 155
} - 156
self.provider = provider.to_string(); - 157
self.model = model.to_string(); - 158
} - 159
- 160
/// Effective (provider, model) attribution for one attempt. - 161
pub fn attempt_leg<'a>(&'a self, a: &'a DispatchAttempt) -> (&'a str, &'a str) { - 162
( - 163
a.provider.as_deref().unwrap_or(self.provider.as_str()), - 164
a.model.as_deref().unwrap_or(self.model.as_str()), - 165
) - 166
} - 167
- 168
pub fn record( - 169
&mut self, - 170
reason: AttemptReason, - 171
domain: FailureDomain, - 172
settlement: Settlement, - 173
latency_ms: u64, - 174
usage: Option<Usage>, - 175
error: Option<String>, - 176
) { - 177
let ordinal = self.attempts.len() as u32; - 178
if settlement == Settlement::Ok { - 179
self.winning_attempt = Some(ordinal); - 180
} - 181
self.attempts.push(DispatchAttempt { - 182
ordinal, - 183
reason, - 184
domain, - 185
settlement, - 186
latency_ms, - 187
usage, - 188
error, - 189
provider: None, - 190
model: None, - 191
}); - 192
} - 193
- 194
pub fn settle_cancelled(&mut self) { - 195
if let Some(last) = self.attempts.last_mut() { - 196
last.settlement = Settlement::Cancelled; - 197
} - 198
} - 199
} - 200
- 201
/// Classify an error into (domain, settlement). Pre-dispatch rejections - 202
/// (auth, bad request, rate limit, overload) deterministically failed; - 203
/// transport/mid-stream failures may have consumed provider compute, so - 204
/// their settlement is Unknown — fail-closed against paid fallback later. - 205
pub fn classify_error(e: &LlmError) -> (FailureDomain, Settlement) { - 206
match e { - 207
LlmError::Auth(_) => (FailureDomain::Account, Settlement::Failed), - 208
LlmError::RateLimit { .. } => (FailureDomain::Account, Settlement::Failed), - 209
LlmError::Overloaded(_) => (FailureDomain::Provider, Settlement::Failed), - 210
LlmError::InvalidRequest(_) => (FailureDomain::Request, Settlement::Failed), - 211
LlmError::Api { .. } => (FailureDomain::Request, Settlement::Unknown), - 212
LlmError::Network(_) => (FailureDomain::Network, Settlement::Unknown), - 213
LlmError::Parse(_) => (FailureDomain::Model, Settlement::Unknown), - 214
LlmError::Context(_) => (FailureDomain::Context, Settlement::Failed), - 215
LlmError::Aborted { .. } => (FailureDomain::Unknown, Settlement::Cancelled), - 216
} - 217
} - 218
- 219
/// Hard cap on provider dispatches for one unit of work. Exhaustion fails - 220
/// closed before another paid call goes out. Single-ladder default in the - 221
/// agent codifies today's worst case: - 222
/// `(max_retries + 1) * (run_retry_attempts + 1)`; Phase B tightens the - 223
/// formula to `ladder.len() + repair_allowance` once ladders exist. - 224
#[derive(Debug, Clone)] - 225
pub struct DispatchBudget { - 226
limit: u32, - 227
used: u32, - 228
} - 229
- 230
impl DispatchBudget { - 231
pub fn new(limit: u32) -> Self { - 232
DispatchBudget { limit, used: 0 } - 233
} - 234
- 235
pub fn remaining(&self) -> u32 { - 236
self.limit.saturating_sub(self.used) - 237
} - 238
- 239
pub fn used(&self) -> u32 { - 240
self.used - 241
} - 242
- 243
pub fn limit(&self) -> u32 { - 244
self.limit - 245
} - 246
- 247
/// Consume one dispatch. Err when the ceiling is exhausted: callers - 248
/// must not dispatch past it. - 249
pub fn consume(&mut self) -> Result<(), DispatchCeiling> { - 250
if self.used >= self.limit { - 251
return Err(DispatchCeiling { limit: self.limit }); - 252
} - 253
self.used += 1; - 254
Ok(()) - 255
} - 256
} - 257
- 258
#[derive(Debug, Clone, Copy, PartialEq, Eq)] - 259
pub struct DispatchCeiling { - 260
pub limit: u32, - 261
} - 262
- 263
impl std::fmt::Display for DispatchCeiling { - 264
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - 265
write!(f, "dispatch ceiling of {} exhausted", self.limit) - 266
} - 267
} - 268
- 269
/// Per-work accounting threaded through the reliability helper: the shared - 270
/// ceiling plus the receipt under construction. - 271
pub struct StepLedger { - 272
pub budget: DispatchBudget, - 273
pub receipt: WorkReceipt, - 274
/// Wall-clock time from request send to the first `StreamEvent` of the - 275
/// winning attempt, measured by the caller's stream-consuming loop - 276
/// (docs/design/68-context-engine.md §1 "Feedback": `prefill_tps` needs - 277
/// this on providers, like OpenAI/Anthropic, that do not report their - 278
/// own prefill duration the way Ollama does). `None` until a dispatch - 279
/// succeeds; overwritten by each new attempt, never accumulated. - 280
pub last_first_token_ms: Option<u64>, - 281
} - 282
- 283
impl StepLedger { - 284
pub fn new(purpose: WorkPurpose, provider: &str, model: &str, ceiling: u32) -> Self { - 285
StepLedger { - 286
budget: DispatchBudget::new(ceiling), - 287
receipt: WorkReceipt::new(purpose, provider, model), - 288
last_first_token_ms: None, - 289
} - 290
} - 291
- 292
/// Hand the finished receipt to the caller, leaving an empty shell in - 293
/// place (used when a ledger outlives its first work unit's write). - 294
pub fn take_receipt(&mut self) -> WorkReceipt { - 295
let purpose = self.receipt.purpose; - 296
let provider = self.receipt.provider.clone(); - 297
let model = self.receipt.model.clone(); - 298
std::mem::replace( - 299
&mut self.receipt, - 300
WorkReceipt::new(purpose, provider, model), - 301
) - 302
} - 303
- 304
/// Time one dispatch and record its outcome from start to end. - 305
pub async fn timed( - 306
&mut self, - 307
reason: AttemptReason, - 308
f: impl std::future::Future<Output = Result<crate::AssistantMessage, LlmError>>, - 309
) -> Result<crate::AssistantMessage, LlmError> { - 310
if let Err(c) = self.budget.consume() { - 311
return Err(LlmError::Network(c.to_string())); - 312
} - 313
let started = Instant::now(); - 314
let outcome = f.await; - 315
match &outcome { - 316
Ok(value) => { - 317
self.receipt.record( - 318
reason, - 319
FailureDomain::Unknown, - 320
Settlement::Ok, - 321
started.elapsed().as_millis() as u64, - 322
Some(value.usage.clone()), - 323
None, - 324
); - 325
} - 326
Err(e) => { - 327
// Deadline classification happens at the watchdog site via - 328
// `record_deadline`; everything else classifies here. - 329
if !matches!(e, LlmError::Aborted { .. }) { - 330
let (domain, settlement) = classify_error(e); - 331
self.receipt.record( - 332
reason, - 333
domain, - 334
settlement, - 335
started.elapsed().as_millis() as u64, - 336
None, - 337
Some(e.to_string()), - 338
); - 339
} - 340
} - 341
} - 342
outcome - 343
} - 344
- 345
/// Record a watchdog-deadline failure explicitly (the timeout wrapper - 346
/// erases the inner future's result, so classification must be manual). - 347
pub fn record_deadline(&mut self, reason: AttemptReason, secs: u64, err: &LlmError) { - 348
self.receipt.record( - 349
reason, - 350
FailureDomain::Deadline, - 351
Settlement::Unknown, - 352
secs * 1000, - 353
None, - 354
Some(err.to_string()), - 355
); - 356
} - 357
} - 358
- 359
#[cfg(test)] - 360
mod tests { - 361
#![allow(clippy::unwrap_used, clippy::expect_used)] - 362
use super::*; - 363
- 364
#[test] - 365
fn classify_maps_domains_and_settlements() { - 366
assert_eq!( - 367
classify_error(&LlmError::RateLimit { - 368
message: "slow down".into(), - 369
retry_after_secs: Some(3), - 370
}), - 371
(FailureDomain::Account, Settlement::Failed) - 372
); - 373
assert_eq!( - 374
classify_error(&LlmError::Parse("truncated".into())), - 375
(FailureDomain::Model, Settlement::Unknown) - 376
); - 377
assert_eq!( - 378
classify_error(&LlmError::Network("conn reset".into())), - 379
(FailureDomain::Network, Settlement::Unknown) - 380
); - 381
assert_eq!( - 382
classify_error(&LlmError::Auth("bad key".into())), - 383
(FailureDomain::Account, Settlement::Failed) - 384
); - 385
assert_eq!( - 386
classify_error(&LlmError::Context("over budget".into())), - 387
(FailureDomain::Context, Settlement::Failed) - 388
); - 389
} - 390
- 391
#[test] - 392
fn budget_fails_closed_at_limit() { - 393
let mut b = DispatchBudget::new(2); - 394
assert_eq!(b.remaining(), 2); - 395
assert!(b.consume().is_ok()); - 396
assert!(b.consume().is_ok()); - 397
assert_eq!(b.remaining(), 0); - 398
let err = b.consume().unwrap_err(); - 399
assert_eq!(err.limit, 2); - 400
assert_eq!(err.to_string(), "dispatch ceiling of 2 exhausted"); - 401
} - 402
- 403
#[test] - 404
fn receipt_tracks_winning_ordinal_and_cancellation() { - 405
let mut r = WorkReceipt::new(WorkPurpose::Execute, "p", "m"); - 406
r.record( - 407
AttemptReason::Initial, - 408
FailureDomain::Account, - 409
Settlement::Failed, - 410
10, - 411
None, - 412
Some("429".into()), - 413
); - 414
r.record( - 415
AttemptReason::EnduranceRetry, - 416
FailureDomain::Deadline, - 417
Settlement::Unknown, - 418
20, - 419
None, - 420
Some("deadline".into()), - 421
); - 422
assert_eq!(r.attempts.len(), 2); - 423
assert!(r.winning_attempt.is_none()); - 424
r.record( - 425
AttemptReason::EnduranceRetry, - 426
FailureDomain::Unknown, - 427
Settlement::Ok, - 428
30, - 429
Some(Usage::default()), - 430
None, - 431
); - 432
assert_eq!(r.winning_attempt, Some(2)); - 433
r.settle_cancelled(); - 434
assert_eq!(r.attempts[2].settlement, Settlement::Cancelled); - 435
} - 436
- 437
#[test] - 438
fn receipt_json_round_trip_preserves_everything() { - 439
let mut r = WorkReceipt::new(WorkPurpose::Summarize, "anthropic", "claude-x"); - 440
r.record( - 441
AttemptReason::Initial, - 442
FailureDomain::Network, - 443
Settlement::Unknown, - 444
1234, - 445
None, - 446
Some("reset by peer".into()), - 447
); - 448
r.record( - 449
AttemptReason::Retry, - 450
FailureDomain::Unknown, - 451
Settlement::Ok, - 452
4321, - 453
Some(Usage { - 454
input_tokens: 11, - 455
output_tokens: 7, - 456
..Default::default() - 457
}), - 458
None, - 459
); - 460
let json = serde_json::to_string(&r).unwrap(); - 461
let back: WorkReceipt = serde_json::from_str(&json).unwrap(); - 462
assert_eq!(back.purpose, WorkPurpose::Summarize); - 463
assert_eq!(back.provider, "anthropic"); - 464
assert_eq!(back.model, "claude-x"); - 465
assert_eq!(back.winning_attempt, Some(1)); - 466
assert_eq!(back.attempts.len(), 2); - 467
assert_eq!(back.attempts[0].domain, FailureDomain::Network); - 468
assert_eq!(back.attempts[0].settlement, Settlement::Unknown); - 469
assert_eq!( - 470
back.attempts[1].usage.as_ref().map(|u| u.output_tokens), - 471
Some(7) - 472
); - 473
} - 474
- 475
#[test] - 476
fn stamp_leg_preserves_prior_leg_attribution() { - 477
let mut r = WorkReceipt::new(WorkPurpose::Execute, "anthropic", "claude-x"); - 478
r.record( - 479
AttemptReason::Initial, - 480
FailureDomain::Provider, - 481
Settlement::Failed, - 482
10, - 483
None, - 484
Some("overload".into()), - 485
); - 486
r.stamp_leg("openai", "gpt-x"); - 487
r.record( - 488
AttemptReason::RouteFallback, - 489
FailureDomain::Network, - 490
Settlement::Ok, - 491
20, - 492
None, - 493
None, - 494
); - 495
assert_eq!(r.provider, "openai"); - 496
assert_eq!(r.model, "gpt-x"); - 497
let (p0, m0) = r.attempt_leg(&r.attempts[0]); - 498
assert_eq!((p0, m0), ("anthropic", "claude-x")); - 499
let (p1, m1) = r.attempt_leg(&r.attempts[1]); - 500
assert_eq!((p1, m1), ("openai", "gpt-x")); - 501
} - 502
- 503
#[test] - 504
fn a_receipt_written_before_per_leg_attribution_still_loads() { - 505
let legacy = serde_json::json!({ - 506
"purpose": "execute", - 507
"model": "m", - 508
"attempts": [{ - 509
"ordinal": 0, - 510
"reason": "initial", - 511
"domain": "network", - 512
"settlement": "unknown", - 513
"latency_ms": 5 - 514
}] - 515
}); - 516
let back: WorkReceipt = serde_json::from_value(legacy).unwrap(); - 517
assert_eq!(back.provider, ""); - 518
assert_eq!(back.attempt_leg(&back.attempts[0]), ("", "m")); - 519
} - 520
} - 521
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.