- 19001
.rev() - 19002
.collect::<Vec<_>>(); - 19003
return ( - 19004
StatusCode::BAD_GATEWAY, - 19005
Json(serde_json::json!({ - 19006
"error": format!("preview process exited before becoming ready ({status})"), - 19007
"lines": output, - 19008
})), - 19009
) - 19010
.into_response(); - 19011
} - 19012
- 19013
( - 19014
StatusCode::OK, - 19015
Json(serde_json::json!({ "started": true, "listening": listening })), - 19016
) - 19017
.into_response() - 19018
} - 19019
- 19020
async fn stop_launch( - 19021
State(state): State<AppState>, - 19022
Path(id): Path<String>, - 19023
Json(body): Json<LaunchNameBody>, - 19024
) -> StatusCode { - 19025
let removed = state - 19026
.procs - 19027
.lock() - 19028
.unwrap_or_else(std::sync::PoisonError::into_inner) - 19029
.remove(&proc_key(&id, body.candidate_id.as_deref(), &body.name)); - 19030
match removed { - 19031
Some(mut p) => { - 19032
vak_tools::bash::kill_process_group(&p.child.id()); - 19033
let _ = p.child.kill().await; - 19034
let _ = p.child.wait().await; - 19035
StatusCode::OK - 19036
} - 19037
None => StatusCode::NOT_FOUND, - 19038
} - 19039
} - 19040
- 19041
async fn launch_logs( - 19042
State(state): State<AppState>, - 19043
Path(id): Path<String>, - 19044
axum::extract::Query(q): axum::extract::Query<LaunchNameBody>, - 19045
) -> Json<serde_json::Value> { - 19046
let procs = state - 19047
.procs - 19048
.lock() - 19049
.unwrap_or_else(std::sync::PoisonError::into_inner); - 19050
match procs.get(&proc_key(&id, q.candidate_id.as_deref(), &q.name)) { - 19051
Some(p) => { - 19052
let lines: Vec<String> = p - 19053
.logs - 19054
.lock() - 19055
.unwrap_or_else(std::sync::PoisonError::into_inner) - 19056
.iter() - 19057
.cloned() - 19058
.collect(); - 19059
Json(serde_json::json!({ "lines": lines })) - 19060
} - 19061
None => Json(serde_json::json!({ "lines": [], "error": "not running" })), - 19062
} - 19063
} - 19064
- 19065
#[cfg(test)] - 19066
#[allow(clippy::unwrap_used, clippy::expect_used)] - 19067
mod scheduler_pure_tests { - 19068
use super::{TaskDef, cron_slot_missed, recover_interrupted_tasks, stdout_section}; - 19069
use chrono::TimeZone; - 19070
use chrono::Utc; - 19071
use std::collections::HashMap; - 19072
- 19073
fn local(y: i32, mo: u32, d: u32, h: u32, mi: u32) -> chrono::DateTime<chrono::Local> { - 19074
chrono::Local - 19075
.with_ymd_and_hms(y, mo, d, h, mi, 0) - 19076
.single() - 19077
.unwrap() - 19078
} - 19079
- 19080
fn utc(dt: chrono::DateTime<chrono::Local>) -> chrono::DateTime<Utc> { - 19081
dt.with_timezone(&Utc) - 19082
} - 19083
- 19084
#[test] - 19085
fn restart_recovery_marks_only_interrupted_tasks() { - 19086
let make = |id: &str, status: Option<&str>| TaskDef { - 19087
id: id.into(), - 19088
name: id.into(), - 19089
prompt: "check in".into(), - 19090
interval_secs: 3600, - 19091
enabled: true, - 19092
cwd: std::path::PathBuf::from("/tmp"), - 19093
created_at: Utc::now(), - 19094
last_run_at: None, - 19095
last_session_id: None, - 19096
last_summary: None, - 19097
last_result_id: None, - 19098
last_run_status: status.map(str::to_owned), - 19099
last_delivery_state: Some("pending".into()), - 19100
last_wt: None, - 19101
deliver_to: None, - 19102
schedule: None, - 19103
timezone: None, - 19104
due_at: None, - 19105
script: None, - 19106
model_pin: None, - 19107
agent_id: None, - 19108
agent_revision: None, - 19109
}; - 19110
let mut tasks = HashMap::from([ - 19111
("running".into(), make("running", Some("working"))), - 19112
("done".into(), make("done", Some("complete"))), - 19113
]); - 19114
assert!(recover_interrupted_tasks(&mut tasks)); - 19115
assert_eq!( - 19116
tasks["running"].last_run_status.as_deref(), - 19117
Some("interrupted") - 19118
); - 19119
assert_eq!(tasks["done"].last_run_status.as_deref(), Some("complete")); - 19120
} - 19121
- 19122
#[test] - 19123
fn missed_slot_matrix() { - 19124
let every_min = "* * * * *"; - 19125
// Ran at the current slot → its next slot is in the future. - 19126
assert!(!cron_slot_missed( - 19127
every_min, - 19128
utc(local(2026, 8, 24, 10, 30)), - 19129
local(2026, 8, 24, 10, 30), - 19130
)); - 19131
// Ran yesterday; today's slot already passed → missed. - 19132
assert!(cron_slot_missed( - 19133
"0 12 * * *", - 19134
utc(local(2026, 8, 23, 12, 0)), - 19135
local(2026, 8, 24, 13, 0), - 19136
)); - 19137
// Ran after the latest slot (manual run-now covers it) → not missed. - 19138
assert!(!cron_slot_missed( - 19139
"0 12 * * *", - 19140
utc(local(2026, 8, 24, 12, 30)), - 19141
local(2026, 8, 24, 13, 0), - 19142
)); - 19143
// The slot exactly one step after the last run is due right now. - 19144
assert!(cron_slot_missed( - 19145
"*/15 * * * *", - 19146
utc(local(2026, 8, 24, 10, 30)), - 19147
local(2026, 8, 24, 10, 45), - 19148
)); - 19149
// Bad expression never reports a miss (parked markers handle it). - 19150
assert!(!cron_slot_missed( - 19151
"99 * * * *", - 19152
utc(local(2026, 8, 23, 12, 0)), - 19153
local(2026, 8, 24, 13, 0), - 19154
)); - 19155
} - 19156
- 19157
#[test] - 19158
fn stdout_section_extracts_only_stdout() { - 19159
assert_eq!( - 19160
stdout_section("[stdout]\nhello\nworld\n\n[stderr]\noops\n"), - 19161
"hello\nworld\n" - 19162
); - 19163
assert_eq!( - 19164
stdout_section( - 19165
"[working directory: /tmp]\n[file: /tmp/res.html]\n[stdout]\nhello\nworld\n\n[stderr]\noops\n" - 19166
), - 19167
"hello\nworld\n" - 19168
); - 19169
assert_eq!(stdout_section("(no output)"), ""); - 19170
assert_eq!(stdout_section(""), ""); - 19171
} - 19172
} - 19173
- 19174
#[cfg(test)] - 19175
#[allow(clippy::unwrap_used, clippy::expect_used)] - 19176
mod configuration_control_tests { - 19177
use super::*; - 19178
- 19179
fn control_state(dir: &std::path::Path) -> AppState { - 19180
vak_config::paths::isolate_home_for_tests(); - 19181
let core = Core::new(dir.to_path_buf()).unwrap(); - 19182
core.set_sessions_home(dir.join("home")); - 19183
AppState::new(core) - 19184
} - 19185
- 19186
// ---- remembering an approval (finding 02) ------------------------------ - 19187
- 19188
/// Put a gate into a session's pending map the way `HttpApprover` does, - 19189
/// so the answer path can be exercised without a provider. - 19190
fn park_gate(handle: &Arc<SessionHandle>, tool: &str, args_json: &str) -> String { - 19191
let id = uuid::Uuid::now_v7().to_string(); - 19192
let (respond, _rx) = oneshot::channel(); - 19193
handle - 19194
.pending - 19195
.lock() - 19196
.unwrap_or_else(std::sync::PoisonError::into_inner) - 19197
.insert( - 19198
id.clone(), - 19199
ApprovalRequest { - 19200
id: id.clone(), - 19201
tool: tool.into(), - 19202
args_json: args_json.into(), - 19203
reason: "needs approval".into(), - 19204
requested_at: chrono::Utc::now(), - 19205
respond: Arc::new(Mutex::new(Some(respond))), - 19206
answered_by: Arc::new(Mutex::new(None)), - 19207
delegated_to: Arc::new(Mutex::new(None)), - 19208
}, - 19209
); - 19210
id - 19211
} - 19212
- 19213
async fn answer_json( - 19214
state: &AppState, - 19215
session: &str, - 19216
req: &str, - 19217
body: ApprovalBody, - 19218
) -> serde_json::Value { - 19219
let response = answer_approval( - 19220
State(state.clone()), - 19221
axum::extract::Path((session.to_string(), req.to_string())), - 19222
Json(body), - 19223
) - 19224
.await; - 19225
assert_eq!(response.status(), StatusCode::OK); - 19226
let bytes = axum::body::to_bytes(response.into_body(), 64 * 1024) - 19227
.await - 19228
.unwrap(); - 19229
serde_json::from_slice(&bytes).unwrap() - 19230
} - 19231
- 19232
/// The mechanism `08-permissions.md` has described since the engine - 19233
/// shipped, and which had no caller on any surface until now. - 19234
#[tokio::test] - 19235
async fn remembering_an_approval_writes_a_scoped_rule_that_applies_at_once() { - 19236
vak_config::paths::isolate_home_for_tests(); - 19237
let dir = tempfile::tempdir().unwrap(); - 19238
let core = Core::new_with_trust(dir.path().to_path_buf(), true).unwrap(); - 19239
core.set_sessions_home(dir.path().join("home")); - 19240
let state = AppState::new(core.clone()); - 19241
let session = core.start_session().await.unwrap(); - 19242
let id = session.header().unwrap().session_id.clone(); - 19243
let handle = register_handle( - 19244
&state, - 19245
id.clone(), - 19246
session, - 19247
core.cwd().clone(), - 19248
core.clone(), - 19249
); - 19250
- 19251
let req = park_gate(&handle, "bash", r#"{"command":"cargo test --lib"}"#); - 19252
let json = answer_json( - 19253
&state, - 19254
&id, - 19255
&req, - 19256
ApprovalBody { - 19257
approve: true, - 19258
remember: true, - 19259
}, - 19260
) - 19261
.await; - 19262
assert_eq!(json["approved"], true); - 19263
assert_eq!(json["learned_rule"], "+bash(cargo *)"); - 19264
assert!(json["learn_error"].is_null(), "{json}"); - 19265
- 19266
// The next engine build sees it, with no restart. - 19267
let engine = core - 19268
.build_permission_engine(&core.extra_allow_snapshot()) - 19269
.unwrap(); - 19270
assert!(matches!( - 19271
engine.evaluate( - 19272
"bash", - 19273
&serde_json::json!({ "command": "cargo build" }), - 19274
vak_permission::Mode::WorkspaceWrite, - 19275
core.cwd() - 19276
), - 19277
vak_permission::Decision::Allow - 19278
)); - 19279
} - 19280
- 19281
/// A call that cannot be narrowed safely is still approved — the run is - 19282
/// waiting on it — and simply not remembered, with the reason reported. - 19283
#[tokio::test] - 19284
async fn a_call_that_cannot_be_narrowed_is_approved_but_not_remembered() { - 19285
vak_config::paths::isolate_home_for_tests(); - 19286
let dir = tempfile::tempdir().unwrap(); - 19287
let core = Core::new_with_trust(dir.path().to_path_buf(), true).unwrap(); - 19288
core.set_sessions_home(dir.path().join("home")); - 19289
let state = AppState::new(core.clone()); - 19290
let session = core.start_session().await.unwrap(); - 19291
let id = session.header().unwrap().session_id.clone(); - 19292
let handle = register_handle( - 19293
&state, - 19294
id.clone(), - 19295
session, - 19296
core.cwd().clone(), - 19297
core.clone(), - 19298
); - 19299
- 19300
let req = park_gate(&handle, "bash", r#"{"command":"echo $(whoami)"}"#); - 19301
let json = answer_json( - 19302
&state, - 19303
&id, - 19304
&req, - 19305
ApprovalBody { - 19306
approve: true, - 19307
remember: true, - 19308
}, - 19309
) - 19310
.await; - 19311
assert_eq!(json["approved"], true, "the gate is still answered"); - 19312
assert!(json["learned_rule"].is_null()); - 19313
assert!( - 19314
json["learn_error"] - 19315
.as_str() - 19316
.unwrap() - 19317
.contains("cannot be narrowed"), - 19318
"{json}" - 19319
); - 19320
assert!(core.extra_allow_snapshot().is_empty()); - 19321
} - 19322
- 19323
/// Remembering a refusal would be a deny rule, which is a different and - 19324
/// much heavier decision than answering one gate. - 19325
#[tokio::test] - 19326
async fn a_refusal_is_never_remembered() { - 19327
vak_config::paths::isolate_home_for_tests(); - 19328
let dir = tempfile::tempdir().unwrap(); - 19329
let core = Core::new_with_trust(dir.path().to_path_buf(), true).unwrap(); - 19330
core.set_sessions_home(dir.path().join("home")); - 19331
let state = AppState::new(core.clone()); - 19332
let session = core.start_session().await.unwrap(); - 19333
let id = session.header().unwrap().session_id.clone(); - 19334
let handle = register_handle( - 19335
&state, - 19336
id.clone(), - 19337
session, - 19338
core.cwd().clone(), - 19339
core.clone(), - 19340
); - 19341
- 19342
let req = park_gate(&handle, "bash", r#"{"command":"rm -rf /"}"#); - 19343
let json = answer_json( - 19344
&state, - 19345
&id, - 19346
&req, - 19347
ApprovalBody { - 19348
approve: false, - 19349
remember: true, - 19350
}, - 19351
) - 19352
.await; - 19353
assert_eq!(json["approved"], false); - 19354
assert!(json["learned_rule"].is_null()); - 19355
assert!(core.extra_allow_snapshot().is_empty()); - 19356
} - 19357
- 19358
// ---- unattended runs (findings 08 and 09) ------------------------------ - 19359
- 19360
/// A gate raised where nobody is subscribed used to emit an SSE event - 19361
/// into the void and then block on `rx.await` forever, holding the - 19362
/// session handle open until the process restarted. - 19363
#[tokio::test] - 19364
async fn an_unattended_http_approver_refuses_instead_of_waiting() { - 19365
let events_tx = events::EventBus::new(); - 19366
let approver = HttpApprover { - 19367
events_tx, - 19368
pending: Arc::new(Mutex::new(HashMap::new())), - 19369
session_id: "s".into(), - 19370
activity_buffer: Arc::new(Mutex::new(Vec::new())), - 19371
answerable: false, - 19372
}; - 19373
assert!(!Approver::answerable(&approver)); - 19374
// Returns immediately; without the guard this would block until the - 19375
// 15-minute deadline, which the test would never reach. - 19376
assert!(!approver.approve("bash", "{}", "needs approval").await); - 19377
} - 19378
- 19379
#[tokio::test] - 19380
async fn resolved_approval_activity_names_the_verified_decision_maker() { - 19381
let pending = Arc::new(Mutex::new(HashMap::new())); - 19382
let activity_buffer = Arc::new(Mutex::new(Vec::new())); - 19383
let approver = HttpApprover { - 19384
events_tx: events::EventBus::new(), - 19385
pending: pending.clone(), - 19386
session_id: "session-approval".into(), - 19387
activity_buffer: activity_buffer.clone(), - 19388
answerable: true, - 19389
}; - 19390
let task = tokio::spawn(async move { approver.approve("write", "{}", "save draft").await }); - 19391
let request = loop { - 19392
if let Some(request) = pending.lock().unwrap().values().next().cloned() { - 19393
break request; - 19394
} - 19395
tokio::task::yield_now().await; - 19396
}; - 19397
*request.answered_by.lock().unwrap() = Some(("person-1".into(), "Asha".into())); - 19398
request.respond(true); - 19399
assert!(task.await.unwrap()); - 19400
let activities = activity_buffer.lock().unwrap(); - 19401
assert_eq!(activities.len(), 2); - 19402
assert!(!activities[0].data.contains_key("actor_id")); - 19403
assert_eq!( - 19404
activities[1].data.get("actor_id").map(String::as_str), - 19405
Some("person-1") - 19406
); - 19407
assert_eq!( - 19408
activities[1].data.get("actor_name").map(String::as_str), - 19409
Some("Asha") - 19410
); - 19411
} - 19412
- 19413
/// `Core::approver_answerable` is stamped before a run and the approver - 19414
/// is installed at dispatch. They used to be independent, with a comment - 19415
/// asking hosts to keep them in step; the scheduler did not. A - 19416
/// disagreement is now corrected in favour of the approver and recorded. - 19417
#[tokio::test] - 19418
async fn a_stamped_answerability_loses_to_the_installed_approver() { - 19419
let dir = tempfile::tempdir().unwrap(); - 19420
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19421
core.set_sessions_home(dir.path().join("home")); - 19422
- 19423
// Default is "attended"; AutoDeny says otherwise. - 19424
assert!(core.approver_answerable()); - 19425
let corrected = core.clone().with_approver(&vak_agent::AutoDeny); - 19426
assert!(!corrected.approver_answerable()); - 19427
- 19428
// And the other direction, for a surface that stamped false. - 19429
let stamped = core.clone().with_approver_answerable(false); - 19430
assert!( - 19431
stamped - 19432
.with_approver(&vak_agent::AutoApprove) - 19433
.approver_answerable() - 19434
); - 19435
} - 19436
- 19437
// ---- gateway approval policy (finding 01) ------------------------------ - 19438
- 19439
/// The setting was readable on three screens and writable nowhere, which - 19440
/// is why every `capability_unreachable` in the audit log had a remedy - 19441
/// no surface could perform. - 19442
#[tokio::test] - 19443
async fn forwarding_can_be_turned_on_and_survives_a_reload() { - 19444
let dir = tempfile::tempdir().unwrap(); - 19445
let state = control_state(dir.path()); - 19446
assert_eq!(state.gateway.approvals_mode(), "deny"); - 19447
- 19448
let response = put_gateway_approvals( - 19449
State(state.clone()), - 19450
Json(GatewayApprovalsBody { - 19451
mode: "forward".into(), - 19452
approver: Some("telegram:12345".into()), - 19453
timeout_secs: Some(60), - 19454
scope: None, - 19455
}), - 19456
) - 19457
.await; - 19458
assert_eq!(response.status(), StatusCode::OK); - 19459
- 19460
// Live, without a restart. - 19461
assert_eq!(state.gateway.approvals_mode(), "forward"); - 19462
assert_eq!( - 19463
state.gateway.approver_target().as_deref(), - 19464
Some("telegram:12345") - 19465
); - 19466
assert_eq!(state.gateway.approval_timeout().as_secs(), 60); - 19467
- 19468
// And on disk, so the next process starts the same way. - 19469
let reloaded = vak_config::load_with_trust(dir.path(), true).unwrap(); - 19470
assert_eq!(reloaded.gateway.approvals, "forward"); - 19471
assert_eq!(reloaded.gateway.approver.as_deref(), Some("telegram:12345")); - 19472
} - 19473
- 19474
/// The loader degrades an unbacked `forward` to `deny` with a warning, - 19475
/// which is right for a bad file and wrong for a button press: the - 19476
/// operator would see success and get the opposite setting. - 19477
#[tokio::test] - 19478
async fn forwarding_without_a_chat_is_refused_rather_than_silently_denied() { - 19479
let dir = tempfile::tempdir().unwrap(); - 19480
let state = control_state(dir.path()); - 19481
for approver in [None, Some("not-a-chat-address".to_string())] { - 19482
let response = put_gateway_approvals( - 19483
State(state.clone()), - 19484
Json(GatewayApprovalsBody { - 19485
mode: "forward".into(), - 19486
approver, - 19487
timeout_secs: None, - 19488
scope: None, - 19489
}), - 19490
) - 19491
.await; - 19492
assert_eq!(response.status(), StatusCode::BAD_REQUEST); - 19493
} - 19494
assert_eq!(state.gateway.approvals_mode(), "deny", "nothing changed"); - 19495
} - 19496
- 19497
/// Going back to `deny` must not leave the old target behind for a later - 19498
/// `forward` to pick up silently. - 19499
#[tokio::test] - 19500
async fn returning_to_deny_clears_the_approver() { - 19501
let dir = tempfile::tempdir().unwrap(); - 19502
let state = control_state(dir.path()); - 19503
let ok = |body| put_gateway_approvals(State(state.clone()), Json(body)); - 19504
assert_eq!( - 19505
ok(GatewayApprovalsBody { - 19506
mode: "forward".into(), - 19507
approver: Some("telegram:1".into()), - 19508
timeout_secs: None, - 19509
scope: None, - 19510
}) - 19511
.await - 19512
.status(), - 19513
StatusCode::OK - 19514
); - 19515
assert_eq!( - 19516
ok(GatewayApprovalsBody { - 19517
mode: "deny".into(), - 19518
approver: None, - 19519
timeout_secs: None, - 19520
scope: None, - 19521
}) - 19522
.await - 19523
.status(), - 19524
StatusCode::OK - 19525
); - 19526
assert!(state.gateway.approver_target().is_none()); - 19527
let reloaded = vak_config::load_with_trust(dir.path(), true).unwrap(); - 19528
assert!(reloaded.gateway.approver.is_none()); - 19529
} - 19530
- 19531
// ---- permission rules (finding 02) ------------------------------------- - 19532
- 19533
#[tokio::test] - 19534
async fn rules_are_written_validated_and_applied_to_the_next_engine() { - 19535
let dir = tempfile::tempdir().unwrap(); - 19536
let state = control_state(dir.path()); - 19537
let args = serde_json::json!({ "command": "rm -rf /" }); - 19538
- 19539
let response = put_permission_rules( - 19540
State(state.clone()), - 19541
Json(PermissionRulesBody { - 19542
allow: None, - 19543
ask: None, - 19544
deny: Some(vec!["Bash(rm *)".into()]), - 19545
scope: None, - 19546
agent: None, - 19547
}), - 19548
) - 19549
.await; - 19550
assert_eq!(response.status(), StatusCode::OK); - 19551
- 19552
let engine = state.core.build_permission_engine(&[]).unwrap(); - 19553
assert!(matches!( - 19554
engine.evaluate( - 19555
"bash", - 19556
&args, - 19557
vak_permission::Mode::FullAccess, - 19558
state.core.cwd() - 19559
), - 19560
vak_permission::Decision::Deny { .. } - 19561
)); - 19562
} - 19563
- 19564
/// A half-applied rule set is a permission decision nobody chose, so one - 19565
/// bad spec rejects the whole request and writes nothing. - 19566
#[tokio::test] - 19567
async fn one_malformed_rule_rejects_the_whole_write() { - 19568
let dir = tempfile::tempdir().unwrap(); - 19569
let state = control_state(dir.path()); - 19570
let response = put_permission_rules( - 19571
State(state.clone()), - 19572
Json(PermissionRulesBody { - 19573
allow: None, - 19574
ask: None, - 19575
deny: Some(vec!["Bash(git *)".into(), "Bash((((".into()]), - 19576
scope: None, - 19577
agent: None, - 19578
}), - 19579
) - 19580
.await; - 19581
assert_eq!(response.status(), StatusCode::BAD_REQUEST); - 19582
let (_, _, deny) = state.core.effective_permission_rules(); - 19583
assert!(deny.is_empty(), "nothing may be written: {deny:?}"); - 19584
} - 19585
- 19586
// ---- global writes shadowed by a project pin (finding 07) -------------- - 19587
// - 19588
// Those tests write the Shared layer, which every test in this binary - 19589
// reads through `Core::new`, so they run in a binary of their own: - 19590
// `tests/shared_config_layer.rs`. No test here may write that layer. - 19591
- 19592
#[tokio::test] - 19593
async fn cross_process_mode_refresh_revokes_live_capability_before_apply() { - 19594
vak_config::paths::isolate_home_for_tests(); - 19595
let dir = tempfile::tempdir().unwrap(); - 19596
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19597
core.set_sessions_home(dir.path().join("home")); - 19598
let state = AppState::new(core.clone()); - 19599
let session = core.start_session().await.unwrap(); - 19600
let id = session.header().unwrap().session_id.clone(); - 19601
let handle = register_handle(&state, id, session, core.cwd().clone(), core.clone()); - 19602
assert!(!handle.cancel.lock().unwrap().is_cancelled()); - 19603
- 19604
vak_config::persist_project_preferences( - 19605
dir.path(), - 19606
None, - 19607
None, - 19608
Some(19), - 19609
Some(vak_config::PermissionMode::ReadOnly), - 19610
None, - 19611
None, - 19612
) - 19613
.unwrap(); - 19614
refresh_control_plane(&state); - 19615
- 19616
assert_eq!(core.effective_max_turns(), 19); - 19617
assert_eq!( - 19618
core.effective_permission_mode(), - 19619
vak_config::PermissionMode::ReadOnly - 19620
); - 19621
assert!(handle.cancel.lock().unwrap().is_cancelled()); - 19622
} - 19623
- 19624
/// The bug this locks in: before `effective_memory_*` existed, - 19625
/// `Core::config().memory.*` was read directly at every call site, so - 19626
/// a live PATCH — or another process persisting a change to disk — - 19627
/// silently did nothing until the process restarted. Mirrors - 19628
/// `cross_process_mode_refresh_revokes_live_capability_before_apply`'s - 19629
/// shape for the memory tier instead of permission mode. - 19630
#[tokio::test] - 19631
async fn cross_process_memory_refresh_takes_effect_without_restart() { - 19632
let dir = tempfile::tempdir().unwrap(); - 19633
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19634
core.set_sessions_home(dir.path().join("home")); - 19635
assert!(core.effective_memory_search_enabled(), "default is on"); - 19636
- 19637
vak_config::persist_project_memory_prefs(dir.path(), Some(false), None, None, None) - 19638
.unwrap(); - 19639
core.refresh_persisted_preferences().unwrap(); - 19640
- 19641
assert!( - 19642
!core.effective_memory_search_enabled(), - 19643
"a disk change from another process must reach an already-running Core" - 19644
); - 19645
// Untouched flags keep their default, proving the write was - 19646
// scoped to exactly the one field this call named. - 19647
assert!(core.effective_memory_write_enabled()); - 19648
} - 19649
- 19650
/// `PATCH /config` end to end: persists to disk, applies live - 19651
/// immediately (no restart), and a field the PATCH didn't mention - 19652
/// keeps its current value rather than reverting to whatever was on - 19653
/// disk before this call. - 19654
#[tokio::test] - 19655
async fn patch_config_memory_flags_apply_live_and_persist() { - 19656
let dir = tempfile::tempdir().unwrap(); - 19657
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19658
core.set_sessions_home(dir.path().join("home")); - 19659
let state = AppState::new(core.clone()); - 19660
- 19661
let response = patch_config( - 19662
State(state.clone()), - 19663
Json(ConfigPatch { - 19664
memory_write_enabled: Some(false), - 19665
..Default::default() - 19666
}), - 19667
) - 19668
.await; - 19669
assert_eq!(response.status(), StatusCode::OK); - 19670
- 19671
assert!( - 19672
!core.effective_memory_write_enabled(), - 19673
"must apply live without a restart" - 19674
); - 19675
assert!( - 19676
core.effective_memory_search_enabled(), - 19677
"a field this PATCH never mentioned must keep its value" - 19678
); - 19679
- 19680
// Persisted to disk, not just the in-process override — a fresh - 19681
// Core over the same cwd sees it too. - 19682
let fresh = Core::new(dir.path().to_path_buf()).unwrap(); - 19683
fresh.set_sessions_home(dir.path().join("home")); - 19684
assert!(!fresh.effective_memory_write_enabled()); - 19685
} - 19686
- 19687
/// Same live-without-restart guarantee as memory, for the `workers` - 19688
/// toggle newly surfaced in the admin console's Settings page — it was - 19689
/// previously read from `Core::config()` directly at both call sites, - 19690
/// so a PATCH would have silently done nothing. - 19691
#[tokio::test] - 19692
async fn patch_config_workers_applies_live_and_persists() { - 19693
let dir = tempfile::tempdir().unwrap(); - 19694
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19695
core.set_sessions_home(dir.path().join("home")); - 19696
let state = AppState::new(core.clone()); - 19697
assert!(core.effective_workers(), "default is on"); - 19698
- 19699
let response = patch_config( - 19700
State(state.clone()), - 19701
Json(ConfigPatch { - 19702
workers: Some(false), - 19703
..Default::default() - 19704
}), - 19705
) - 19706
.await; - 19707
assert_eq!(response.status(), StatusCode::OK); - 19708
assert!( - 19709
!core.effective_workers(), - 19710
"must apply live without a restart" - 19711
); - 19712
- 19713
let fresh = Core::new(dir.path().to_path_buf()).unwrap(); - 19714
fresh.set_sessions_home(dir.path().join("home")); - 19715
assert!(!fresh.effective_workers(), "must be persisted to disk too"); - 19716
} - 19717
- 19718
/// `PATCH /finops` sets a cap live and persists it; an explicit `null` - 19719
/// clears a previously-set cap rather than being indistinguishable - 19720
/// from the field being absent (the exact bug `deserialize_present` - 19721
/// exists to prevent, exercised here for a fresh field). - 19722
#[tokio::test] - 19723
async fn patch_finops_sets_and_clears_caps_live_and_persisted() { - 19724
let dir = tempfile::tempdir().unwrap(); - 19725
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19726
core.set_sessions_home(dir.path().join("home")); - 19727
let state = AppState::new(core.clone()); - 19728
assert_eq!(core.effective_finops_max_run_usd(), None); - 19729
- 19730
let status = patch_finops( - 19731
State(state.clone()), - 19732
Json(FinopsPatch { - 19733
max_run_usd: Some(Some(5.0)), - 19734
max_day_usd: None, - 19735
agent: None, - 19736
}), - 19737
) - 19738
.await; - 19739
assert_eq!(status.status(), StatusCode::OK); - 19740
assert_eq!(core.effective_finops_max_run_usd(), Some(5.0)); - 19741
- 19742
let fresh = Core::new(dir.path().to_path_buf()).unwrap(); - 19743
fresh.set_sessions_home(dir.path().join("home")); - 19744
assert_eq!( - 19745
fresh.effective_finops_max_run_usd(), - 19746
Some(5.0), - 19747
"must be persisted to disk too" - 19748
); - 19749
- 19750
// Explicit null clears it back to "no cap". - 19751
let status = patch_finops( - 19752
State(state.clone()), - 19753
Json(FinopsPatch { - 19754
max_run_usd: Some(None), - 19755
max_day_usd: None, - 19756
agent: None, - 19757
}), - 19758
) - 19759
.await; - 19760
assert_eq!(status.status(), StatusCode::OK); - 19761
assert_eq!( - 19762
core.effective_finops_max_run_usd(), - 19763
None, - 19764
"explicit null must clear the cap, not be a no-op" - 19765
); - 19766
} - 19767
- 19768
#[tokio::test] - 19769
async fn patch_finops_rejects_negative_cap() { - 19770
let dir = tempfile::tempdir().unwrap(); - 19771
let core = Core::new(dir.path().to_path_buf()).unwrap(); - 19772
core.set_sessions_home(dir.path().join("home")); - 19773
let state = AppState::new(core.clone()); - 19774
let status = patch_finops( - 19775
State(state), - 19776
Json(FinopsPatch { - 19777
max_run_usd: Some(Some(-1.0)), - 19778
max_day_usd: None, - 19779
agent: None, - 19780
}), - 19781
) - 19782
.await; - 19783
assert_eq!(status.status(), StatusCode::BAD_REQUEST); - 19784
} - 19785
} - 19786
- 19787
#[cfg(test)] - 19788
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 19789
mod sandbox_promotion_tests { - 19790
use super::*; - 19791
use std::collections::VecDeque; - 19792
use vak_llm::stream; - 19793
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; - 19794
use vak_llm::{EventStream, LlmError}; - 19795
- 19796
struct RevisionProvider { - 19797
replies: Mutex<VecDeque<AssistantMessage>>, - 19798
} - 19799
- 19800
#[async_trait::async_trait] - 19801
impl Provider for RevisionProvider { - 19802
fn name(&self) -> &str { - 19803
"revision-test" - 19804
} - 19805
- 19806
async fn stream( - 19807
&self, - 19808
_request: ChatRequest, - 19809
_cancel: CancellationToken, - 19810
) -> Result<EventStream, LlmError> { - 19811
let reply = self - 19812
.replies - 19813
.lock() - 19814
.ok() - 19815
.and_then(|mut replies| replies.pop_front()); - 19816
let (mut sink, rx) = stream::channel(64); - 19817
match reply { - 19818
Some(message) => { - 19819
sink.push(stream::StreamEvent::Start { - 19820
partial: message.clone(), - 19821
}); - 19822
sink.close_message(message).await; - 19823
} - 19824
None => { - 19825
sink.close_error(LlmError::Parse("revision test replies exhausted".into())) - 19826
.await - 19827
} - 19828
} - 19829
Ok(rx) - 19830
} - 19831
} - 19832
- 19833
fn revision_message(content: Vec<ContentBlock>, stop_reason: StopReason) -> AssistantMessage { - 19834
AssistantMessage { - 19835
content, - 19836
stop_reason, - 19837
usage: Usage::default(), - 19838
model: "test-model".into(), - 19839
response_id: None, - 19840
} - 19841
} - 19842
- 19843
fn seed_bound_result( - 19844
core: &Core, - 19845
session_id: &str, - 19846
execution_id: &str, - 19847
) -> vak_session::SessionLog { - 19848
let path = core - 19849
.sessions_home() - 19850
.join("sessions") - 19851
.join(vak_core::memory::hash_cwd(core.cwd())) - 19852
.join(format!("{session_id}.jsonl")); - 19853
let mut log = vak_session::SessionLog::create( - 19854
path, - 19855
vak_session::types::SessionHeader { - 19856
agent: Some(vak_core::vak_agent_identity()), - 19857
session_id: session_id.into(), - 19858
created_at: chrono::Utc::now(), - 19859
cwd: core.cwd().to_path_buf(), - 19860
parent_session_id: None, - 19861
contract_id: None, - 19862
work_item_id: None, - 19863
conversation: Some(vak_session::types::ConversationContext::local( - 19864
session_id, "test", - 19865
)), - 19866
contract: vak_session::types::FrozenContract { - 19867
app_version: "test".into(), - 19868
provider: "test".into(), - 19869
model: "test".into(), - 19870
route_ladder: Vec::new(), - 19871
route_objective: String::new(), - 19872
route_annotations: Vec::new(), - 19873
system_prompt: String::new(), - 19874
permission_mode: "workspace-write".into(), - 19875
capabilities: Vec::new(), - 19876
prompt_layers: Vec::new(), - 19877
}, - 19878
}, - 19879
) - 19880
.unwrap(); - 19881
log.append_message(vak_session::types::MessageRecord { - 19882
message: vak_llm::Message::user_text("Create the result"), - 19883
meta: None, - 19884
}) - 19885
.unwrap(); - 19886
log.append_message(vak_session::types::MessageRecord { - 19887
message: vak_llm::Message::assistant(vec![ - 19888
vak_llm::ContentBlock::ToolUse { - 19889
id: execution_id.into(), - 19890
name: "bash".into(), - 19891
input: serde_json::json!({"command": "create result"}), - 19892
}, - 19893
vak_llm::ContentBlock::text("The result is ready."), - 19894
]), - 19895
meta: None, - 19896
}) - 19897
.unwrap(); - 19898
log - 19899
} - 19900
- 19901
/// Target verification runs in the broker worker, so a test that - 19902
/// freezes or promotes a candidate needs the real worker binary. - 19903
fn pin_test_tool_worker(core: &Core) { - 19904
let worker = std::env::current_exe() - 19905
.unwrap() - 19906
.parent() - 19907
.unwrap() - 19908
.parent() - 19909
.unwrap() - 19910
.join("vak-tool-worker"); - 19911
assert!( - 19912
worker.is_file(), - 19913
"build vak-tool-worker (cargo build -p vak-server --bins) to run candidate verification" - 19914
); - 19915
core.set_tool_worker_exe(worker); - 19916
} - 19917
- 19918
async fn export_candidate(state: &AppState) -> vak_sandbox::CandidateRecord { - 19919
pin_test_tool_worker(&state.core); - 19920
append_session_sandbox_event( - 19921
&state.core.sessions_home(), - 19922
"session-1", - 19923
&AgentEvent::Sandbox(vak_tools::SandboxEvent::ExecutionStarted { - 19924
execution_id: "exec-1".into(), - 19925
owner_session_id: Some("session-1".into()), - 19926
tool: "bash".into(), - 19927
code_preview: "create result".into(), - 19928
language: "bash".into(), - 19929
scratch_dir: ".vak/scratch/e1".into(), - 19930
}), - 19931
); - 19932
let response = export_sandbox_candidate( - 19933
State(state.clone()), - 19934
Path("session-1".into()), - 19935
Json(SandboxCandidateBody { - 19936
execution_id: "exec-1".into(), - 19937
source: ".vak/scratch/e1".into(), - 19938
destination: ".".into(), - 19939
}), - 19940
) - 19941
.await; - 19942
assert_eq!(response.status(), StatusCode::OK); - 19943
let record: vak_sandbox::DurableRecord = serde_json::from_slice( - 19944
&axum::body::to_bytes(response.into_body(), 64 * 1024) - 19945
.await - 19946
.unwrap(), - 19947
) - 19948
.unwrap(); - 19949
let vak_sandbox::DurableRecord::Candidate(record) = record else { - 19950
panic!("candidate response") - 19951
}; - 19952
record - 19953
} - 19954
- 19955
/// A ledger in which the Agent made `calls` (`office_apply` id and - 19956
/// arguments), each succeeding, then answered. - 19957
fn seed_office_calls(core: &Core, session_id: &str, calls: &[(&str, serde_json::Value)]) { - 19958
let mut log = seed_bound_result(core, session_id, "unused"); - 19959
log.append_message(vak_session::types::MessageRecord { - 19960
message: vak_llm::Message::user_text("Update the budget"), - 19961
meta: None, - 19962
}) - 19963
.unwrap(); - 19964
for (id, input) in calls { - 19965
log.append_message(vak_session::types::MessageRecord { - 19966
message: vak_llm::Message::assistant(vec![vak_llm::ContentBlock::ToolUse { - 19967
id: (*id).into(), - 19968
name: "office_apply".into(), - 19969
input: input.clone(), - 19970
}]), - 19971
meta: None, - 19972
}) - 19973
.unwrap(); - 19974
log.append_message(vak_session::types::MessageRecord { - 19975
message: vak_llm::Message { - 19976
role: vak_llm::Role::User, - 19977
content: vec![vak_llm::ContentBlock::ToolResult { - 19978
tool_use_id: (*id).into(), - 19979
content: "Draft written.".into(), - 19980
is_error: false, - 19981
}], - 19982
}, - 19983
meta: None, - 19984
}) - 19985
.unwrap(); - 19986
} - 19987
log.append_message(vak_session::types::MessageRecord { - 19988
message: vak_llm::Message::assistant(vec![vak_llm::ContentBlock::text( - 19989
"The draft is ready for review.", - 19990
)]), - 19991
meta: None, - 19992
}) - 19993
.unwrap(); - 19994
} - 19995
- 19996
async fn body_json(response: axum::response::Response) -> serde_json::Value { - 19997
serde_json::from_slice( - 19998
&axum::body::to_bytes(response.into_body(), 1024 * 1024) - 19999
.await - 20000
.unwrap(),
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.