- 1
//! The append-only commitment ledger and its projection. - 2
//! - 3
//! Events are never rewritten; the current [`Commitment`] is always a fold - 4
//! over them, exactly as `vak-session`'s work projector folds `WorkEvent`s. - 5
//! That is what makes "what did this agent commit to, under whose authority, - 6
//! and what closed it" answerable months later rather than a matter of trust. - 7
//! - 8
//! One rule is enforced *at append time* rather than at read time: a closure - 9
//! claiming success is refused unless the evidence actually supports it. A - 10
//! ledger that can record a lie is not an audit trail. - 11
- 12
use std::io::Write; - 13
use std::path::{Path, PathBuf}; - 14
- 15
use serde::{Deserialize, Serialize}; - 16
- 17
use vak_intent::{Envelope, Reading, Satisfaction}; - 18
use vak_session::types::{CriterionResult, EvidenceRef, WorkCriterion}; - 19
- 20
use crate::types::{ - 21
Advancement, Closure, ClosureRefusal, Commitment, CommitmentSpec, CriterionState, Economics, - 22
Episode, Phase, Suspension, Verdict, - 23
}; - 24
- 25
/// One thing that happened to a commitment. - 26
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] - 27
#[serde(tag = "kind", rename_all = "kebab-case")] - 28
pub enum EventKind { - 29
Opened { - 30
spec: Box<CommitmentSpec>, - 31
}, - 32
/// Resumed after a gap; the engagement was re-resolved against current - 33
/// reality and this is what had changed underneath it. - 34
Readmitted { - 35
#[serde(default)] - 36
drift: Vec<String>, - 37
reading: Box<Reading>, - 38
}, - 39
EpisodeStarted { - 40
episode_id: String, - 41
session_id: String, - 42
}, - 43
EpisodeEnded { - 44
episode_id: String, - 45
advancement: Advancement, - 46
#[serde(default)] - 47
spend_usd: f64, - 48
}, - 49
Suspended { - 50
suspension: Suspension, - 51
}, - 52
Resumed { - 53
reason: String, - 54
}, - 55
Blocked { - 56
blocker: String, - 57
}, - 58
Unblocked { - 59
reason: String, - 60
}, - 61
/// A criterion was evaluated. `strength` records *how* it was established, - 62
/// which is what the closure invariant later checks. - 63
CriterionEvaluated { - 64
criterion_id: String, - 65
result: CriterionResult, - 66
strength: Satisfaction, - 67
}, - 68
EnvelopeGranted { - 69
envelope: Box<Envelope>, - 70
}, - 71
EnvelopeRevoked { - 72
envelope_id: String, - 73
by: String, - 74
}, - 75
QuestionAnswered { - 76
question_id: String, - 77
answer: String, - 78
}, - 79
Superseded { - 80
by: String, - 81
reason: String, - 82
}, - 83
Closed { - 84
verdict: Verdict, - 85
strength: Satisfaction, - 86
#[serde(default)] - 87
evidence: Vec<EvidenceRef>, - 88
note: String, - 89
}, - 90
} - 91
- 92
impl EventKind { - 93
pub fn as_str(&self) -> &'static str { - 94
match self { - 95
EventKind::Opened { .. } => "opened", - 96
EventKind::Readmitted { .. } => "readmitted", - 97
EventKind::EpisodeStarted { .. } => "episode-started", - 98
EventKind::EpisodeEnded { .. } => "episode-ended", - 99
EventKind::Suspended { .. } => "suspended", - 100
EventKind::Resumed { .. } => "resumed", - 101
EventKind::Blocked { .. } => "blocked", - 102
EventKind::Unblocked { .. } => "unblocked", - 103
EventKind::CriterionEvaluated { .. } => "criterion-evaluated", - 104
EventKind::EnvelopeGranted { .. } => "envelope-granted", - 105
EventKind::EnvelopeRevoked { .. } => "envelope-revoked", - 106
EventKind::QuestionAnswered { .. } => "question-answered", - 107
EventKind::Superseded { .. } => "superseded", - 108
EventKind::Closed { .. } => "closed", - 109
} - 110
} - 111
} - 112
- 113
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] - 114
pub struct Event { - 115
pub event_id: String, - 116
pub commitment_id: String, - 117
pub ts: chrono::DateTime<chrono::Utc>, - 118
#[serde(flatten)] - 119
pub kind: EventKind, - 120
} - 121
- 122
impl Event { - 123
pub fn new(commitment_id: impl Into<String>, kind: EventKind) -> Self { - 124
Event { - 125
event_id: uuid::Uuid::now_v7().to_string(), - 126
commitment_id: commitment_id.into(), - 127
ts: chrono::Utc::now(), - 128
kind, - 129
} - 130
} - 131
} - 132
- 133
#[derive(Debug, thiserror::Error)] - 134
pub enum LedgerError { - 135
#[error("io error: {0}")] - 136
Io(#[from] std::io::Error), - 137
#[error("serialization error: {0}")] - 138
Serde(#[from] serde_json::Error), - 139
#[error("unknown commitment '{0}'")] - 140
UnknownCommitment(String), - 141
#[error("commitment '{0}' is already closed")] - 142
AlreadyClosed(String), - 143
#[error(transparent)] - 144
Closure(#[from] ClosureRefusal), - 145
} - 146
- 147
/// Held while one check-and-append runs. The OS releases the lock when the - 148
/// file closes — on drop, and when a process dies holding it — so there is - 149
/// no stale lock to detect and none to break. - 150
struct LedgerLock { - 151
_file: std::fs::File, - 152
} - 153
- 154
/// Append-only JSONL of commitment events, one file per home. - 155
/// - 156
/// Sits beside `routing-evidence.jsonl` in the sessions home for the same - 157
/// reason: it is runtime evidence about the agent's own behaviour, not - 158
/// conversation content, and it must survive any individual session. - 159
pub struct CommitmentLedger { - 160
path: PathBuf, - 161
} - 162
- 163
impl CommitmentLedger { - 164
pub fn new(sessions_home: &Path) -> Self { - 165
CommitmentLedger { - 166
path: sessions_home.join("commitments.jsonl"), - 167
} - 168
} - 169
- 170
pub fn path(&self) -> &Path { - 171
&self.path - 172
} - 173
- 174
/// Append one event, after checking it is legal against current state. - 175
/// - 176
/// The check happens here rather than in a caller because there are many - 177
/// callers and exactly one ledger: a rule enforced at the write boundary - 178
/// cannot be bypassed by a surface that forgot about it. The check and - 179
/// the append happen under an exclusive lock, so two writers — a running - 180
/// turn and `vak commit close`, or the upkeep tick — cannot both read - 181
/// "not closed" and both close it. - 182
/// - 183
/// A closure's strength is the runtime's to state, not the caller's: it - 184
/// is recomputed here from the criteria as they stand, so a `Closed` - 185
/// event can never record stronger evidence than the commitment holds. - 186
pub fn append(&self, event: &Event) -> Result<(), LedgerError> { - 187
let _guard = self.lock()?; - 188
let current = self - 189
.get(&event.commitment_id)? - 190
.ok_or_else(|| LedgerError::UnknownCommitment(event.commitment_id.clone()))?; - 191
if current.phase.is_terminal() { - 192
return Err(LedgerError::AlreadyClosed(event.commitment_id.clone())); - 193
} - 194
if let EventKind::Closed { verdict, .. } = &event.kind { - 195
// The closure invariant, as a hard error rather than a lint. - 196
let achieved = current.may_close(*verdict)?; - 197
let mut event = event.clone(); - 198
if let EventKind::Closed { strength, .. } = &mut event.kind { - 199
*strength = achieved; - 200
} - 201
return self.append_unchecked(&event); - 202
} - 203
self.append_unchecked(event) - 204
} - 205
- 206
/// A cross-process lock on the ledger, held for one check-and-append. - 207
/// The lock file itself is never removed: deleting it while another - 208
/// writer waits on it would hand the two of them different files. - 209
fn lock(&self) -> Result<LedgerLock, LedgerError> { - 210
let path = self.path.with_extension("lock"); - 211
if let Some(parent) = path.parent() { - 212
std::fs::create_dir_all(parent)?; - 213
} - 214
let file = std::fs::OpenOptions::new() - 215
.create(true) - 216
.truncate(false) - 217
.write(true) - 218
.open(&path)?; - 219
file.lock()?; - 220
Ok(LedgerLock { _file: file }) - 221
} - 222
- 223
/// Append without the closure check. Used by the projector's own tests and - 224
/// by recovery tooling that is deliberately reconstructing history. - 225
fn append_unchecked(&self, event: &Event) -> Result<(), LedgerError> { - 226
if let Some(parent) = self.path.parent() { - 227
std::fs::create_dir_all(parent)?; - 228
} - 229
let line = serde_json::to_string(event)?; - 230
let mut file = std::fs::OpenOptions::new() - 231
.create(true) - 232
.append(true) - 233
.open(&self.path)?; - 234
writeln!(file, "{line}")?; - 235
Ok(()) - 236
} - 237
- 238
/// Every event, oldest first. A corrupt line — torn, not UTF-8, or not - 239
/// an event — is skipped rather than trusted, and never ends the read: - 240
/// a malformed row must not become a state change, and must not hide - 241
/// every row written after it either. - 242
pub fn events(&self) -> Vec<Event> { - 243
let Ok(bytes) = std::fs::read(&self.path) else { - 244
return Vec::new(); - 245
}; - 246
bytes - 247
.split(|byte| *byte == b'\n') - 248
.filter_map(|line| serde_json::from_slice::<Event>(line).ok()) - 249
.collect() - 250
} - 251
- 252
/// Events for one commitment. - 253
pub fn events_for(&self, commitment_id: &str) -> Vec<Event> { - 254
self.events() - 255
.into_iter() - 256
.filter(|event| event.commitment_id == commitment_id) - 257
.collect() - 258
} - 259
- 260
/// Project one commitment's current state. - 261
pub fn get(&self, commitment_id: &str) -> Result<Option<Commitment>, LedgerError> { - 262
let events = self.events_for(commitment_id); - 263
Ok(project(&events)) - 264
} - 265
- 266
/// Every commitment, most recently updated first. - 267
pub fn all(&self) -> Vec<Commitment> { - 268
let mut by_id: std::collections::BTreeMap<String, Vec<Event>> = - 269
std::collections::BTreeMap::new(); - 270
for event in self.events() { - 271
by_id - 272
.entry(event.commitment_id.clone()) - 273
.or_default() - 274
.push(event); - 275
} - 276
let mut out: Vec<Commitment> = by_id - 277
.values() - 278
.filter_map(|events| project(events)) - 279
.collect(); - 280
out.sort_by_key(|commitment| std::cmp::Reverse(commitment.updated_at)); - 281
out - 282
} - 283
- 284
/// Commitments the portfolio scheduler could act on now. - 285
pub fn open(&self) -> Vec<Commitment> { - 286
self.all() - 287
.into_iter() - 288
.filter(|commitment| !commitment.phase.is_terminal()) - 289
.collect() - 290
} - 291
- 292
/// Open a new commitment and return its id. - 293
pub fn open_commitment(&self, spec: CommitmentSpec) -> Result<String, LedgerError> { - 294
let commitment_id = uuid::Uuid::now_v7().to_string(); - 295
self.append_unchecked(&Event::new( - 296
commitment_id.clone(), - 297
EventKind::Opened { - 298
spec: Box::new(spec), - 299
}, - 300
))?; - 301
Ok(commitment_id) - 302
} - 303
} - 304
- 305
/// Fold events into current state. - 306
/// - 307
/// Returns `None` when the first event is not an `Opened`, which is how a - 308
/// truncated or partially-corrupt ledger declines to invent a commitment - 309
/// rather than projecting a plausible-looking fiction. - 310
pub fn project(events: &[Event]) -> Option<Commitment> { - 311
let first = events.first()?; - 312
let EventKind::Opened { spec } = &first.kind else { - 313
return None; - 314
}; - 315
let spec = (**spec).clone(); - 316
let criteria = spec - 317
.criteria - 318
.iter() - 319
.map(|criterion| CriterionState { - 320
criterion_id: criterion.criterion_id.clone(), - 321
statement: criterion.statement.clone(), - 322
required: criterion.required, - 323
result: None, - 324
strength: None, - 325
evaluated_at: None, - 326
}) - 327
.collect(); - 328
- 329
let mut commitment = Commitment { - 330
commitment_id: first.commitment_id.clone(), - 331
opened_at: first.ts, - 332
spec, - 333
phase: Phase::Active, - 334
criteria, - 335
episodes: Vec::new(), - 336
suspension: None, - 337
blocker: None, - 338
envelope: None, - 339
closure: None, - 340
superseded_by: None, - 341
spend_usd: 0.0, - 342
consecutive_stalls: 0, - 343
drift: Vec::new(), - 344
updated_at: first.ts, - 345
}; - 346
- 347
for event in events.iter().skip(1) { - 348
commitment.updated_at = event.ts; - 349
match &event.kind { - 350
EventKind::Opened { .. } => { - 351
// A second open for the same id is a ledger fault, not a - 352
// reset. Ignore it rather than losing the accumulated history. - 353
} - 354
EventKind::Readmitted { drift, reading } => { - 355
commitment.drift = drift.clone(); - 356
commitment.spec.reading = (**reading).clone(); - 357
} - 358
EventKind::EpisodeStarted { - 359
episode_id, - 360
session_id, - 361
} => { - 362
// A new episode is someone working it again: whatever - 363
// blocked or suspended the last one is no longer the state. - 364
commitment.phase = Phase::Active; - 365
commitment.suspension = None; - 366
commitment.blocker = None; - 367
commitment.episodes.push(Episode { - 368
episode_id: episode_id.clone(), - 369
session_id: session_id.clone(), - 370
started_at: event.ts, - 371
ended_at: None, - 372
advancement: None, - 373
spend_usd: 0.0, - 374
}); - 375
} - 376
EventKind::EpisodeEnded { - 377
episode_id, - 378
advancement, - 379
spend_usd, - 380
} => { - 381
if let Some(episode) = commitment - 382
.episodes - 383
.iter_mut() - 384
.find(|episode| episode.episode_id == *episode_id) - 385
{ - 386
episode.ended_at = Some(event.ts); - 387
episode.advancement = Some(advancement.clone()); - 388
episode.spend_usd = *spend_usd; - 389
} - 390
commitment.spend_usd += spend_usd; - 391
if advancement.is_stall() { - 392
commitment.consecutive_stalls += 1; - 393
} else { - 394
// Any real advancement — including `Learned` — clears the - 395
// streak. Exploration is not stalling. - 396
commitment.consecutive_stalls = 0; - 397
} - 398
if let Advancement::Blocked { blocker } = advancement { - 399
commitment.phase = Phase::Blocked; - 400
commitment.blocker = Some(blocker.clone()); - 401
} - 402
} - 403
EventKind::Suspended { suspension } => { - 404
commitment.phase = Phase::Suspended; - 405
commitment.suspension = Some(suspension.clone()); - 406
} - 407
EventKind::Resumed { .. } => { - 408
commitment.phase = Phase::Active; - 409
commitment.suspension = None; - 410
} - 411
EventKind::Blocked { blocker } => { - 412
commitment.phase = Phase::Blocked; - 413
commitment.blocker = Some(blocker.clone()); - 414
} - 415
EventKind::Unblocked { .. } => { - 416
commitment.phase = Phase::Active; - 417
commitment.blocker = None; - 418
} - 419
EventKind::CriterionEvaluated { - 420
criterion_id, - 421
result, - 422
strength, - 423
} => { - 424
if let Some(criterion) = commitment - 425
.criteria - 426
.iter_mut() - 427
.find(|criterion| criterion.criterion_id == *criterion_id) - 428
{ - 429
criterion.result = Some(result.clone()); - 430
criterion.strength = Some(*strength); - 431
criterion.evaluated_at = Some(event.ts); - 432
} - 433
if commitment.phase == Phase::Active { - 434
commitment.phase = Phase::Satisfying; - 435
} - 436
} - 437
EventKind::EnvelopeGranted { envelope } => { - 438
commitment.envelope = Some((**envelope).clone()); - 439
} - 440
EventKind::EnvelopeRevoked { envelope_id, .. } => { - 441
if let Some(envelope) = commitment.envelope.as_mut() - 442
&& envelope.envelope_id == *envelope_id - 443
{ - 444
envelope.revoked_at = Some(event.ts); - 445
} - 446
} - 447
EventKind::QuestionAnswered { question_id, .. } => { - 448
let answered = matches!( - 449
&commitment.suspension, - 450
Some(Suspension::Human { question_id: id, .. }) if id == question_id - 451
); - 452
if answered { - 453
commitment.phase = Phase::Active; - 454
commitment.suspension = None; - 455
} - 456
} - 457
EventKind::Superseded { by, reason } => { - 458
commitment.superseded_by = Some(by.clone()); - 459
commitment.phase = Phase::Closed; - 460
commitment.closure = Some(Closure { - 461
verdict: Verdict::Superseded, - 462
strength: commitment.achieved_strength(), - 463
closed_at: event.ts, - 464
evidence: Vec::new(), - 465
note: reason.clone(), - 466
}); - 467
} - 468
EventKind::Closed { - 469
verdict, - 470
strength, - 471
evidence, - 472
note, - 473
} => { - 474
commitment.phase = Phase::Closed; - 475
commitment.closure = Some(Closure { - 476
verdict: *verdict, - 477
strength: *strength, - 478
closed_at: event.ts, - 479
evidence: evidence.clone(), - 480
note: note.clone(), - 481
}); - 482
} - 483
} - 484
} - 485
Some(commitment) - 486
} - 487
- 488
/// Build a spec from a resolved reading. - 489
/// - 490
/// `criteria` come from whatever proposed them — the planner, a flow, or a - 491
/// human. The model may propose criteria; it never marks them passed. - 492
pub fn spec_from_reading( - 493
objective: impl Into<String>, - 494
reading: Reading, - 495
criteria: Vec<WorkCriterion>, - 496
cwd: std::path::PathBuf, - 497
economics: Economics, - 498
) -> CommitmentSpec { - 499
let min_satisfaction = reading.evidence.min_satisfaction(); - 500
CommitmentSpec { - 501
objective: objective.into(), - 502
reading, - 503
criteria, - 504
min_satisfaction, - 505
economics, - 506
cwd, - 507
supersedes: None, - 508
thread_id: None, - 509
audience_id: None, - 510
} - 511
} - 512
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.