- 18650
StatusCode::CONFLICT, - 18651
format!("{command} is not available in the preview environment"), - 18652
) - 18653
.into_response(); - 18654
} - 18655
let runtime_home = prepared.join(".vak-runtime-home"); - 18656
let _ = std::fs::create_dir_all(&runtime_home); - 18657
let environment = vec![ - 18658
( - 18659
"HOME".to_string(), - 18660
runtime_home.to_string_lossy().into_owned(), - 18661
), - 18662
( - 18663
"XDG_CACHE_HOME".to_string(), - 18664
runtime_home.join("cache").to_string_lossy().into_owned(), - 18665
), - 18666
( - 18667
"npm_config_cache".to_string(), - 18668
runtime_home.join("npm").to_string_lossy().into_owned(), - 18669
), - 18670
]; - 18671
let command_display = std::iter::once(command) - 18672
.chain(args.iter().map(String::as_str)) - 18673
.collect::<Vec<_>>() - 18674
.join(" "); - 18675
if let Err(error) = append_preview_preparation( - 18676
&state, - 18677
&saved, - 18678
vak_sandbox::EnvironmentState::Preparing, - 18679
&command_display, - 18680
"dependency preparation started", - 18681
) { - 18682
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18683
return (StatusCode::INTERNAL_SERVER_ERROR, error).into_response(); - 18684
} - 18685
let child = vak_tools::broker::spawn_persistent_worker( - 18686
&state.core.tool_worker_exe(), - 18687
&prepared, - 18688
command, - 18689
&args, - 18690
&environment, - 18691
preview_sandbox(&state.core).as_deref(), - 18692
) - 18693
.await; - 18694
let mut child = match child { - 18695
Ok(child) => child, - 18696
Err(error) => { - 18697
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18698
let _ = append_preview_preparation( - 18699
&state, - 18700
&saved, - 18701
vak_sandbox::EnvironmentState::Failed, - 18702
&command_display, - 18703
&error, - 18704
); - 18705
return (StatusCode::INTERNAL_SERVER_ERROR, error).into_response(); - 18706
} - 18707
}; - 18708
let stdout = child - 18709
.stdout - 18710
.take() - 18711
.map(|stream| tokio::spawn(read_capped_output(stream, 65_536))); - 18712
let stderr = child - 18713
.stderr - 18714
.take() - 18715
.map(|stream| tokio::spawn(read_capped_output(stream, 65_536))); - 18716
let status = tokio::time::timeout(std::time::Duration::from_secs(300), child.wait()).await; - 18717
let status = match status { - 18718
Ok(Ok(status)) => status, - 18719
_ => { - 18720
vak_tools::bash::kill_process_group(&child.id()); - 18721
let _ = child.wait().await; - 18722
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18723
let _ = append_preview_preparation( - 18724
&state, - 18725
&saved, - 18726
vak_sandbox::EnvironmentState::Failed, - 18727
&command_display, - 18728
"dependency preparation timed out", - 18729
); - 18730
return ( - 18731
StatusCode::GATEWAY_TIMEOUT, - 18732
"dependency preparation timed out", - 18733
) - 18734
.into_response(); - 18735
} - 18736
}; - 18737
let stdout = match stdout { - 18738
Some(task) => task.await.unwrap_or_default(), - 18739
None => Vec::new(), - 18740
}; - 18741
let stderr = match stderr { - 18742
Some(task) => task.await.unwrap_or_default(), - 18743
None => Vec::new(), - 18744
}; - 18745
let evidence = format!( - 18746
"{}{}", - 18747
String::from_utf8_lossy(&stdout), - 18748
String::from_utf8_lossy(&stderr) - 18749
); - 18750
if !status.success() || verify_launch_tree(&saved.candidate, &prepared).is_err() { - 18751
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18752
let _ = append_preview_preparation( - 18753
&state, - 18754
&saved, - 18755
vak_sandbox::EnvironmentState::Failed, - 18756
&command_display, - 18757
&evidence, - 18758
); - 18759
return (StatusCode::BAD_GATEWAY, Json(serde_json::json!({ "error": "dependency preparation failed", "evidence": evidence }))).into_response(); - 18760
} - 18761
if let Err(error) = std::fs::write( - 18762
prepared.join(".vak-candidate-digest"), - 18763
&saved.candidate_digest, - 18764
) { - 18765
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18766
let _ = append_preview_preparation( - 18767
&state, - 18768
&saved, - 18769
vak_sandbox::EnvironmentState::Failed, - 18770
&command_display, - 18771
format!("prepared runtime marker could not be written: {error}"), - 18772
); - 18773
return (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(); - 18774
} - 18775
if let Err(error) = append_preview_preparation( - 18776
&state, - 18777
&saved, - 18778
vak_sandbox::EnvironmentState::Ready, - 18779
&command_display, - 18780
&evidence, - 18781
) { - 18782
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18783
return (StatusCode::INTERNAL_SERVER_ERROR, error).into_response(); - 18784
} - 18785
Json( - 18786
serde_json::json!({ "prepared": true, "candidate_id": candidate_id, "evidence": evidence }), - 18787
) - 18788
.into_response() - 18789
} - 18790
- 18791
async fn start_launch( - 18792
State(state): State<AppState>, - 18793
Path(id): Path<String>, - 18794
Json(body): Json<LaunchNameBody>, - 18795
) -> axum::response::Response { - 18796
use axum::response::IntoResponse; - 18797
- 18798
let Some(workspace) = sandbox_session_workspace(&state, &id) else { - 18799
return StatusCode::NOT_FOUND.into_response(); - 18800
}; - 18801
let launch_root = match launch_root(&state, &id, body.candidate_id.as_deref(), &workspace) { - 18802
Ok(root) => root, - 18803
Err(error) => { - 18804
return ( - 18805
StatusCode::CONFLICT, - 18806
Json(serde_json::json!({ "error": error })), - 18807
) - 18808
.into_response(); - 18809
} - 18810
}; - 18811
let mut servers = match parse_launch_toml(&launch_root) { - 18812
Ok(s) => s, - 18813
Err(e) => { - 18814
return ( - 18815
StatusCode::BAD_REQUEST, - 18816
Json(serde_json::json!({ "error": e })), - 18817
) - 18818
.into_response(); - 18819
} - 18820
}; - 18821
if servers.is_empty() { - 18822
servers = detect_launch(&launch_root); - 18823
} - 18824
let Some(cfg) = servers.iter().find(|s| s.name == body.name) else { - 18825
return ( - 18826
StatusCode::NOT_FOUND, - 18827
Json(serde_json::json!({ "error": "unknown server name" })), - 18828
) - 18829
.into_response(); - 18830
}; - 18831
if !vak_tools::bash::executable_available(&cfg.cmd, &launch_root) { - 18832
return ( - 18833
StatusCode::CONFLICT, - 18834
Json(serde_json::json!({ - 18835
"error": format!("{} is not available in the preview environment", cfg.cmd) - 18836
})), - 18837
) - 18838
.into_response(); - 18839
} - 18840
if matches!(cfg.cmd.as_str(), "npm" | "pnpm" | "yarn" | "bun") - 18841
&& javascript_dependencies_missing(&launch_root) - 18842
{ - 18843
return ( - 18844
StatusCode::CONFLICT, - 18845
Json(serde_json::json!({ - 18846
"error": "project dependencies have not been prepared for this saved version", - 18847
"availability": "needs_preparation" - 18848
})), - 18849
) - 18850
.into_response(); - 18851
} - 18852
if let Some(port) = cfg.port - 18853
&& !port_is_available(port) - 18854
{ - 18855
return ( - 18856
StatusCode::CONFLICT, - 18857
Json(serde_json::json!({ - 18858
"error": format!("port {port} is already in use") - 18859
})), - 18860
) - 18861
.into_response(); - 18862
} - 18863
- 18864
let key = proc_key(&id, body.candidate_id.as_deref(), &cfg.name); - 18865
{ - 18866
let procs = state - 18867
.procs - 18868
.lock() - 18869
.unwrap_or_else(std::sync::PoisonError::into_inner); - 18870
if procs.contains_key(&key) { - 18871
return ( - 18872
StatusCode::CONFLICT, - 18873
Json(serde_json::json!({ "error": "already running" })), - 18874
) - 18875
.into_response(); - 18876
} - 18877
} - 18878
- 18879
let mut child = match vak_tools::broker::spawn_persistent_worker( - 18880
&state.core.tool_worker_exe(), - 18881
&launch_root, - 18882
&cfg.cmd, - 18883
&cfg.args, - 18884
&[], - 18885
preview_sandbox(&state.core).as_deref(), - 18886
) - 18887
.await - 18888
{ - 18889
Ok(child) => child, - 18890
Err(error) => { - 18891
return ( - 18892
StatusCode::INTERNAL_SERVER_ERROR, - 18893
Json(serde_json::json!({ "error": error })), - 18894
) - 18895
.into_response(); - 18896
} - 18897
}; - 18898
- 18899
let logs: Arc<Mutex<std::collections::VecDeque<String>>> = - 18900
Arc::new(Mutex::new(std::collections::VecDeque::with_capacity(500))); - 18901
// Drain stdout+stderr into a bounded ring. - 18902
if let Some(out) = child.stdout.take() { - 18903
let logs_out = logs.clone(); - 18904
tokio::spawn(async move { - 18905
use tokio::io::AsyncReadExt; - 18906
let mut reader = out; - 18907
let mut buf = [0u8; 1024]; - 18908
let mut line = String::new(); - 18909
loop { - 18910
match reader.read(&mut buf).await { - 18911
Ok(0) | Err(_) => break, - 18912
Ok(n) => { - 18913
line.push_str(&String::from_utf8_lossy(&buf[..n])); - 18914
while let Some(pos) = line.find('\n') { - 18915
let l: String = line.drain(..=pos).collect(); - 18916
let mut g = logs_out - 18917
.lock() - 18918
.unwrap_or_else(std::sync::PoisonError::into_inner); - 18919
if g.len() >= 500 { - 18920
g.pop_front(); - 18921
} - 18922
g.push_back(l.trim_end().to_string()); - 18923
} - 18924
} - 18925
} - 18926
} - 18927
}); - 18928
} - 18929
if let Some(err) = child.stderr.take() { - 18930
let logs_err = logs.clone(); - 18931
tokio::spawn(async move { - 18932
use tokio::io::AsyncReadExt; - 18933
let mut reader = err; - 18934
let mut buf = [0u8; 1024]; - 18935
let mut line = String::new(); - 18936
loop { - 18937
match reader.read(&mut buf).await { - 18938
Ok(0) | Err(_) => break, - 18939
Ok(n) => { - 18940
line.push_str(&String::from_utf8_lossy(&buf[..n])); - 18941
while let Some(pos) = line.find('\n') { - 18942
let l: String = line.drain(..=pos).collect(); - 18943
let mut g = logs_err - 18944
.lock() - 18945
.unwrap_or_else(std::sync::PoisonError::into_inner); - 18946
if g.len() >= 500 { - 18947
g.pop_front(); - 18948
} - 18949
g.push_back(l.trim_end().to_string()); - 18950
} - 18951
} - 18952
} - 18953
} - 18954
}); - 18955
} - 18956
- 18957
state - 18958
.procs - 18959
.lock() - 18960
.unwrap_or_else(std::sync::PoisonError::into_inner) - 18961
.insert( - 18962
key.clone(), - 18963
ManagedProc { - 18964
child, - 18965
logs: logs.clone(), - 18966
}, - 18967
); - 18968
- 18969
// Give the server a moment to bind its port so the preview iframe works - 18970
// immediately after start. - 18971
let listening = match cfg.port { - 18972
Some(p) => wait_for_port(p, std::time::Duration::from_secs(15)).await, - 18973
None => { - 18974
tokio::time::sleep(std::time::Duration::from_millis(250)).await; - 18975
false - 18976
} - 18977
}; - 18978
let exited = { - 18979
let mut procs = state - 18980
.procs - 18981
.lock() - 18982
.unwrap_or_else(std::sync::PoisonError::into_inner); - 18983
let finished = procs - 18984
.get_mut(&key) - 18985
.and_then(|process| process.child.try_wait().ok().flatten()); - 18986
if finished.is_some() { - 18987
procs.remove(&key); - 18988
} - 18989
finished - 18990
}; - 18991
if let Some(status) = exited { - 18992
let output = logs - 18993
.lock() - 18994
.unwrap_or_else(std::sync::PoisonError::into_inner) - 18995
.iter() - 18996
.rev() - 18997
.take(8) - 18998
.cloned() - 18999
.collect::<Vec<_>>() - 19000
.into_iter() - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.