- 3299
.map_err(|error| format!("managed work projection is invalid: {error}"))? - 3300
else { - 3301
return Ok(()); - 3302
}; - 3303
if projection.status != vak_session::types::WorkContractStatus::AwaitingInput { - 3304
return Ok(()); - 3305
} - 3306
let Some(answer) = answer.trim().strip_prefix("answer:").map(str::trim) else { - 3307
return Ok(()); - 3308
}; - 3309
let Some(assumption) = - 3310
projection.contract.assumptions.iter().find(|assumption| { - 3311
assumption.requires_confirmation && assumption.resolution.is_none() - 3312
}) - 3313
else { - 3314
return Ok(()); - 3315
}; - 3316
if answer.is_empty() { - 3317
return Err(format!( - 3318
"cannot resolve assumption '{}' with an empty answer", - 3319
assumption.assumption_id - 3320
)); - 3321
} - 3322
session - 3323
.append_work(vak_session::types::WorkEvent { - 3324
contract_id: projection.contract.contract_id.clone(), - 3325
revision: projection.contract.revision, - 3326
kind: vak_session::types::WorkEventKind::AssumptionResolved { - 3327
assumption_id: assumption.assumption_id.clone(), - 3328
resolution: answer.into(), - 3329
}, - 3330
}) - 3331
.map_err(|error| format!("managed assumption resolution failed: {error}"))?; - 3332
let Some(updated) = session - 3333
.work_projection() - 3334
.map_err(|error| format!("managed work projection is invalid: {error}"))? - 3335
else { - 3336
return Ok(()); - 3337
}; - 3338
if updated.status == vak_session::types::WorkContractStatus::AwaitingInput - 3339
&& updated - 3340
.contract - 3341
.assumptions - 3342
.iter() - 3343
.filter(|assumption| assumption.requires_confirmation) - 3344
.all(|assumption| assumption.resolution.is_some()) - 3345
{ - 3346
session - 3347
.append_work(vak_session::types::WorkEvent { - 3348
contract_id: updated.contract.contract_id, - 3349
revision: updated.contract.revision, - 3350
kind: vak_session::types::WorkEventKind::ContractStatusChanged { - 3351
from: vak_session::types::WorkContractStatus::AwaitingInput, - 3352
to: vak_session::types::WorkContractStatus::Active, - 3353
reason: "all required assumptions resolved from chat".into(), - 3354
}, - 3355
}) - 3356
.map_err(|error| format!("managed work activation failed: {error}"))?; - 3357
} - 3358
Ok(()) - 3359
} - 3360
- 3361
async fn start_managed_contract( - 3362
&self, - 3363
prompt: &str, - 3364
source_entry_id: String, - 3365
cancel: &CancellationToken, - 3366
events: &mpsc::Sender<AgentEvent>, - 3367
) -> Result<(), String> { - 3368
if self - 3369
.session - 3370
.lock() - 3371
.await - 3372
.work_projection() - 3373
.map_err(|error| format!("managed work projection is invalid: {error}"))? - 3374
.is_some_and(|work| { - 3375
!matches!( - 3376
work.status, - 3377
vak_session::types::WorkContractStatus::Completed - 3378
| vak_session::types::WorkContractStatus::Failed - 3379
| vak_session::types::WorkContractStatus::Cancelled - 3380
| vak_session::types::WorkContractStatus::Unverified - 3381
) - 3382
}) - 3383
{ - 3384
return Ok(()); - 3385
} - 3386
let session_id = self - 3387
.session - 3388
.lock() - 3389
.await - 3390
.header() - 3391
.map(|header| header.session_id.clone()) - 3392
.ok_or_else(|| "managed work requires a session header".to_string())?; - 3393
// config.model reflects the per-turn effective route set by run_turn_inner. - 3394
let model = self.config.model.clone(); - 3395
let authoring_request = ChatRequest { - 3396
model: model.clone(), - 3397
system: Some("You author durable work contracts. Return only one strict JSON object with keys objective, constraints, assumptions, criteria, and items. Each item must have item_id, title, instructions, dependencies, owner, required, readonly, path_claims, and criterion_ids. Owner must be one of parent_agent, worker, flow, tool, or human. Criterion kind must be one of shell, file_exists, file_contains, tool_succeeded, flow_completed, external_receipt, or semantic. Do not include markdown or commentary.".into()), - 3398
messages: vec![Message::user_text(prompt)], - 3399
tools: Vec::new(), - 3400
max_tokens: self.config.max_output.min(8_000) as u32, - 3401
temperature: None, - 3402
cache: None, - 3403
previous_response_id: None, - 3404
think: None, - 3405
effort: None, - 3406
}; - 3407
let mut ledger = StepLedger::new( - 3408
WorkPurpose::Plan, - 3409
self.provider.name(), - 3410
&model, - 3411
self.config.dispatch_ceiling.min(4), - 3412
); - 3413
let response = self - 3414
.complete_with_reliability(&authoring_request, cancel, events, false, &mut ledger) - 3415
.await - 3416
.map_err(|error| format!("managed contract authoring failed: {error}"))?; - 3417
self.session - 3418
.lock() - 3419
.await - 3420
.append_receipt(ledger.take_receipt()) - 3421
.map_err(|error| format!("managed contract receipt failed: {error}"))?; - 3422
let authored: AuthoredContract = - 3423
serde_json::from_str(&response.text_content()).map_err(|error| { - 3424
format!("managed contract authoring returned invalid JSON: {error}") - 3425
})?; - 3426
if authored.items.is_empty() || authored.items.len() > self.config.max_work_items { - 3427
return Err(format!( - 3428
"managed contract must contain between one and {} items", - 3429
self.config.max_work_items - 3430
)); - 3431
} - 3432
if authored.objective.trim().is_empty() || authored.objective.chars().count() > 16_000 { - 3433
return Err("managed contract objective is empty or too long".into()); - 3434
} - 3435
let contract_id = format!( - 3436
"work-{session_id}-{}", - 3437
chrono::Utc::now().timestamp_millis() - 3438
); - 3439
let contract = vak_session::types::WorkContract { - 3440
contract_id: contract_id.clone(), - 3441
revision: 0, - 3442
source_entry_id, - 3443
objective: authored.objective, - 3444
constraints: authored.constraints, - 3445
assumptions: authored.assumptions, - 3446
criteria: authored.criteria, - 3447
items: authored.items, - 3448
}; - 3449
vak_session::validate_contract_for_admission(&contract) - 3450
.map_err(|error| format!("managed contract validation failed: {error}"))?; - 3451
validate_work_paths(&contract)?; - 3452
let mut session = self.session.lock().await; - 3453
session - 3454
.append_work(vak_session::types::WorkEvent { - 3455
contract_id: contract_id.clone(), - 3456
revision: 0, - 3457
kind: vak_session::types::WorkEventKind::ContractCreated { contract }, - 3458
}) - 3459
.map_err(|error| format!("managed contract write failed: {error}"))?; - 3460
let awaiting_input = session - 3461
.work_projection() - 3462
.ok() - 3463
.flatten() - 3464
.is_some_and(|work| { - 3465
work.contract.assumptions.iter().any(|assumption| { - 3466
assumption.requires_confirmation && assumption.resolution.is_none() - 3467
}) - 3468
}); - 3469
session - 3470
.append_work(vak_session::types::WorkEvent { - 3471
contract_id: contract_id.clone(), - 3472
revision: 0, - 3473
kind: vak_session::types::WorkEventKind::ContractStatusChanged { - 3474
from: vak_session::types::WorkContractStatus::Draft, - 3475
to: if awaiting_input { - 3476
vak_session::types::WorkContractStatus::AwaitingInput - 3477
} else { - 3478
vak_session::types::WorkContractStatus::Active - 3479
}, - 3480
reason: if awaiting_input { - 3481
"required assumptions need confirmation".into() - 3482
} else { - 3483
"managed execution admitted".into() - 3484
}, - 3485
}, - 3486
}) - 3487
.map_err(|error| format!("managed contract activation failed: {error}"))?; - 3488
Ok(()) - 3489
} - 3490
- 3491
async fn emit_work_state(&self, events: &mpsc::Sender<AgentEvent>) { - 3492
let projection = self.session.lock().await.work_projection().ok().flatten(); - 3493
if let Some(projection) = projection { - 3494
let _ = events.send(AgentEvent::WorkState { projection }).await; - 3495
} - 3496
} - 3497
- 3498
async fn managed_work_gate( - 3499
&mut self, - 3500
cancel: &CancellationToken, - 3501
events: &mpsc::Sender<AgentEvent>, - 3502
) -> Option<String> { - 3503
if !self.config.work_enabled || self.config.work_mode != WorkMode::Managed { - 3504
return None; - 3505
} - 3506
let session = self.session.lock().await; - 3507
let projection = match session.work_projection() { - 3508
Ok(Some(projection)) => projection, - 3509
Ok(None) => return Some("managed run has no durable work contract".into()), - 3510
Err(error) => return Some(format!("managed work projection is invalid: {error}")), - 3511
}; - 3512
if matches!( - 3513
projection.status, - 3514
vak_session::types::WorkContractStatus::Completed - 3515
| vak_session::types::WorkContractStatus::Failed - 3516
| vak_session::types::WorkContractStatus::Cancelled - 3517
| vak_session::types::WorkContractStatus::Unverified - 3518
) { - 3519
return None; - 3520
} - 3521
let criteria = projection.contract.criteria.clone(); - 3522
drop(session); - 3523
self.verify_managed_criteria(&criteria, &projection, cancel, events) - 3524
.await; - 3525
let mut session = self.session.lock().await; - 3526
let projection = match session.work_projection() { - 3527
Ok(Some(projection)) => projection, - 3528
Ok(None) => return Some("managed run lost its work contract".into()), - 3529
Err(error) => return Some(format!("managed work projection is invalid: {error}")), - 3530
}; - 3531
for item in &projection.contract.items { - 3532
let Some(state) = projection.items.get(&item.item_id) else { - 3533
return Some(format!("managed work item '{}' has no state", item.item_id)); - 3534
}; - 3535
if state.status == vak_session::types::WorkItemStatus::ReadyForVerification { - 3536
let event = vak_session::types::WorkEvent { - 3537
contract_id: projection.contract.contract_id.clone(), - 3538
revision: projection.contract.revision, - 3539
kind: vak_session::types::WorkEventKind::ItemVerified { - 3540
item_id: item.item_id.clone(), - 3541
attempt: state.attempt, - 3542
}, - 3543
}; - 3544
if let Err(error) = session.append_work(event) { - 3545
return Some(format!( - 3546
"managed verification pending for '{}': {error}", - 3547
item.item_id - 3548
)); - 3549
} - 3550
} - 3551
} - 3552
let projection = match session.work_projection() { - 3553
Ok(Some(projection)) => projection, - 3554
Ok(None) => return Some("managed run lost its work contract".into()), - 3555
Err(error) => return Some(format!("managed work projection is invalid: {error}")), - 3556
}; - 3557
let incomplete: Vec<&str> = projection - 3558
.contract - 3559
.items - 3560
.iter() - 3561
.filter(|item| item.required) - 3562
.filter_map(|item| { - 3563
let state = projection.items.get(&item.item_id)?; - 3564
(!matches!(state.status, vak_session::types::WorkItemStatus::Succeeded)) - 3565
.then_some(item.item_id.as_str()) - 3566
}) - 3567
.collect(); - 3568
if !incomplete.is_empty() { - 3569
return Some(format!( - 3570
"managed work is not complete; required items pending: {}. Use the work tool and do not claim completion.", - 3571
incomplete.join(", ") - 3572
)); - 3573
} - 3574
let missing_criteria: Vec<&str> = projection - 3575
.contract - 3576
.criteria - 3577
.iter() - 3578
.filter(|criterion| criterion.required) - 3579
.filter_map(|criterion| { - 3580
(!matches!( - 3581
projection.criteria.get(&criterion.criterion_id), - 3582
Some(vak_session::types::CriterionResult::Passed { .. }) - 3583
)) - 3584
.then_some(criterion.criterion_id.as_str()) - 3585
}) - 3586
.collect(); - 3587
if !missing_criteria.is_empty() { - 3588
return Some(format!( - 3589
"managed work cannot complete; required criteria are not proven: {}", - 3590
missing_criteria.join(", ") - 3591
)); - 3592
} - 3593
if projection.status == vak_session::types::WorkContractStatus::Active { - 3594
let contract_id = projection.contract.contract_id.clone(); - 3595
let revision = projection.contract.revision; - 3596
if let Err(error) = session.append_work(vak_session::types::WorkEvent { - 3597
contract_id: contract_id.clone(), - 3598
revision, - 3599
kind: vak_session::types::WorkEventKind::ContractStatusChanged { - 3600
from: vak_session::types::WorkContractStatus::Active, - 3601
to: vak_session::types::WorkContractStatus::Verifying, - 3602
reason: "required work items verified".into(), - 3603
}, - 3604
}) { - 3605
return Some(format!("managed verification status failed: {error}")); - 3606
} - 3607
if let Err(error) = session.append_work(vak_session::types::WorkEvent { - 3608
contract_id, - 3609
revision, - 3610
kind: vak_session::types::WorkEventKind::ContractStatusChanged { - 3611
from: vak_session::types::WorkContractStatus::Verifying, - 3612
to: vak_session::types::WorkContractStatus::Completed, - 3613
reason: "all required work items verified".into(), - 3614
}, - 3615
}) { - 3616
return Some(format!("managed completion status failed: {error}")); - 3617
} - 3618
} - 3619
drop(session); - 3620
self.emit_work_state(events).await; - 3621
None - 3622
} - 3623
- 3624
async fn verify_managed_criteria( - 3625
&mut self, - 3626
criteria: &[vak_session::types::WorkCriterion], - 3627
projection: &vak_session::work::WorkProjection, - 3628
cancel: &CancellationToken, - 3629
events: &mpsc::Sender<AgentEvent>, - 3630
) { - 3631
let cwd = self - 3632
.session - 3633
.lock() - 3634
.await - 3635
.header() - 3636
.map(|header| header.contract_cwd()) - 3637
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| ".".into())); - 3638
for criterion in criteria { - 3639
let result = match &criterion.kind { - 3640
vak_session::types::CriterionKind::Shell { command } => self - 3641
.run_audit_command(command, cancel) - 3642
.await - 3643
.map(|_| vak_session::types::CriterionResult::Passed { - 3644
evidence: format!("shell:{command}"), - 3645
}) - 3646
.unwrap_or_else(|reason| vak_session::types::CriterionResult::Failed { - 3647
reason, - 3648
}), - 3649
vak_session::types::CriterionKind::FileExists { path } => { - 3650
let Some(resolved) = workspace_criterion_path(&cwd, path) else { - 3651
self.record_managed_criterion( - 3652
criterion, - 3653
vak_session::types::CriterionResult::Failed { - 3654
reason: format!("path is outside the workspace: {}", path.display()), - 3655
}, - 3656
) - 3657
.await; - 3658
continue; - 3659
}; - 3660
if resolved.exists() { - 3661
vak_session::types::CriterionResult::Passed { - 3662
evidence: format!("file_exists:{}", path.display()), - 3663
} - 3664
} else { - 3665
vak_session::types::CriterionResult::Failed { - 3666
reason: format!("file does not exist: {}", path.display()), - 3667
} - 3668
} - 3669
} - 3670
vak_session::types::CriterionKind::FileContains { path, pattern } => { - 3671
let Some(resolved) = workspace_criterion_path(&cwd, path) else { - 3672
self.record_managed_criterion( - 3673
criterion, - 3674
vak_session::types::CriterionResult::Failed { - 3675
reason: format!("path is outside the workspace: {}", path.display()), - 3676
}, - 3677
) - 3678
.await; - 3679
continue; - 3680
}; - 3681
match std::fs::read_to_string(&resolved) { - 3682
Ok(content) if content.contains(pattern) => { - 3683
vak_session::types::CriterionResult::Passed { - 3684
evidence: format!("file_contains:{}", path.display()), - 3685
} - 3686
} - 3687
Ok(_) => vak_session::types::CriterionResult::Failed { - 3688
reason: format!("pattern not found in {}", path.display()), - 3689
}, - 3690
Err(error) => vak_session::types::CriterionResult::Failed { - 3691
reason: format!("cannot read {}: {error}", path.display()), - 3692
}, - 3693
} - 3694
} - 3695
vak_session::types::CriterionKind::Semantic => continue, - 3696
vak_session::types::CriterionKind::ToolSucceeded { tool } => { - 3697
let mut succeeded = false; - 3698
for item_id in self.criterion_item_ids(projection, &criterion.criterion_id) { - 3699
if self - 3700
.tool_succeeded_for_item(projection, &item_id, tool) - 3701
.await - 3702
{ - 3703
succeeded = true; - 3704
break; - 3705
} - 3706
} - 3707
if succeeded { - 3708
vak_session::types::CriterionResult::Passed { - 3709
evidence: format!("tool_succeeded:{tool}"), - 3710
} - 3711
} else { - 3712
vak_session::types::CriterionResult::Unknown { - 3713
reason: format!("no successful '{tool}' tool result exists yet"), - 3714
} - 3715
} - 3716
} - 3717
vak_session::types::CriterionKind::FlowCompleted { flow } => { - 3718
if self.criterion_item_ids(projection, &criterion.criterion_id).into_iter().any( - 3719
|item_id| { - 3720
projection.items.get(&item_id).is_some_and(|item| { - 3721
item.evidence.iter().any(|evidence| { - 3722
matches!(evidence, vak_session::types::EvidenceRef::FlowNode { flow: name, node_id, .. } if name == flow && node_id == "__flow_completed__") - 3723
}) - 3724
}) - 3725
}, - 3726
) { - 3727
vak_session::types::CriterionResult::Passed { - 3728
evidence: format!("flow_completed:{flow}"), - 3729
} - 3730
} else { - 3731
vak_session::types::CriterionResult::Unknown { - 3732
reason: format!("flow '{flow}' has no linked completion evidence"), - 3733
} - 3734
} - 3735
} - 3736
vak_session::types::CriterionKind::ExternalReceipt { integration } => { - 3737
if self.criterion_item_ids(projection, &criterion.criterion_id).into_iter().any( - 3738
|item_id| { - 3739
projection.items.get(&item_id).is_some_and(|item| { - 3740
item.evidence.iter().any(|evidence| { - 3741
matches!(evidence, vak_session::types::EvidenceRef::ExternalOperation { integration: name, .. } if name == integration) - 3742
}) - 3743
}) - 3744
}, - 3745
) { - 3746
vak_session::types::CriterionResult::Passed { - 3747
evidence: format!("external_receipt:{integration}"), - 3748
} - 3749
} else { - 3750
vak_session::types::CriterionResult::Unknown { - 3751
reason: format!("integration '{integration}' has no receipt evidence"), - 3752
} - 3753
} - 3754
} - 3755
}; - 3756
self.record_managed_criterion(criterion, result).await; - 3757
} - 3758
let semantic: Vec<(vak_session::types::WorkCriterion, String)> = criteria - 3759
.iter() - 3760
.filter_map(|criterion| match &criterion.kind { - 3761
vak_session::types::CriterionKind::Semantic => Some(( - 3762
criterion.clone(), - 3763
format!("[{}] {}", criterion.criterion_id, criterion.statement), - 3764
)), - 3765
_ => None, - 3766
}) - 3767
.collect(); - 3768
if semantic.is_empty() || cancel.is_cancelled() { - 3769
return; - 3770
} - 3771
let judge_criteria: Vec<String> = semantic.iter().map(|(_, text)| text.clone()).collect(); - 3772
match self.run_judge(&judge_criteria, cancel, events).await { - 3773
Ok(verdicts) => { - 3774
for (criterion, expected) in semantic { - 3775
let result = verdicts - 3776
.iter() - 3777
.find(|verdict| verdict.criterion == expected) - 3778
.map(|verdict| match verdict.verdict.as_str() { - 3779
"pass" => vak_session::types::CriterionResult::Passed { - 3780
evidence: verdict.evidence.clone(), - 3781
}, - 3782
"fail" => vak_session::types::CriterionResult::Failed { - 3783
reason: verdict.evidence.clone(), - 3784
}, - 3785
_ => vak_session::types::CriterionResult::Unknown { - 3786
reason: verdict.evidence.clone(), - 3787
}, - 3788
}) - 3789
.unwrap_or_else(|| vak_session::types::CriterionResult::Unknown { - 3790
reason: "judge returned no verdict for this criterion".into(), - 3791
}); - 3792
self.record_managed_criterion(&criterion, result).await; - 3793
} - 3794
} - 3795
Err(reason) => { - 3796
for (criterion, _) in semantic { - 3797
self.record_managed_criterion( - 3798
&criterion, - 3799
vak_session::types::CriterionResult::Unknown { - 3800
reason: reason.clone(), - 3801
}, - 3802
) - 3803
.await; - 3804
} - 3805
} - 3806
} - 3807
} - 3808
- 3809
async fn record_managed_criterion( - 3810
&self, - 3811
criterion: &vak_session::types::WorkCriterion, - 3812
result: vak_session::types::CriterionResult, - 3813
) { - 3814
let contract_id = self - 3815
.session - 3816
.lock() - 3817
.await - 3818
.work_projection() - 3819
.ok() - 3820
.flatten() - 3821
.map(|projection| { - 3822
( - 3823
projection.contract.contract_id, - 3824
projection.contract.revision, - 3825
) - 3826
}); - 3827
if let Some((contract_id, revision)) = contract_id { - 3828
let _ = self - 3829
.session - 3830
.lock() - 3831
.await - 3832
.append_work(vak_session::types::WorkEvent { - 3833
contract_id, - 3834
revision, - 3835
kind: vak_session::types::WorkEventKind::VerificationRecorded { - 3836
criterion_id: criterion.criterion_id.clone(), - 3837
result, - 3838
}, - 3839
}); - 3840
} - 3841
} - 3842
- 3843
fn criterion_item_ids( - 3844
&self, - 3845
projection: &vak_session::work::WorkProjection, - 3846
criterion_id: &str, - 3847
) -> Vec<String> { - 3848
projection - 3849
.contract - 3850
.items - 3851
.iter() - 3852
.filter(|item| item.criterion_ids.iter().any(|id| id == criterion_id)) - 3853
.map(|item| item.item_id.clone()) - 3854
.collect() - 3855
} - 3856
- 3857
async fn tool_succeeded_for_item( - 3858
&self, - 3859
projection: &vak_session::work::WorkProjection, - 3860
item_id: &str, - 3861
tool: &str, - 3862
) -> bool { - 3863
let Some(item) = projection.items.get(item_id) else { - 3864
return false; - 3865
}; - 3866
let tool_result_ids: std::collections::HashSet<&str> = item - 3867
.evidence - 3868
.iter() - 3869
.filter_map(|evidence| match evidence { - 3870
vak_session::types::EvidenceRef::ToolResult { tool_use_id, .. } => { - 3871
Some(tool_use_id.as_str()) - 3872
} - 3873
_ => None, - 3874
}) - 3875
.collect(); - 3876
if tool_result_ids.is_empty() { - 3877
return false; - 3878
} - 3879
let session = self.session.lock().await; - 3880
let mut tool_names = std::collections::HashMap::new(); - 3881
for entry in session.chain_to_root() { - 3882
if let vak_session::types::EntryPayload::Message(record) = &entry.payload { - 3883
for block in &record.message.content { - 3884
if let vak_llm::ContentBlock::ToolUse { id, name, .. } = block { - 3885
tool_names.insert(id.as_str(), name.as_str()); - 3886
} - 3887
} - 3888
} - 3889
} - 3890
session.chain_to_root().iter().any(|entry| { - 3891
let vak_session::types::EntryPayload::Message(record) = &entry.payload else { - 3892
return false; - 3893
}; - 3894
record.message.content.iter().any(|block| { - 3895
matches!( - 3896
block, - 3897
vak_llm::ContentBlock::ToolResult { - 3898
tool_use_id, - 3899
is_error: false, - 3900
.. - 3901
} if tool_result_ids.contains(tool_use_id.as_str()) - 3902
&& tool_names.get(tool_use_id.as_str()).copied() == Some(tool) - 3903
) - 3904
}) - 3905
}) - 3906
} - 3907
- 3908
/// Goal audit gate (Phase H): runs when the model claims completion. - 3909
/// Some(reason) rejects the claim and continues the run; None lets it - 3910
/// end. Completion is recorded from audit, never self-report — and the - 3911
/// audit budget is capped so this can never trap a run. - 3912
async fn goal_gate( - 3913
&mut self, - 3914
_response: &AssistantMessage, - 3915
cancel: &CancellationToken, - 3916
events: &mpsc::Sender<AgentEvent>, - 3917
) -> Option<String> { - 3918
let goal_active = self.active_goal.is_some(); - 3919
if !goal_active { - 3920
return None; - 3921
} - 3922
- 3923
let mut findings = String::new(); - 3924
- 3925
// 1) Regression obligations: everything proven green must stay so. - 3926
for cmd in self.obligations.clone() { - 3927
if cancel.is_cancelled() { - 3928
return None; - 3929
} - 3930
match self.run_audit_command(&cmd, cancel).await { - 3931
Ok(()) => {} - 3932
Err(err) => { - 3933
findings.push_str(&format!( - 3934
"REGRESSION: previously-green command failed now:\n $ {cmd}\n {err}\n" - 3935
)); - 3936
} - 3937
} - 3938
} - 3939
- 3940
let criteria = self - 3941
.active_goal - 3942
.as_ref() - 3943
.map(|g| g.criteria.clone()) - 3944
.unwrap_or_default(); - 3945
- 3946
// 2) Deterministic shell criteria. - 3947
let mut judged_criteria: Vec<String> = Vec::new(); - 3948
for criterion in &criteria { - 3949
if goal::is_shell_criterion(criterion) { - 3950
let cmd = goal::shell_command(criterion); - 3951
match self.run_audit_command(cmd, cancel).await { - 3952
Ok(()) => {} - 3953
Err(err) => { - 3954
findings.push_str(&format!("CRITERION FAILED: {criterion}\n {err}\n")); - 3955
} - 3956
} - 3957
judged_criteria.push(criterion.clone()); - 3958
} - 3959
} - 3960
- 3961
// 3) Judge call for remaining free-text criteria — only worth a - 3962
// model dispatch when deterministic checks already passed. - 3963
let text_criteria: Vec<String> = criteria - 3964
.iter() - 3965
.filter(|c| !goal::is_shell_criterion(c)) - 3966
.cloned() - 3967
.collect(); - 3968
if !text_criteria.is_empty() && findings.is_empty() { - 3969
match self.run_judge(&text_criteria, cancel, events).await { - 3970
Ok(verdicts) => { - 3971
for v in verdicts { - 3972
if v.verdict != "pass" { - 3973
findings.push_str(&format!( - 3974
"JUDGE {}: evidence: {}\n", - 3975
v.verdict.to_uppercase(), - 3976
v.evidence - 3977
)); - 3978
} - 3979
} - 3980
} - 3981
Err(e) => { - 3982
// Fail closed: an unavailable auditor cannot confirm done. - 3983
findings.push_str(&format!("AUDIT UNAVAILABLE: {e}\n")); - 3984
} - 3985
} - 3986
} - 3987
- 3988
if findings.is_empty() { - 3989
let mut session = self.session.lock().await; - 3990
let goal_snapshot = self.active_goal.clone(); - 3991
let sid = session - 3992
.header() - 3993
.map(|h| h.session_id.clone()) - 3994
.unwrap_or_default(); - 3995
if let Some(g) = goal_snapshot.as_ref() { - 3996
let _ = session.append_goal(vak_session::types::GoalEntry { - 3997
goal_id: format!("goal-{sid}"), - 3998
objective: g.objective.clone(), - 3999
criteria: g.criteria.clone(), - 4000
status: vak_session::types::GoalStatus::Done { audited: true }, - 4001
}); - 4002
} - 4003
drop(session); - 4004
self.active_goal = None; - 4005
return None; - 4006
} - 4007
- 4008
let audits_left = self - 4009
.active_goal - 4010
.as_mut() - 4011
.map(|g| { - 4012
g.audits_left = g.audits_left.saturating_sub(1); - 4013
g.audits_left - 4014
}) - 4015
.unwrap_or(0); - 4016
if audits_left > 0 { - 4017
Some(format!( - 4018
"[goal-audit] not verified yet:\n{findings}\nAddress these and finish again." - 4019
)) - 4020
} else { - 4021
// Budget exhausted: degrade to Unverified rather than trapping. - 4022
let mut session = self.session.lock().await; - 4023
let goal_snapshot = self.active_goal.clone(); - 4024
let sid = session - 4025
.header() - 4026
.map(|h| h.session_id.clone()) - 4027
.unwrap_or_default(); - 4028
if let Some(g) = goal_snapshot.as_ref() { - 4029
let _ = session.append_goal(vak_session::types::GoalEntry { - 4030
goal_id: format!("goal-{sid}"), - 4031
objective: g.objective.clone(), - 4032
criteria: g.criteria.clone(), - 4033
status: vak_session::types::GoalStatus::Unverified { - 4034
reason: findings.clone(), - 4035
}, - 4036
}); - 4037
} - 4038
drop(session); - 4039
self.active_goal = None; - 4040
None - 4041
} - 4042
} - 4043
- 4044
/// Runs one brokered bash command for auditing; Err = failure text. - 4045
async fn run_audit_command(&self, cmd: &str, cancel: &CancellationToken) -> Result<(), String> { - 4046
let tool = self - 4047
.config - 4048
.tools - 4049
.iter() - 4050
.find(|t| t.name() == "bash") - 4051
.ok_or_else(|| "no bash tool available for verification".to_string())?; - 4052
let cwd = self - 4053
.session - 4054
.lock() - 4055
.await - 4056
.header() - 4057
.map(|h| h.contract_cwd()) - 4058
.unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| ".".into())); - 4059
authorize( - 4060
&self.config, - 4061
&PendingToolCall { - 4062
id: "managed-verification".into(), - 4063
name: "bash".into(), - 4064
input: serde_json::json!({"command": cmd}), - 4065
}, - 4066
&cwd, - 4067
&self.run_call_counts, - 4068
&self.config.tools, - 4069
) - 4070
.await?; - 4071
let Some(sandbox) = self - 4072
.config - 4073
.sandbox - 4074
.as_ref() - 4075
.and_then(|sandbox| sandbox.read_only_variant()) - 4076
else { - 4077
return Err("shell verification requires a read-only sandbox".into()); - 4078
}; - 4079
let ctx = vak_tools::ToolContext { - 4080
cwd, - 4081
cancel: cancel.child_token(), - 4082
sandbox: Some(sandbox), - 4083
sandbox_sink: None, - 4084
agent_id: None, - 4085
new_documents: Vec::new(), - 4086
}; - 4087
let input = serde_json::json!({ "command": cmd }); - 4088
let out = match tokio::time::timeout( - 4089
std::time::Duration::from_secs(120), - 4090
tool.execute(&input, &ctx), - 4091
) - 4092
.await - 4093
{ - 4094
Ok(o) => o, - 4095
Err(_) => return Err("timed out after 120s".to_string()), - 4096
}; - 4097
if out.is_error { - 4098
Err(vak_session::bash_digest(&out.content)) - 4099
} else { - 4100
Ok(()) - 4101
} - 4102
} - 4103
- 4104
/// One skeptical judge dispatch over the transcript digest. - 4105
async fn run_judge( - 4106
&mut self, - 4107
text_criteria: &[String], - 4108
cancel: &CancellationToken, - 4109
events: &mpsc::Sender<AgentEvent>, - 4110
) -> Result<Vec<goal::CriterionVerdict>, String> { - 4111
// config.model reflects the per-turn effective route set by run_turn_inner. - 4112
let model = self.config.model.clone(); - 4113
let digest = { - 4114
let session = self.session.lock().await; - 4115
let msgs = session.derive_messages(); - 4116
goal::transcript_digest(&msgs, 24_000) - 4117
}; - 4118
let objective = self - 4119
.active_goal - 4120
.as_ref() - 4121
.map(|g| g.objective.clone()) - 4122
.unwrap_or_default(); - 4123
// MEA: environment facts over transcript claims. Unavailable delta - 4124
// is stated as such to the judge (UNKNOWN, never fabricated). - 4125
let workspace_delta = match &self.config.workspace_delta { - 4126
Some(p) => p.summary().ok(), - 4127
None => None, - 4128
}; - 4129
let req = goal::audit_request( - 4130
&model, - 4131
goal::audit_prompt( - 4132
&objective, - 4133
text_criteria, - 4134
&digest, - 4135
workspace_delta.as_deref(), - 4136
), - 4137
); - 4138
let mut ledger = StepLedger::new( - 4139
WorkPurpose::Verify, - 4140
self.provider.name(), - 4141
&model, - 4142
self.config.dispatch_ceiling.min(4), - 4143
); - 4144
let reply = self - 4145
.complete_with_reliability(&req, cancel, events, false, &mut ledger) - 4146
.await - 4147
.map_err(|e| e.to_string())?; - 4148
let _ = self - 4149
.session - 4150
.lock() - 4151
.await - 4152
.append_receipt(ledger.take_receipt()); - 4153
goal::parse_verdicts(&reply.text_content()) - 4154
} - 4155
- 4156
/// Reset-with-handoff rescue: one structured summary replaces the whole - 4157
/// projection; returns the handoff markdown. - 4158
async fn write_handoff( - 4159
&self, - 4160
est_tokens: u64, - 4161
original_prompt: &str, - 4162
cancel: &CancellationToken, - 4163
events: &mpsc::Sender<AgentEvent>, - 4164
) -> Result<String, LlmError> { - 4165
// config.model reflects the per-turn effective route set by run_turn_inner. - 4166
let model = self.config.model.clone(); - 4167
let digest = { - 4168
let session = self.session.lock().await; - 4169
let msgs = session.derive_messages(); - 4170
goal::transcript_digest(&msgs, 20_000) - 4171
}; - 4172
let objective_line = if self.active_goal.is_some() || !original_prompt.is_empty() { - 4173
format!("Original task: {original_prompt}\n\n") - 4174
} else { - 4175
String::new() - 4176
}; - 4177
let req = goal::handoff_request(&model, format!("{objective_line}{digest}")); - 4178
let mut ledger = StepLedger::new( - 4179
WorkPurpose::Summarize, - 4180
self.provider.name(), - 4181
&model, - 4182
self.config.dispatch_ceiling.min(3), - 4183
); - 4184
let _ = events - 4185
.send(AgentEvent::ContextCompacting { - 4186
estimated_tokens: est_tokens, - 4187
}) - 4188
.await; - 4189
let reply = self - 4190
.complete_with_reliability(&req, cancel, events, false, &mut ledger) - 4191
.await?; - 4192
let _ = self - 4193
.session - 4194
.lock() - 4195
.await - 4196
.append_receipt(ledger.take_receipt()); - 4197
Ok(reply.text_content()) - 4198
} - 4199
- 4200
/// Internal premature-completion gate. Returns a continuation reason - 4201
/// when the stop policy fires and budget remains. - 4202
async fn stop_gate( - 4203
&self, - 4204
prompt: &str, - 4205
response: &AssistantMessage, - 4206
receipts: &stop_policy::ReceiptSummary, - 4207
verification_stale: bool, - 4208
blocks_left: &mut u32, - 4209
user_completion_released: bool, - 4210
) -> Option<String> { - 4211
let policy = self.config.stop_policy.as_ref()?; - 4212
if user_completion_released { - 4213
return None; - 4214
} - 4215
let reason = policy.evaluate_receipts( - 4216
prompt, - 4217
&response.text_content(), - 4218
self.config.outcome.as_ref(), - 4219
receipts, - 4220
verification_stale, - 4221
)?; - 4222
if !matches!(reason, BlockReason::UserCompletionRequired) && *blocks_left == 0 { - 4223
return None; - 4224
} - 4225
if matches!(reason, BlockReason::UserCompletionRequired) { - 4226
return Some(reason.message()); - 4227
} - 4228
*blocks_left -= 1; - 4229
Some(reason.message()) - 4230
} - 4231
- 4232
/// Appends the continue nudge (model-visible => logged) and reports - 4233
/// whether the loop may continue within max_turns. - 4234
async fn guard_continue( - 4235
&mut self, - 4236
reason: String, - 4237
events: &mpsc::Sender<AgentEvent>, - 4238
turn: usize, - 4239
) -> bool { - 4240
if turn + 1 >= self.config.max_turns { - 4241
return false; - 4242
} - 4243
let _ = events - 4244
.send(AgentEvent::StopHookContinuation { - 4245
reason: reason.clone(), - 4246
}) - 4247
.await; - 4248
let _ = self - 4249
.session - 4250
.lock() - 4251
.await - 4252
.append_message(MessageRecord::control( - 4253
vak_intent::control::ControlKind::StopGuard, - 4254
format!("[stop-guard]: {reason}\nPlease continue."), - 4255
)); - 4256
let _ = events.send(AgentEvent::DraftDiscarded { turn }).await; - 4257
true - 4258
} - 4259
- 4260
/// Recovers the session ledger after a run (server/API consumers). - 4261
#[allow(clippy::panic)] - 4262
pub async fn into_session(self) -> SessionLog { - 4263
let mut config = self.config; - 4264
config.tools.clear(); - 4265
config.flow_dispatcher = None; - 4266
drop(config); - 4267
match Arc::try_unwrap(self.session) { - 4268
Ok(session) => session.into_inner(), - 4269
Err(_) => { - 4270
panic!("managed flow dispatcher retained the session after the agent stopped") - 4271
} - 4272
} - 4273
} - 4274
- 4275
/// The path a successful call of `name` with `input` delivers for - 4276
/// review, when `name` is one of this agent's tools - 4277
/// (`Tool::delivered_file`). - 4278
fn delivered_file(&self, name: &str, input: &Value) -> Option<String> { - 4279
self.config - 4280
.tools - 4281
.iter() - 4282
.find(|tool| tool.name() == name) - 4283
.and_then(|tool| tool.delivered_file(input)) - 4284
} - 4285
- 4286
/// Whether `name` is one of this agent's tools and it declares that a - 4287
/// successful call presents a card (`Tool::presents_cards`). - 4288
fn tool_presents_cards(&self, name: &str) -> bool { - 4289
self.config - 4290
.tools - 4291
.iter() - 4292
.any(|tool| tool.name() == name && tool.presents_cards()) - 4293
} - 4294
- 4295
/// Appends the model's response and returns its ledger entry id, so - 4296
/// callers that need to attribute something back to this exact message - 4297
/// (a fence-path `Presentation`, docs/design/68-context-engine.md §10) - 4298
/// don't have to re-derive it. `None` only on a session write failure.
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.