- 1
//! Wiring durable commitments into the run loop. - 2
//! - 3
//! `vak-commit` owns the lifecycle, the satisfaction lattice and the ledger; - 4
//! it knows nothing about sessions, providers or the agent loop. This module - 5
//! is the seam that opens a commitment when a turn's reading calls for one, - 6
//! records the episode that advanced it, and evaluates its criteria against - 7
//! the world. - 8
//! - 9
//! # Why episodes are recorded here rather than in the agent - 10
//! - 11
//! An episode's outcome is a judgement about *the commitment*, not about the - 12
//! turn: a turn that ended cleanly may still have moved nothing. Only the - 13
//! caller that holds both the commitment and the turn's result can say which - 14
//! it was, so the classification lives at this seam and `vak-agent` stays - 15
//! unaware that commitments exist at all. - 16
- 17
use std::path::Path; - 18
- 19
use vak_commit::CriterionEvaluator; - 20
use vak_commit::{ - 21
Advancement, CommitmentLedger, CommitmentSpec, Economics, Evaluation, Event, EventKind, Verdict, - 22
}; - 23
use vak_intent::{Intent, Satisfaction}; - 24
use vak_session::types::{CriterionKind, CriterionResult, WorkCriterion}; - 25
- 26
/// A commitment this turn is serving, and the episode it opened on it. - 27
#[derive(Debug, Clone)] - 28
pub struct EpisodeHandle { - 29
pub commitment_id: String, - 30
pub episode_id: String, - 31
/// The strand this episode serves. - 32
pub strand_id: String, - 33
} - 34
- 35
/// Economics from configuration, resolved once. - 36
pub fn economics(config: &vak_config::Config) -> Economics { - 37
Economics { - 38
lifetime_budget_usd: config.commitment.lifetime_budget_usd, - 39
expires_at: config - 40
.commitment - 41
.default_ttl_days - 42
.map(|days| chrono::Utc::now() + chrono::Duration::days(i64::from(days))), - 43
review_every_hours: config.commitment.review_every_hours, - 44
stall_limit: config.commitment.stall_limit, - 45
} - 46
} - 47
- 48
/// Criteria proposed from the reading alone. - 49
/// - 50
/// Deliberately thin. The runtime can only propose what it can also *check*, - 51
/// and from a bare request it can check almost nothing — so this yields at - 52
/// most a semantic placeholder and leaves the real criteria to the planner, a - 53
/// flow, or a human. Inventing a `Shell` criterion by guessing at a test - 54
/// command would manufacture `Observed` evidence out of a guess, which is - 55
/// precisely the laundering the satisfaction lattice exists to prevent. - 56
pub fn seed_criteria(intent: &Intent, objective: &str) -> Vec<WorkCriterion> { - 57
if intent.reading.evidence.min_satisfaction().rank() < Satisfaction::Cited.rank() { - 58
return Vec::new(); - 59
} - 60
vec![WorkCriterion { - 61
criterion_id: "objective".into(), - 62
statement: objective.to_string(), - 63
// Semantic, so it carries `Asserted` strength and therefore cannot by - 64
// itself close work held to `Verified` or `Audited`. That is the - 65
// intended outcome: the commitment stays open, visibly, until someone - 66
// supplies a criterion the runtime can actually check. - 67
kind: CriterionKind::Semantic, - 68
required: true, - 69
}] - 70
} - 71
- 72
/// One line a human would recognise months later. - 73
pub fn objective_from(prompt: &str) -> String { - 74
let first = prompt.trim().lines().next().unwrap_or("").trim(); - 75
let mut out: String = first.chars().take(120).collect(); - 76
if first.chars().count() > 120 { - 77
out.push('…'); - 78
} - 79
if out.is_empty() { - 80
"untitled work".into() - 81
} else { - 82
out - 83
} - 84
} - 85
- 86
/// What a turn will do with the commitment ledger, decided before it writes - 87
/// anything: which strands work on a commitment that already exists, and - 88
/// which open a new one. - 89
#[derive(Debug, Clone, Default)] - 90
pub struct EpisodePlan { - 91
pub strands: Vec<PlannedStrand>, - 92
} - 93
- 94
#[derive(Debug, Clone)] - 95
pub struct PlannedStrand { - 96
pub strand_id: String, - 97
pub target: PlannedTarget, - 98
} - 99
- 100
#[derive(Debug, Clone)] - 101
pub enum PlannedTarget { - 102
/// Continue the open commitment on this strand's thread. - 103
Existing { - 104
commitment_id: String, - 105
/// Its envelope, when one is live: what the human pre-authorized - 106
/// for this work. - 107
envelope: Option<vak_intent::Envelope>, - 108
}, - 109
/// Open a new commitment, superseding the one the strand replaces. - 110
Open { supersedes: Option<String> }, - 111
} - 112
- 113
impl EpisodePlan { - 114
/// Live envelopes by strand id, for `vak_intent::apply_envelopes`. Only - 115
/// an existing commitment can carry one: a grant is given to work that - 116
/// already exists, never to work this turn is about to open. - 117
pub fn envelopes(&self) -> std::collections::BTreeMap<String, vak_intent::Envelope> { - 118
self.strands - 119
.iter() - 120
.filter_map(|planned| match &planned.target { - 121
PlannedTarget::Existing { - 122
envelope: Some(envelope), - 123
.. - 124
} => Some((planned.strand_id.clone(), envelope.clone())), - 125
_ => None, - 126
}) - 127
.collect() - 128
} - 129
- 130
/// The commitments whose live envelopes a gated action may fall inside. - 131
pub fn enveloped_commitments(&self) -> Vec<String> { - 132
self.strands - 133
.iter() - 134
.filter_map(|planned| match &planned.target { - 135
PlannedTarget::Existing { - 136
commitment_id, - 137
envelope: Some(_), - 138
} => Some(commitment_id.clone()), - 139
_ => None, - 140
}) - 141
.collect() - 142
} - 143
} - 144
- 145
/// Decide, without writing, what each strand of this turn does with the - 146
/// commitment ledger. - 147
/// - 148
/// One commitment per *thread*, not per turn: - 149
/// - a strand on a thread that already has an open commitment works on it, - 150
/// whatever its own horizon reads as — "now also check staging" is part of - 151
/// the weekly job it continues; - 152
/// - a strand that replaces a thread with an open commitment opens the - 153
/// successor and supersedes it — `/goal replace` on durable work yields - 154
/// durable work; - 155
/// - any other strand opens a commitment only when it is itself durable and - 156
/// its reading is confident enough to carry an obligation. - 157
/// - 158
/// Empty when commitments are disabled or nothing here is durable. - 159
pub fn plan_episodes( - 160
sessions_home: &Path, - 161
config: &vak_config::Config, - 162
intent: &Intent, - 163
now: chrono::DateTime<chrono::Utc>, - 164
) -> EpisodePlan { - 165
if !config.commitment.enabled || intent.strands.is_empty() { - 166
return EpisodePlan::default(); - 167
} - 168
let open = CommitmentLedger::new(sessions_home).open(); - 169
let on_thread = |thread: &str| { - 170
open.iter() - 171
.find(|commitment| commitment.spec.thread_id.as_deref() == Some(thread)) - 172
}; - 173
let mut plan = EpisodePlan::default(); - 174
for strand in &intent.strands { - 175
let durable = strand.engagement.posture.open_commitment - 176
&& strand - 177
.reading - 178
.may_open_commitment(config.intent.accept_confidence); - 179
let target = if let Some(existing) = on_thread(&strand.thread_id) { - 180
PlannedTarget::Existing { - 181
commitment_id: existing.commitment_id.clone(), - 182
envelope: existing - 183
.envelope - 184
.clone() - 185
.filter(|envelope| envelope.is_live(now)), - 186
} - 187
} else if let Some(replaced) = strand.lineage.replaced_thread().and_then(on_thread) { - 188
PlannedTarget::Open { - 189
supersedes: Some(replaced.commitment_id.clone()), - 190
} - 191
} else if durable { - 192
PlannedTarget::Open { supersedes: None } - 193
} else { - 194
continue; - 195
}; - 196
plan.strands.push(PlannedStrand { - 197
strand_id: strand.strand_id.clone(), - 198
target, - 199
}); - 200
} - 201
plan - 202
} - 203
- 204
/// Carry out an [`EpisodePlan`]: open what it opens, supersede what it - 205
/// replaces, and start one episode per planned strand. - 206
/// - 207
/// A ledger failure costs the audit row, never the turn: the strand is - 208
/// skipped and logged. Episodes come back in strand order; the first is the - 209
/// turn's primary. - 210
#[allow(clippy::too_many_arguments)] - 211
pub fn begin_episodes( - 212
sessions_home: &Path, - 213
config: &vak_config::Config, - 214
intent: &Intent, - 215
plan: &EpisodePlan, - 216
prompt: &str, - 217
session_id: &str, - 218
cwd: &Path, - 219
audience_id: Option<&str>, - 220
) -> Vec<EpisodeHandle> { - 221
let ledger = CommitmentLedger::new(sessions_home); - 222
let mut handles = Vec::new(); - 223
for planned in &plan.strands { - 224
let Some(strand) = intent - 225
.strands - 226
.iter() - 227
.find(|strand| strand.strand_id == planned.strand_id) - 228
else { - 229
continue; - 230
}; - 231
let commitment_id = match &planned.target { - 232
PlannedTarget::Existing { commitment_id, .. } => commitment_id.clone(), - 233
PlannedTarget::Open { supersedes } => { - 234
let Some(new_id) = open_for_strand( - 235
&ledger, - 236
config, - 237
strand, - 238
prompt, - 239
cwd, - 240
supersedes.clone(), - 241
audience_id, - 242
) else { - 243
continue; - 244
}; - 245
if let Some(old) = supersedes { - 246
let _ = ledger.append(&Event::new( - 247
old, - 248
EventKind::Superseded { - 249
by: new_id.clone(), - 250
reason: "replaced by an explicit /goal replace".into(), - 251
}, - 252
)); - 253
} - 254
new_id - 255
} - 256
}; - 257
close_orphan_episode(&ledger, &commitment_id, session_id); - 258
let episode_id = uuid::Uuid::now_v7().to_string(); - 259
if let Err(error) = ledger.append(&Event::new( - 260
&commitment_id, - 261
EventKind::EpisodeStarted { - 262
episode_id: episode_id.clone(), - 263
session_id: session_id.to_string(), - 264
}, - 265
)) { - 266
eprintln!("[commit] could not start an episode: {error}"); - 267
continue; - 268
} - 269
handles.push(EpisodeHandle { - 270
commitment_id, - 271
episode_id, - 272
strand_id: strand.strand_id.clone(), - 273
}); - 274
} - 275
handles - 276
} - 277
- 278
/// Open a commitment for one strand. - 279
fn open_for_strand( - 280
ledger: &CommitmentLedger, - 281
config: &vak_config::Config, - 282
strand: &vak_intent::Strand, - 283
prompt: &str, - 284
cwd: &Path, - 285
supersedes: Option<String>, - 286
audience_id: Option<&str>, - 287
) -> Option<String> { - 288
let objective = objective_from(if strand.text.is_empty() { - 289
prompt - 290
} else { - 291
&strand.text - 292
}); - 293
let seed = Intent { - 294
reading: strand.reading.clone(), - 295
strands: Vec::new(), - 296
engagement: strand.engagement.clone(), - 297
provenance: vak_intent::Provenance::new( - 298
vak_intent::Tier::Signals, - 299
vak_intent::RESOLVER_VERSION, - 300
Vec::new(), - 301
), - 302
}; - 303
let spec = CommitmentSpec { - 304
criteria: seed_criteria(&seed, &objective), - 305
min_satisfaction: strand.reading.evidence.min_satisfaction(), - 306
economics: economics(config), - 307
cwd: cwd.to_path_buf(), - 308
supersedes, - 309
thread_id: Some(strand.thread_id.clone()), - 310
audience_id: audience_id.map(str::to_string), - 311
objective, - 312
reading: strand.reading.clone(), - 313
}; - 314
match ledger.open_commitment(spec) { - 315
Ok(id) => Some(id), - 316
Err(error) => { - 317
eprintln!("[commit] could not open a commitment: {error}"); - 318
None - 319
} - 320
} - 321
} - 322
- 323
/// A resumed session may have crashed after EpisodeStarted but before - 324
/// EpisodeEnded. Close that exact orphan as blocked before opening the new - 325
/// episode; never replay its effects implicitly. - 326
fn close_orphan_episode(ledger: &CommitmentLedger, commitment_id: &str, session_id: &str) { - 327
if let Ok(Some(commitment)) = ledger.get(commitment_id) - 328
&& let Some(episode) = commitment - 329
.episodes - 330
.iter() - 331
.rev() - 332
.find(|episode| episode.ended_at.is_none() && episode.session_id == session_id) - 333
{ - 334
let _ = ledger.append(&Event::new( - 335
commitment_id, - 336
EventKind::EpisodeEnded { - 337
episode_id: episode.episode_id.clone(), - 338
advancement: Advancement::Blocked { - 339
blocker: "recovered after an interrupted process; review before retry".into(), - 340
}, - 341
spend_usd: 0.0, - 342
}, - 343
)); - 344
} - 345
} - 346
- 347
/// The commitment ledger rendered for the model: `ContextProfile::Full`. - 348
/// - 349
/// One block per commitment this turn serves — objective, criteria and - 350
/// their standing, an open question if the work is suspended on one. Terse - 351
/// and factual, like the intent note it is appended to; the model reads the - 352
/// state of its obligations instead of reconstructing them from history. - 353
pub fn prompt_projection(sessions_home: &Path, episodes: &[EpisodeHandle]) -> Option<String> { - 354
if episodes.is_empty() { - 355
return None; - 356
} - 357
let ledger = CommitmentLedger::new(sessions_home); - 358
let mut lines = Vec::new(); - 359
for episode in episodes { - 360
let Ok(Some(commitment)) = ledger.get(&episode.commitment_id) else { - 361
continue; - 362
}; - 363
let met = commitment - 364
.criteria - 365
.iter() - 366
.filter(|c| { - 367
matches!( - 368
c.result, - 369
Some(vak_session::types::CriterionResult::Passed { .. }) - 370
) - 371
}) - 372
.count(); - 373
lines.push(format!( - 374
"Commitment {} ({}): {} — {} of {} criteria met, {} episode(s) so far, closes at {} evidence.", - 375
commitment.commitment_id, - 376
commitment.phase.as_str(), - 377
commitment.spec.objective, - 378
met, - 379
commitment.criteria.len(), - 380
commitment.episodes.len(), - 381
commitment.spec.min_satisfaction.as_str() - 382
)); - 383
for criterion in &commitment.criteria { - 384
let result = match &criterion.result { - 385
Some(vak_session::types::CriterionResult::Passed { .. }) => "passed", - 386
Some(vak_session::types::CriterionResult::Failed { .. }) => "failed", - 387
Some(vak_session::types::CriterionResult::Unknown { .. }) => "unknown", - 388
None => "not yet evaluated", - 389
}; - 390
let standing = match criterion.strength { - 391
Some(strength) if criterion.result.is_some() => { - 392
format!("{result} ({})", strength.as_str()) - 393
} - 394
_ => result.to_string(), - 395
}; - 396
lines.push(format!(" - {}: {standing}", criterion.statement)); - 397
} - 398
if let Some(vak_commit::Suspension::Human { question, .. }) = &commitment.suspension { - 399
lines.push(format!(" open question: {question}")); - 400
} - 401
if let Some(blocker) = &commitment.blocker { - 402
lines.push(format!(" blocked: {blocker}")); - 403
} - 404
} - 405
(!lines.is_empty()).then(|| lines.join("\n")) - 406
} - 407
- 408
/// What a finished turn did for its commitment. - 409
/// - 410
/// The distinction that matters is `Learned` versus `Stalled`. A turn that - 411
/// produced a substantive answer but moved no criterion has still reduced - 412
/// uncertainty and must not count against the stall breaker; a turn that - 413
/// produced nothing has. - 414
pub fn classify( - 415
outcome: &vak_agent::TurnOutcome, - 416
tool_calls: usize, - 417
criteria_moved: Vec<String>, - 418
) -> Advancement { - 419
if !criteria_moved.is_empty() { - 420
return Advancement::Advanced { - 421
criteria_moved, - 422
evidence: Vec::new(), - 423
}; - 424
} - 425
match outcome { - 426
vak_agent::TurnOutcome::Completed { response } => { - 427
let said_something = response - 428
.content - 429
.iter() - 430
.any(|block| matches!(block, vak_llm::ContentBlock::Text { text } if text.trim().len() > 40)); - 431
// Tool activity alone is not progress: retries, repeated reads, - 432
// and failed probes are common in a long run. Only a substantive - 433
// response or an explicitly moved criterion clears the stall - 434
// streak. - 435
if said_something { - 436
Advancement::Learned { - 437
fact: format!( - 438
"episode completed with {tool_calls} tool call(s) and no criterion movement" - 439
), - 440
} - 441
} else { - 442
Advancement::Stalled { - 443
reason: "episode produced neither output nor effect".into(), - 444
} - 445
} - 446
} - 447
vak_agent::TurnOutcome::Aborted { .. } => Advancement::Blocked { - 448
blocker: "cancelled".into(), - 449
}, - 450
vak_agent::TurnOutcome::Failed { error } => Advancement::Blocked { - 451
blocker: error.to_string(), - 452
}, - 453
// Hitting the turn ceiling is the textbook motion-without-progress - 454
// case, and the one the stall breaker exists to catch. - 455
vak_agent::TurnOutcome::MaxTurnsReached => Advancement::Stalled { - 456
reason: "exhausted the turn budget without moving a criterion".into(), - 457
}, - 458
} - 459
} - 460
- 461
/// Close out this turn's episode. - 462
pub fn end_episode( - 463
sessions_home: &Path, - 464
handle: &EpisodeHandle, - 465
advancement: Advancement, - 466
spend_usd: f64, - 467
) { - 468
let ledger = CommitmentLedger::new(sessions_home); - 469
if let Err(error) = ledger.append(&Event::new( - 470
&handle.commitment_id, - 471
EventKind::EpisodeEnded { - 472
episode_id: handle.episode_id.clone(), - 473
advancement, - 474
spend_usd, - 475
}, - 476
)) { - 477
eprintln!("[commit] could not close the episode: {error}"); - 478
} - 479
} - 480
- 481
/// Suspend a commitment on a question nobody here can answer. - 482
/// - 483
/// This is the `Defer` half of the human-in-the-loop contract: an unattended - 484
/// surface used to fail closed unconditionally, which is right for a one-shot - 485
/// turn and destroys month-long work that merely needed to wait. - 486
pub fn defer_for_human( - 487
sessions_home: &Path, - 488
commitment_id: &str, - 489
question: &str, - 490
addressed_to: Option<String>, - 491
escalation: vak_intent::Escalation, - 492
) -> Result<String, vak_commit::LedgerError> { - 493
let question_id = uuid::Uuid::now_v7().to_string(); - 494
CommitmentLedger::new(sessions_home).append(&Event::new( - 495
commitment_id, - 496
EventKind::Suspended { - 497
suspension: vak_commit::Suspension::Human { - 498
question_id: question_id.clone(), - 499
question: question.to_string(), - 500
addressed_to, - 501
escalation, - 502
}, - 503
}, - 504
))?; - 505
Ok(question_id) - 506
} - 507
- 508
/// Record a criterion evaluation the runtime performed. - 509
pub fn record_evaluation( - 510
sessions_home: &Path, - 511
commitment_id: &str, - 512
evaluation: &Evaluation, - 513
) -> Result<(), vak_commit::LedgerError> { - 514
CommitmentLedger::new(sessions_home).append(&Event::new( - 515
commitment_id, - 516
EventKind::CriterionEvaluated { - 517
criterion_id: evaluation.criterion_id.clone(), - 518
result: evaluation.result.clone(), - 519
strength: evaluation.strength, - 520
}, - 521
)) - 522
} - 523
- 524
/// Sweep commitments whose economics have run out. - 525
/// - 526
/// Expiry is an explicit `Expired` verdict rather than a deletion, so the - 527
/// record says the work lapsed instead of quietly forgetting it existed. - 528
/// Budget exhaustion deliberately does **not** close anything: a human can - 529
/// raise the ceiling, and everything the commitment established is still - 530
/// good, so the scheduler simply stops offering it. - 531
pub fn sweep_expired(sessions_home: &Path) -> Vec<String> { - 532
let ledger = CommitmentLedger::new(sessions_home); - 533
let now = chrono::Utc::now(); - 534
let mut closed = Vec::new(); - 535
for commitment in ledger.open() { - 536
if !commitment.is_expired(now) { - 537
continue; - 538
} - 539
let strength = commitment.achieved_strength(); - 540
if ledger - 541
.append(&Event::new( - 542
&commitment.commitment_id, - 543
EventKind::Closed { - 544
verdict: Verdict::Expired, - 545
strength, - 546
evidence: Vec::new(), - 547
note: "past its relevance window without closing".into(), - 548
}, - 549
)) - 550
.is_ok() - 551
{ - 552
closed.push(commitment.commitment_id); - 553
} - 554
} - 555
closed - 556
} - 557
- 558
/// One maintenance pass over the portfolio. - 559
/// - 560
/// Runs on the server's ordinary tick, independently of whether the heartbeat - 561
/// is enabled: heartbeat is an opt-in *model* pass and costs tokens, while - 562
/// everything here is filesystem and clock work that costs none. Tying durable - 563
/// work's upkeep to an opt-in prober would mean a commitment stopped being - 564
/// durable the moment someone turned the prober off. - 565
/// - 566
/// Everything it does is either a state transition the ledger already - 567
/// authorises or an evaluation the runtime performs itself. It never - 568
/// dispatches a model and never closes work as fulfilled. - 569
#[derive(Debug, Default, PartialEq)] - 570
pub struct Maintenance { - 571
/// Closed `Expired` because their relevance window passed. - 572
pub expired: Vec<String>, - 573
/// Woken because a scheduled time arrived. - 574
pub resumed: Vec<String>, - 575
/// Woken because a predicate the runtime can check became true. - 576
pub satisfied: Vec<String>, - 577
/// Deferred questions whose escalation policy came due. - 578
pub escalated: Vec<String>, - 579
} - 580
- 581
impl Maintenance { - 582
pub fn is_empty(&self) -> bool { - 583
self.expired.is_empty() - 584
&& self.resumed.is_empty() - 585
&& self.satisfied.is_empty() - 586
&& self.escalated.is_empty() - 587
} - 588
- 589
fn extend(&mut self, other: Maintenance) { - 590
self.expired.extend(other.expired); - 591
self.resumed.extend(other.resumed); - 592
self.satisfied.extend(other.satisfied); - 593
self.escalated.extend(other.escalated); - 594
} - 595
} - 596
- 597
/// One upkeep pass over every Agent's portfolio. - 598
/// - 599
/// Commitments live in the Agent that holds them - 600
/// (`<data home>/agents/<agent_id>/`, docs/design/64-agent-owned-platform.md), - 601
/// so a pass over one home would leave every other Agent's scheduled wakes - 602
/// and deadlines unkept. The data home's own ledger is included for work - 603
/// admitted without an Agent. - 604
pub async fn maintain_all(shared_data_home: &Path) -> Maintenance { - 605
let mut homes = vec![shared_data_home.to_path_buf()]; - 606
if let Ok(entries) = std::fs::read_dir(shared_data_home.join("agents")) { - 607
let mut agents: Vec<_> = entries - 608
.flatten() - 609
.map(|entry| entry.path()) - 610
.filter(|path| path.join("commitments.jsonl").is_file()) - 611
.collect(); - 612
agents.sort(); - 613
homes.extend(agents); - 614
} - 615
let mut report = Maintenance::default(); - 616
for home in homes { - 617
report.extend(maintain(&home).await); - 618
} - 619
report - 620
} - 621
- 622
/// Advance whatever the clock and the workspace now permit, for the - 623
/// portfolio in one home. A predicate is checked against the workspace its - 624
/// commitment belongs to, never the caller's. - 625
pub async fn maintain(sessions_home: &Path) -> Maintenance { - 626
let ledger = CommitmentLedger::new(sessions_home); - 627
let now = chrono::Utc::now(); - 628
let mut report = Maintenance { - 629
expired: sweep_expired(sessions_home), - 630
..Maintenance::default() - 631
}; - 632
- 633
for commitment in ledger.open() { - 634
let Some(suspension) = commitment.suspension.clone() else { - 635
continue; - 636
}; - 637
match suspension { - 638
vak_commit::Suspension::Schedule { at: Some(at), .. } if now >= at => { - 639
if ledger - 640
.append(&Event::new( - 641
&commitment.commitment_id, - 642
EventKind::Resumed { - 643
reason: "scheduled time reached".into(), - 644
}, - 645
)) - 646
.is_ok() - 647
{ - 648
report.resumed.push(commitment.commitment_id.clone()); - 649
} - 650
} - 651
// A predicate is checked by the runtime for free. This is the - 652
// path that lets "watch X and tell me when Y" cost nothing at all - 653
// while Y stays false. - 654
vak_commit::Suspension::Predicate { criterion } => { - 655
let evaluation = WorkspaceEvaluator { - 656
cwd: &commitment.spec.cwd, - 657
} - 658
.evaluate(&criterion) - 659
.await; - 660
if evaluation.passed() - 661
&& record_evaluation(sessions_home, &commitment.commitment_id, &evaluation) - 662
.is_ok() - 663
&& ledger - 664
.append(&Event::new( - 665
&commitment.commitment_id, - 666
EventKind::Resumed { - 667
reason: format!("condition met: {}", criterion.statement), - 668
}, - 669
)) - 670
.is_ok() - 671
{ - 672
report.satisfied.push(commitment.commitment_id.clone()); - 673
} - 674
} - 675
vak_commit::Suspension::Commitment { commitment_id } => { - 676
if ledger - 677
.get(&commitment_id) - 678
.ok() - 679
.flatten() - 680
.is_some_and(|dependency| { - 681
dependency - 682
.closure - 683
.as_ref() - 684
.is_some_and(|closure| closure.verdict == Verdict::Fulfilled) - 685
}) - 686
&& ledger - 687
.append(&Event::new( - 688
&commitment.commitment_id, - 689
EventKind::Resumed { - 690
reason: format!("dependency {commitment_id} was fulfilled"), - 691
}, - 692
)) - 693
.is_ok() - 694
{ - 695
report.resumed.push(commitment.commitment_id.clone()); - 696
} - 697
} - 698
vak_commit::Suspension::Human { - 699
question_id, - 700
escalation, - 701
.. - 702
} => { - 703
if let Some(id) = - 704
escalate_if_due(&ledger, &commitment, &question_id, &escalation, now) - 705
{ - 706
report.escalated.push(id); - 707
} - 708
} - 709
_ => {} - 710
} - 711
} - 712
report - 713
} - 714
- 715
/// Apply a deferred question's escalation policy once its deadline passes. - 716
/// - 717
/// A question with no policy waits forever by design; that is a decision, not - 718
/// a leak. What must never happen is a silent default standing in for consent - 719
/// on work that cannot be undone, so `AssumeConservative` is refused there - 720
/// even if a grant somehow carried it. - 721
fn escalate_if_due( - 722
ledger: &CommitmentLedger, - 723
commitment: &vak_commit::Commitment, - 724
question_id: &str, - 725
escalation: &vak_intent::Escalation, - 726
now: chrono::DateTime<chrono::Utc>, - 727
) -> Option<String> { - 728
let waited_hours = (now - commitment.updated_at).num_minutes() as f64 / 60.0; - 729
let (after, verdict, note): (u32, Option<Verdict>, &str) = match escalation { - 730
vak_intent::Escalation::WaitIndefinitely => return None, - 731
vak_intent::Escalation::AssumeConservative { after_hours } => { - 732
if !escalation.permitted_for(commitment.spec.reading.stakes) { - 733
return None; - 734
} - 735
( - 736
*after_hours, - 737
None, - 738
"no answer; taking the conservative branch", - 739
) - 740
} - 741
vak_intent::Escalation::AbandonAfter { after_hours } => ( - 742
*after_hours, - 743
Some(Verdict::Abandoned), - 744
"no answer within the agreed window", - 745
), - 746
vak_intent::Escalation::Reassign { after_hours, .. } => { - 747
(*after_hours, None, "reassigned after no answer") - 748
} - 749
}; - 750
if waited_hours < f64::from(after) { - 751
return None; - 752
} - 753
let event = match verdict { - 754
Some(verdict) => Event::new( - 755
&commitment.commitment_id, - 756
EventKind::Closed { - 757
verdict, - 758
strength: commitment.achieved_strength(), - 759
evidence: Vec::new(), - 760
note: format!("{note} (question {question_id})"), - 761
}, - 762
), - 763
None => Event::new( - 764
&commitment.commitment_id, - 765
EventKind::Resumed { - 766
reason: format!("{note} (question {question_id})"), - 767
}, - 768
), - 769
}; - 770
ledger - 771
.append(&event) - 772
.ok() - 773
.map(|()| commitment.commitment_id.clone()) - 774
} - 775
- 776
/// A criterion evaluator backed by the workspace filesystem. - 777
/// - 778
/// Only the checks that need no permission gate live here: file existence and - 779
/// content. `Shell`, `ToolSucceeded` and `FlowCompleted` are permissioned - 780
/// effects and must cross the tool broker, so they return `Unknown` until a - 781
/// caller supplies a brokered evaluator — an honest "not determined" rather - 782
/// than a failure the work did not earn. - 783
/// - 784
/// It runs in the server process, outside any sandbox, so a criterion may - 785
/// only name a path inside its commitment's workspace (invariant 10): an - 786
/// absolute path, a `..` step, or a symlink out of the tree is not checked - 787
/// at all, rather than turning a pass/fail into a probe of the machine. - 788
pub struct WorkspaceEvaluator<'a> { - 789
pub cwd: &'a Path, - 790
} - 791
- 792
/// `path` inside `cwd`, or `None` when it names anything outside it. - 793
fn inside_workspace(cwd: &Path, relative: &Path) -> Option<std::path::PathBuf> { - 794
if relative.is_absolute() - 795
|| relative.components().any(|component| { - 796
!matches!( - 797
component, - 798
std::path::Component::Normal(_) | std::path::Component::CurDir - 799
) - 800
}) - 801
{ - 802
return None; - 803
} - 804
let root = cwd.canonicalize().ok()?; - 805
let joined = root.join(relative); - 806
let mut probe = joined.as_path(); - 807
loop { - 808
if let Ok(real) = probe.canonicalize() { - 809
return real.starts_with(&root).then_some(joined); - 810
} - 811
probe = probe.parent()?; - 812
} - 813
} - 814
- 815
impl vak_commit::CriterionEvaluator for WorkspaceEvaluator<'_> { - 816
async fn evaluate(&self, criterion: &WorkCriterion) -> Evaluation { - 817
let outside = || { - 818
Evaluation::unknown( - 819
&criterion.criterion_id, - 820
"names a path outside the workspace; not checked", - 821
) - 822
}; - 823
match &criterion.kind { - 824
CriterionKind::FileExists { path } => { - 825
let Some(resolved) = inside_workspace(self.cwd, path) else { - 826
return outside(); - 827
}; - 828
if resolved.exists() { - 829
Evaluation::observed( - 830
criterion, - 831
CriterionResult::Passed { - 832
evidence: format!("{} exists", resolved.display()), - 833
}, - 834
) - 835
} else { - 836
Evaluation::observed( - 837
criterion, - 838
CriterionResult::Failed { - 839
reason: format!("{} does not exist", resolved.display()), - 840
}, - 841
) - 842
} - 843
} - 844
CriterionKind::FileContains { path, pattern } => { - 845
let Some(resolved) = inside_workspace(self.cwd, path) else { - 846
return outside(); - 847
}; - 848
match std::fs::read_to_string(&resolved) { - 849
Ok(text) if text.contains(pattern) => Evaluation::observed( - 850
criterion, - 851
CriterionResult::Passed { - 852
evidence: format!("{} contains the pattern", resolved.display()), - 853
}, - 854
), - 855
Ok(_) => Evaluation::observed( - 856
criterion, - 857
CriterionResult::Failed { - 858
reason: format!("{} does not contain the pattern", resolved.display()), - 859
}, - 860
), - 861
// Unreadable is not the same as failing: the check could - 862
// not run, and saying otherwise would be as dishonest as - 863
// claiming it passed. - 864
Err(error) => Evaluation::unknown( - 865
&criterion.criterion_id, - 866
format!("could not read {}: {error}", resolved.display()), - 867
), - 868
} - 869
} - 870
_ => Evaluation::unknown( - 871
&criterion.criterion_id, - 872
"needs a brokered evaluator; not checked here", - 873
), - 874
} - 875
} - 876
} - 877
- 878
#[cfg(test)] - 879
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 880
mod tests { - 881
use super::*; - 882
use vak_commit::CriterionEvaluator; - 883
use vak_intent::{Evidence, Horizon, Reading}; - 884
- 885
fn intent_with(evidence: Evidence, horizon: Horizon, confidence: f64) -> Intent { - 886
intent_on( - 887
evidence, - 888
horizon, - 889
confidence, - 890
"t1.0", - 891
vak_intent::Lineage::New, - 892
) - 893
} - 894
- 895
/// A one-strand intent. `strand_id` is also the thread for a new strand - 896
/// and a replacement; a continuation inherits its thread. - 897
fn intent_on( - 898
evidence: Evidence, - 899
horizon: Horizon, - 900
confidence: f64, - 901
strand_id: &str, - 902
lineage: vak_intent::Lineage, - 903
) -> Intent { - 904
let reading = Reading { - 905
horizon, - 906
evidence, - 907
confidence, - 908
axis_confidence: vak_intent::Confidences { - 909
act: confidence, - 910
horizon: confidence, - 911
stakes: confidence, - 912
evidence: confidence, - 913
}, - 914
..Reading::general() - 915
}; - 916
let engagement = vak_intent::derive(&reading, &vak_intent::Authority::default(), false); - 917
let thread_id = lineage.continued_thread().unwrap_or(strand_id).to_string(); - 918
let mut intent = Intent::general(vak_intent::RESOLVER_VERSION); - 919
intent.strands = vec![vak_intent::Strand { - 920
strand_id: strand_id.into(), - 921
thread_id, - 922
text: String::new(), - 923
reading: reading.clone(), - 924
relation: vak_intent::StrandRelation::Independent, - 925
lineage, - 926
engagement: engagement.clone(), - 927
}]; - 928
intent.reading = reading; - 929
intent.engagement = engagement; - 930
intent - 931
} - 932
- 933
fn begin_episode( - 934
sessions_home: &Path, - 935
config: &vak_config::Config, - 936
intent: &Intent, - 937
prompt: &str, - 938
session_id: &str, - 939
cwd: &Path, - 940
) -> Option<EpisodeHandle> { - 941
let plan = plan_episodes(sessions_home, config, intent, chrono::Utc::now()); - 942
begin_episodes( - 943
sessions_home, - 944
config, - 945
intent, - 946
&plan, - 947
prompt, - 948
session_id, - 949
cwd, - 950
Some("local"), - 951
) - 952
.into_iter() - 953
.next() - 954
} - 955
- 956
fn config() -> vak_config::Config { - 957
vak_config::Config::default() - 958
} - 959
- 960
#[test] - 961
fn a_turn_horizon_opens_no_commitment() { - 962
let dir = tempfile::tempdir().unwrap(); - 963
let handle = begin_episode( - 964
dir.path(), - 965
&config(), - 966
&intent_with(Evidence::None, Horizon::Turn, 0.9), - 967
"fix the test", - 968
"s1", - 969
dir.path(), - 970
); - 971
assert!(handle.is_none()); - 972
assert!(CommitmentLedger::new(dir.path()).all().is_empty()); - 973
} - 974
- 975
#[test] - 976
fn a_durable_horizon_opens_one_and_records_an_episode() { - 977
let dir = tempfile::tempdir().unwrap(); - 978
let handle = begin_episode( - 979
dir.path(), - 980
&config(), - 981
&intent_with(Evidence::None, Horizon::Durable, 0.9), - 982
"watch the cloud bill every day", - 983
"s1", - 984
dir.path(), - 985
) - 986
.expect("a durable turn opens a commitment"); - 987
let commitment = CommitmentLedger::new(dir.path()) - 988
.get(&handle.commitment_id) - 989
.unwrap() - 990
.unwrap(); - 991
assert_eq!(commitment.episodes.len(), 1); - 992
assert_eq!(commitment.spec.objective, "watch the cloud bill every day"); - 993
} - 994
- 995
/// A commitment on the durable thread, opened directly. - 996
fn open_on_thread(ledger: &CommitmentLedger, cwd: &Path, thread: &str) -> String { - 997
let mut spec = vak_commit::spec_from_reading( - 998
"resume the migration", - 999
vak_intent::Reading { - 1000
horizon: Horizon::Durable,
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.