- 18464
.ok_or_else(|| "saved draft is unavailable".to_string()) - 18465
} - 18466
- 18467
fn verify_launch_tree( - 18468
candidate: &vak_sandbox::CandidateManifest, - 18469
root: &std::path::Path, - 18470
) -> Result<(), String> { - 18471
for file in candidate - 18472
.files - 18473
.iter() - 18474
.filter(|file| file.operation == vak_sandbox::CandidateOperation::Upsert) - 18475
{ - 18476
let path = confined_path(root, &file.path) - 18477
.ok_or_else(|| format!("saved draft path is invalid: {}", file.path))?; - 18478
let bytes = std::fs::read(path) - 18479
.map_err(|_| format!("saved draft file is unavailable: {}", file.path))?; - 18480
if vak_sandbox::digest(&bytes) != file.candidate_hash { - 18481
return Err(format!("saved draft changed after review: {}", file.path)); - 18482
} - 18483
} - 18484
Ok(()) - 18485
} - 18486
- 18487
fn sandbox_previews_root(state: &AppState) -> std::path::PathBuf { - 18488
state.core.sessions_home().join("sandbox").join("previews") - 18489
} - 18490
- 18491
async fn wait_for_port(port: u16, timeout: std::time::Duration) -> bool { - 18492
let deadline = tokio::time::Instant::now() + timeout; - 18493
while tokio::time::Instant::now() < deadline { - 18494
if tokio::net::TcpStream::connect(("127.0.0.1", port)) - 18495
.await - 18496
.is_ok() - 18497
{ - 18498
return true; - 18499
} - 18500
tokio::time::sleep(std::time::Duration::from_millis(250)).await; - 18501
} - 18502
false - 18503
} - 18504
- 18505
fn port_is_available(port: u16) -> bool { - 18506
std::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, port)).is_ok() - 18507
} - 18508
- 18509
#[derive(serde::Deserialize)] - 18510
struct LaunchNameBody { - 18511
name: String, - 18512
#[serde(default)] - 18513
candidate_id: Option<String>, - 18514
} - 18515
- 18516
fn dependency_install_command( - 18517
root: &std::path::Path, - 18518
) -> Result<(&'static str, Vec<String>), String> { - 18519
let bytes = std::fs::read(root.join("package.json")) - 18520
.map_err(|_| "package.json is unavailable".to_string())?; - 18521
let package = serde_json::from_slice::<serde_json::Value>(&bytes) - 18522
.map_err(|error| format!("package.json is invalid: {error}"))?; - 18523
let manager = javascript_package_manager(&package, |name| root.join(name).is_file()); - 18524
let args = match manager { - 18525
"pnpm" => vec!["install", "--frozen-lockfile", "--ignore-scripts"], - 18526
"yarn" => vec!["install", "--immutable", "--mode=skip-build"], - 18527
"bun" => vec!["install", "--frozen-lockfile", "--ignore-scripts"], - 18528
_ if root.join("package-lock.json").is_file() => { - 18529
vec!["ci", "--ignore-scripts", "--no-audit", "--no-fund"] - 18530
} - 18531
_ => vec!["install", "--ignore-scripts", "--no-audit", "--no-fund"], - 18532
}; - 18533
Ok((manager, args.into_iter().map(str::to_string).collect())) - 18534
} - 18535
- 18536
async fn read_capped_output( - 18537
mut stream: impl tokio::io::AsyncRead + Unpin, - 18538
limit: usize, - 18539
) -> Vec<u8> { - 18540
use tokio::io::AsyncReadExt; - 18541
let mut captured = Vec::new(); - 18542
let mut buffer = [0_u8; 8192]; - 18543
while let Ok(read) = stream.read(&mut buffer).await { - 18544
if read == 0 { - 18545
break; - 18546
} - 18547
let remaining = limit.saturating_sub(captured.len()); - 18548
captured.extend_from_slice(&buffer[..read.min(remaining)]); - 18549
} - 18550
captured - 18551
} - 18552
- 18553
fn append_preview_preparation( - 18554
state: &AppState, - 18555
saved: &vak_sandbox::CandidateRecord, - 18556
status: vak_sandbox::EnvironmentState, - 18557
command: &str, - 18558
evidence: impl Into<String>, - 18559
) -> Result<(), String> { - 18560
let record = - 18561
vak_sandbox::DurableRecord::PreviewPreparation(vak_sandbox::PreviewPreparationRecord { - 18562
record_id: format!("preview-preparation-{}", uuid::Uuid::now_v7()), - 18563
session_id: saved.session_id.clone(), - 18564
result_id: saved.result_id.clone(), - 18565
candidate_id: saved.candidate.candidate_id.clone(), - 18566
candidate_digest: saved.candidate_digest.clone(), - 18567
environment_id: format!("preview:{}", saved.candidate.candidate_id), - 18568
state: status, - 18569
command: command.to_string(), - 18570
evidence: evidence.into(), - 18571
updated_at: chrono::Utc::now().to_rfc3339(), - 18572
}); - 18573
vak_sandbox::append_record(&sandbox_records_path(state), &record) - 18574
.map_err(|error| error.to_string()) - 18575
} - 18576
- 18577
async fn prepare_launch( - 18578
State(state): State<AppState>, - 18579
Path(id): Path<String>, - 18580
Json(body): Json<LaunchNameBody>, - 18581
) -> axum::response::Response { - 18582
use axum::response::IntoResponse; - 18583
let Some(candidate_id) = body.candidate_id.as_deref() else { - 18584
return (StatusCode::BAD_REQUEST, "saved draft id is required").into_response(); - 18585
}; - 18586
let saved = match saved_launch_candidate(&state, &id, candidate_id) { - 18587
Ok(saved) => saved, - 18588
Err(error) => return (StatusCode::NOT_FOUND, error).into_response(), - 18589
}; - 18590
let frozen = sandbox_candidates_root(&state).join(candidate_id); - 18591
if saved.candidate.source_root != frozen - 18592
|| verify_launch_tree(&saved.candidate, &frozen).is_err() - 18593
{ - 18594
return ( - 18595
StatusCode::CONFLICT, - 18596
"saved draft failed integrity verification", - 18597
) - 18598
.into_response(); - 18599
} - 18600
let prepared = sandbox_previews_root(&state).join(candidate_id); - 18601
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18602
if let Err(error) = std::fs::create_dir_all(&prepared) { - 18603
return (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(); - 18604
} - 18605
if saved.candidate.files.iter().any(|file| { - 18606
matches!( - 18607
file.path.as_str(), - 18608
".npmrc" | ".yarnrc" | ".yarnrc.yml" | "bunfig.toml" - 18609
) - 18610
}) { - 18611
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18612
return ( - 18613
StatusCode::BAD_REQUEST, - 18614
"package manager credential/config files cannot enter a prepared preview", - 18615
) - 18616
.into_response(); - 18617
} - 18618
for file in saved - 18619
.candidate - 18620
.files - 18621
.iter() - 18622
.filter(|file| file.operation == vak_sandbox::CandidateOperation::Upsert) - 18623
{ - 18624
let Some(source) = confined_path(&frozen, &file.path) else { - 18625
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18626
return (StatusCode::CONFLICT, "saved draft path is invalid").into_response(); - 18627
}; - 18628
let target = prepared.join(&file.path); - 18629
if let Some(parent) = target.parent() - 18630
&& let Err(error) = std::fs::create_dir_all(parent) - 18631
{ - 18632
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18633
return (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(); - 18634
} - 18635
if let Err(error) = std::fs::copy(source, target) { - 18636
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18637
return (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(); - 18638
} - 18639
} - 18640
let (command, args) = match dependency_install_command(&prepared) { - 18641
Ok(value) => value, - 18642
Err(error) => { - 18643
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18644
return (StatusCode::BAD_REQUEST, error).into_response(); - 18645
} - 18646
}; - 18647
if !vak_tools::bash::executable_available(command, &prepared) { - 18648
let _ = vak_sandbox::remove_frozen_candidate(&prepared); - 18649
return ( - 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(),
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.