- 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, - 1001
..Reading::general() - 1002
}, - 1003
Vec::new(), - 1004
cwd.to_path_buf(), - 1005
Economics::default(), - 1006
); - 1007
spec.thread_id = Some(thread.into()); - 1008
ledger.open_commitment(spec).unwrap() - 1009
} - 1010
- 1011
#[test] - 1012
fn resuming_a_session_closes_its_orphaned_episode_before_starting_again() { - 1013
let dir = tempfile::tempdir().unwrap(); - 1014
let ledger = CommitmentLedger::new(dir.path()); - 1015
let id = open_on_thread(&ledger, dir.path(), "t0.0"); - 1016
ledger - 1017
.append(&Event::new( - 1018
&id, - 1019
EventKind::EpisodeStarted { - 1020
episode_id: "orphan".into(), - 1021
session_id: "s1".into(), - 1022
}, - 1023
)) - 1024
.unwrap(); - 1025
let intent = intent_on( - 1026
Evidence::None, - 1027
Horizon::Durable, - 1028
0.9, - 1029
"t1.0", - 1030
vak_intent::Lineage::Continues { - 1031
thread_id: "t0.0".into(), - 1032
}, - 1033
); - 1034
let next = begin_episode( - 1035
dir.path(), - 1036
&config(), - 1037
&intent, - 1038
"resume the migration", - 1039
"s1", - 1040
dir.path(), - 1041
) - 1042
.unwrap(); - 1043
assert_eq!(next.commitment_id, id); - 1044
let commitment = ledger.get(&id).unwrap().unwrap(); - 1045
assert!(commitment.episodes[0].ended_at.is_some()); - 1046
assert!(matches!( - 1047
commitment.episodes[0].advancement, - 1048
Some(Advancement::Blocked { .. }) - 1049
)); - 1050
assert_ne!(next.episode_id, "orphan"); - 1051
} - 1052
- 1053
/// A strand that continues a durable thread works on that thread's - 1054
/// commitment even when its own reading is a one-turn request: "now also - 1055
/// check staging" is part of the job it continues. - 1056
#[test] - 1057
fn a_continuation_works_on_its_threads_commitment_whatever_its_own_horizon() { - 1058
let dir = tempfile::tempdir().unwrap(); - 1059
let ledger = CommitmentLedger::new(dir.path()); - 1060
let id = open_on_thread(&ledger, dir.path(), "t0.0"); - 1061
let intent = intent_on( - 1062
Evidence::None, - 1063
Horizon::Turn, - 1064
0.9, - 1065
"t1.0", - 1066
vak_intent::Lineage::Continues { - 1067
thread_id: "t0.0".into(), - 1068
}, - 1069
); - 1070
let handle = begin_episode( - 1071
dir.path(), - 1072
&config(), - 1073
&intent, - 1074
"also staging", - 1075
"s2", - 1076
dir.path(), - 1077
) - 1078
.expect("the continuation joins the open commitment"); - 1079
assert_eq!(handle.commitment_id, id); - 1080
assert_eq!(ledger.all().len(), 1, "no twin was opened"); - 1081
} - 1082
- 1083
/// `/goal replace` on durable work opens its successor and supersedes the - 1084
/// original, keyed by the replacement's own thread. - 1085
#[test] - 1086
fn a_replacement_supersedes_the_replaced_threads_commitment() { - 1087
let dir = tempfile::tempdir().unwrap(); - 1088
let ledger = CommitmentLedger::new(dir.path()); - 1089
let old = open_on_thread(&ledger, dir.path(), "t0.0"); - 1090
let intent = intent_on( - 1091
Evidence::None, - 1092
Horizon::Turn, - 1093
0.9, - 1094
"t1.0", - 1095
vak_intent::Lineage::Replaces { - 1096
thread_id: "t0.0".into(), - 1097
}, - 1098
); - 1099
let handle = begin_episode( - 1100
dir.path(), - 1101
&config(), - 1102
&intent, - 1103
"do this instead", - 1104
"s2", - 1105
dir.path(), - 1106
) - 1107
.expect("the replacement opens a successor"); - 1108
assert_ne!(handle.commitment_id, old); - 1109
let old = ledger.get(&old).unwrap().unwrap(); - 1110
assert_eq!( - 1111
old.superseded_by.as_deref(), - 1112
Some(handle.commitment_id.as_str()) - 1113
); - 1114
let new = ledger.get(&handle.commitment_id).unwrap().unwrap(); - 1115
assert_eq!(new.spec.thread_id.as_deref(), Some("t1.0")); - 1116
assert_eq!(new.spec.audience_id.as_deref(), Some("local")); - 1117
} - 1118
- 1119
/// Two turns that each ask for durable work open two commitments: strand - 1120
/// ids carry the turn id, so the second can never mistake the first's - 1121
/// thread for its own. - 1122
#[test] - 1123
fn separate_durable_requests_open_separate_commitments() { - 1124
let dir = tempfile::tempdir().unwrap(); - 1125
for (turn, prompt) in [ - 1126
("t1.0", "watch the bill daily"), - 1127
("t2.0", "watch the logs daily"), - 1128
] { - 1129
let intent = intent_on( - 1130
Evidence::None, - 1131
Horizon::Durable, - 1132
0.9, - 1133
turn, - 1134
vak_intent::Lineage::New, - 1135
); - 1136
begin_episode(dir.path(), &config(), &intent, prompt, "s1", dir.path()).unwrap(); - 1137
} - 1138
assert_eq!(CommitmentLedger::new(dir.path()).open().len(), 2); - 1139
} - 1140
- 1141
/// Only an existing commitment can carry a grant, and only a live one is - 1142
/// offered to the turn. - 1143
#[test] - 1144
fn the_plan_offers_only_live_envelopes_of_existing_commitments() { - 1145
let dir = tempfile::tempdir().unwrap(); - 1146
let ledger = CommitmentLedger::new(dir.path()); - 1147
let id = open_on_thread(&ledger, dir.path(), "t0.0"); - 1148
let now = chrono::Utc::now(); - 1149
let envelope = vak_intent::Envelope { - 1150
envelope_id: "env-1".into(), - 1151
granted_by: "owner".into(), - 1152
granted_at: now, - 1153
expires_at: None, - 1154
spend_limit_usd: Some(5.0), - 1155
path_scope: vec!["docs/**".into()], - 1156
tool_scope: Vec::new(), - 1157
permission_ceiling: vak_intent::PermissionCeiling::WorkspaceWrite, - 1158
escalation: vak_intent::Escalation::WaitIndefinitely, - 1159
revoked_at: None, - 1160
}; - 1161
ledger - 1162
.append(&Event::new( - 1163
&id, - 1164
EventKind::EnvelopeGranted { - 1165
envelope: Box::new(envelope), - 1166
}, - 1167
)) - 1168
.unwrap(); - 1169
let intent = intent_on( - 1170
Evidence::None, - 1171
Horizon::Turn, - 1172
0.9, - 1173
"t1.0", - 1174
vak_intent::Lineage::Continues { - 1175
thread_id: "t0.0".into(), - 1176
}, - 1177
); - 1178
let plan = plan_episodes(dir.path(), &config(), &intent, now); - 1179
assert!(plan.envelopes().contains_key("t1.0")); - 1180
assert_eq!(plan.enveloped_commitments(), vec![id.clone()]); - 1181
- 1182
ledger - 1183
.append(&Event::new( - 1184
&id, - 1185
EventKind::EnvelopeRevoked { - 1186
envelope_id: "env-1".into(), - 1187
by: "owner".into(), - 1188
}, - 1189
)) - 1190
.unwrap(); - 1191
let plan = plan_episodes(dir.path(), &config(), &intent, now); - 1192
assert!( - 1193
plan.envelopes().is_empty(), - 1194
"a revoked grant is not offered" - 1195
); - 1196
} - 1197
- 1198
/// The upkeep evaluator runs outside any sandbox, so a criterion naming a - 1199
/// path outside the workspace is not checked at all. - 1200
#[tokio::test] - 1201
async fn the_workspace_evaluator_refuses_paths_outside_the_workspace() { - 1202
let outer = tempfile::tempdir().unwrap(); - 1203
let workspace = outer.path().join("ws"); - 1204
std::fs::create_dir_all(&workspace).unwrap(); - 1205
std::fs::write(outer.path().join("secret.txt"), "hunter2").unwrap(); - 1206
#[cfg(unix)] - 1207
std::os::unix::fs::symlink(outer.path(), workspace.join("up")).unwrap(); - 1208
let evaluator = WorkspaceEvaluator { cwd: &workspace }; - 1209
let mut paths = vec![ - 1210
"../secret.txt".to_string(), - 1211
outer.path().join("secret.txt").display().to_string(), - 1212
]; - 1213
if cfg!(unix) { - 1214
paths.push("up/secret.txt".into()); - 1215
} - 1216
for path in paths { - 1217
let evaluation = evaluator - 1218
.evaluate(&WorkCriterion { - 1219
criterion_id: "leak".into(), - 1220
statement: "probe".into(), - 1221
kind: CriterionKind::FileContains { - 1222
path: path.clone().into(), - 1223
pattern: "hunter2".into(), - 1224
}, - 1225
required: true, - 1226
}) - 1227
.await; - 1228
assert!( - 1229
matches!(evaluation.result, CriterionResult::Unknown { .. }), - 1230
"{path} was checked: {:?}", - 1231
evaluation.result - 1232
); - 1233
} - 1234
} - 1235
- 1236
/// Upkeep reaches every Agent's portfolio, not only the server's own. - 1237
#[tokio::test] - 1238
async fn upkeep_covers_every_agents_ledger() { - 1239
let data = tempfile::tempdir().unwrap(); - 1240
let agent_home = data.path().join("agents").join("helper"); - 1241
std::fs::create_dir_all(&agent_home).unwrap(); - 1242
let ledger = CommitmentLedger::new(&agent_home); - 1243
let id = open_on_thread(&ledger, data.path(), "t0.0"); - 1244
ledger - 1245
.append(&Event::new( - 1246
&id, - 1247
EventKind::Suspended { - 1248
suspension: vak_commit::Suspension::Schedule { - 1249
at: Some(chrono::Utc::now() - chrono::Duration::minutes(1)), - 1250
cron: None, - 1251
}, - 1252
}, - 1253
)) - 1254
.unwrap(); - 1255
assert_eq!(maintain_all(data.path()).await.resumed, vec![id]); - 1256
} - 1257
- 1258
/// A stray recurrence-ish word must not leave a month-long obligation - 1259
/// behind. Weak horizon evidence declines to open one. - 1260
#[test] - 1261
fn a_weak_horizon_reading_opens_nothing() { - 1262
let dir = tempfile::tempdir().unwrap(); - 1263
let handle = begin_episode( - 1264
dir.path(), - 1265
&config(), - 1266
&intent_with(Evidence::None, Horizon::Durable, 0.3), - 1267
"maybe keep an eye on things", - 1268
"s1", - 1269
dir.path(), - 1270
);
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.