- 17650
if let Some(task) = map.get_mut(&tid) { - 17651
task.last_run_status = - 17652
Some(if is_error { "failed" } else { "complete" }.into()); - 17653
task.last_delivery_state = Some(delivery_state.into()); - 17654
} - 17655
}); - 17656
break; - 17657
} - 17658
} - 17659
}); - 17660
// Subscribe the completion watcher before starting the turn so fast - 17661
// scripted/provider responses cannot publish RunFinished into a void. - 17662
tokio::task::yield_now().await; - 17663
begin_turn(&h, &h.core, &scheduled_prompt, false); - 17664
} - 17665
Ok(child_id) - 17666
} - 17667
- 17668
// ---- Watchdog script tasks (docs/design/29-personal-os.md P2) --------------- - 17669
// - 17670
// A `script:` task NEVER reaches the LLM. The shell line runs through the - 17671
// exact brokered bash tool the agent loop uses (`Core::agent_tools()` → - 17672
// BrokeredTool → `__tool_worker`), so sandboxing, environment scrubbing, - 17673
// process-group isolation, and output caps are identical by construction. - 17674
// Zero provider dispatch is a property of the call graph: nothing here can - 17675
// name a Provider. - 17676
- 17677
const SCRIPT_TIMEOUT_MS: u64 = 120_000; - 17678
/// Upper bound on one background reflection pass (docs/design/29 P1) so a - 17679
/// stuck auxiliary stream cannot pin a session handle indefinitely. - 17680
const REFLECTION_CALL_TIMEOUT: Duration = Duration::from_secs(120); - 17681
/// Fallback delivery surface for error alerts when a watchdog has no - 17682
/// `deliver_to`: failures are never silent. - 17683
pub(crate) const FALLBACK_ALERT_TARGET: &str = "log:vak"; - 17684
- 17685
/// Extract the stdout section from BashTool's combined report - 17686
/// ("[stdout]\n…\n[stderr]\n…" or "(no output)"). A literal "[stderr]" - 17687
/// inside the script's own stdout ends the section early — watchdogs that - 17688
/// print the marker get truncated delivery, never a misparse of stderr. - 17689
fn stdout_section(content: &str) -> &str { - 17690
let stdout_part = if let Some(idx) = content.find("[stdout]\n") { - 17691
&content[idx + "[stdout]\n".len()..] - 17692
} else { - 17693
return ""; - 17694
}; - 17695
match stdout_part.find("\n[stderr]") { - 17696
Some(end) => &stdout_part[..end], - 17697
None => stdout_part, - 17698
} - 17699
} - 17700
- 17701
struct ScriptOutcome { - 17702
/// True when the brokered command exited zero within its watchdog. - 17703
ok: bool, - 17704
/// Trimmed stdout on success; combined failure detail otherwise. - 17705
text: String, - 17706
} - 17707
- 17708
async fn execute_script(core: &Core, cwd: &std::path::Path, script: &str) -> ScriptOutcome { - 17709
let bash = core.agent_tools().into_iter().find(|t| t.name() == "bash"); - 17710
let Some(bash) = bash else { - 17711
return ScriptOutcome { - 17712
ok: false, - 17713
text: "script task failed: no bash tool available".to_string(), - 17714
}; - 17715
}; - 17716
let ctx = vak_tools::ToolContext { - 17717
cwd: cwd.to_path_buf(), - 17718
cancel: CancellationToken::new(), - 17719
sandbox: core.agent_sandbox(), - 17720
sandbox_sink: None, - 17721
agent_id: core.agent_identity().map(|a| a.id.clone()), - 17722
new_documents: Vec::new(), - 17723
}; - 17724
let args = serde_json::json!({ "command": script, "timeout_ms": SCRIPT_TIMEOUT_MS }); - 17725
let out = bash.execute(&args, &ctx).await; - 17726
if out.is_error { - 17727
ScriptOutcome { - 17728
ok: false, - 17729
text: format!("script failed: {}", vak_tools::bounded(out.content).trim()), - 17730
} - 17731
} else { - 17732
ScriptOutcome { - 17733
ok: true, - 17734
text: vak_tools::bounded(stdout_section(&out.content).trim().to_string()), - 17735
} - 17736
} - 17737
} - 17738
- 17739
/// Run one watchdog tick: execute, record, deliver. Empty stdout on - 17740
/// success stays silent (zero tokens, zero noise); any failure delivers a - 17741
/// typed error alert even without a configured target. - 17742
async fn fire_script_task(state: &AppState, task: &TaskDef, script: &str) -> Option<String> { - 17743
// One execution at a time per watchdog: a scheduler tick and run-now - 17744
// must never double-fire (or double-deliver) the same tick. - 17745
if !state - 17746
.script_inflight - 17747
.lock() - 17748
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17749
.insert(task.id.clone()) - 17750
{ - 17751
return None; - 17752
} - 17753
update_tasks(state, |map| { - 17754
if let Some(current) = map.get_mut(&task.id) { - 17755
current.last_run_status = Some("working".into()); - 17756
current.last_delivery_state = Some("pending".into()); - 17757
} - 17758
}); - 17759
let outcome = execute_script(&state.core, &task.cwd, script).await; - 17760
let mut delivery_state = "inbox"; - 17761
// Deliver FIRST: once the summary is visible on the task, its delivery - 17762
// attempt has already been made. With zero transports configured the - 17763
// inbox itself is the sink (docs/design/29 P6): a watchdog summary is - 17764
// never lost just because no chat channel exists. - 17765
if outcome.ok { - 17766
if !outcome.text.is_empty() { - 17767
let title = format!("watchdog '{}'", task.name); - 17768
match task.deliver_to.as_deref() { - 17769
Some(target) => { - 17770
match gateway::deliver_and_record_with_result( - 17771
&state.core, - 17772
target, - 17773
&outcome.text, - 17774
vak_core::inbox::Kind::TaskSummary, - 17775
title, - 17776
None, - 17777
Some(&task.id), - 17778
None, - 17779
) - 17780
.await - 17781
{ - 17782
Ok(state) => delivery_state = state, - 17783
Err(_) => delivery_state = "pending", - 17784
} - 17785
} - 17786
None => { - 17787
let _ = vak_core::inbox::record( - 17788
&state.core.shared_data_home(), - 17789
vak_core::inbox::Kind::TaskSummary, - 17790
&title, - 17791
&outcome.text, - 17792
None, - 17793
Some(&task.id), - 17794
); - 17795
} - 17796
} - 17797
} - 17798
} else { - 17799
eprintln!( - 17800
"[scheduler] watchdog '{}' failed: {}", - 17801
task.name, outcome.text - 17802
); - 17803
let target = task.deliver_to.as_deref().unwrap_or(FALLBACK_ALERT_TARGET); - 17804
match gateway::deliver_and_record_with_result( - 17805
&state.core, - 17806
target, - 17807
&format!("watchdog '{}' alert:\n{}", task.name, outcome.text), - 17808
vak_core::inbox::Kind::TaskSummary, - 17809
format!("failure: watchdog '{}'", task.name), - 17810
None, - 17811
Some(&task.id), - 17812
None, - 17813
) - 17814
.await - 17815
{ - 17816
Ok(state) => delivery_state = state, - 17817
Err(_) => delivery_state = "pending", - 17818
} - 17819
} - 17820
update_tasks(state, |map| { - 17821
if let Some(t) = map.get_mut(&task.id) { - 17822
t.last_run_at = Some(chrono::Utc::now()); - 17823
t.last_summary = Some(if outcome.ok && outcome.text.is_empty() { - 17824
"(silent tick)".to_string() - 17825
} else { - 17826
outcome.text.clone() - 17827
}); - 17828
t.last_run_status = Some(if outcome.ok { "complete" } else { "failed" }.into()); - 17829
t.last_delivery_state = Some(delivery_state.into()); - 17830
} - 17831
}); - 17832
check_budget_alert(state, &task.id).await; - 17833
state - 17834
.script_inflight - 17835
.lock() - 17836
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17837
.remove(&task.id); - 17838
Some(task.id.clone()) - 17839
} - 17840
- 17841
// ---- Scheduler (intervals + cron + catch-up, docs/design/29 P2) ------------- - 17842
- 17843
/// True when the first scheduled slot STRICTLY AFTER `last_run_at` has - 17844
/// already arrived by `now` — i.e. a fire was skipped (typically while the - 17845
/// process was down). A run made after the latest slot (manual run-now) - 17846
/// covers it, so nothing is missed. Never-run tasks are decided by the - 17847
/// caller: without history there is nothing to catch up on, and a freshly - 17848
/// created task waits for its first computed slot. Pure; unit-tested - 17849
/// against fixed instants. - 17850
fn cron_slot_missed( - 17851
expr: &str, - 17852
last_run_at: chrono::DateTime<Utc>, - 17853
now: chrono::DateTime<chrono::Local>, - 17854
) -> bool { - 17855
let last_local = last_run_at.with_timezone(&chrono::Local); - 17856
match vak_core::tasks::cron_next_after(expr, last_local) { - 17857
// Instant comparison: correct across DST folds and gaps. - 17858
Ok(next_due) => next_due <= now, - 17859
Err(_) => false, - 17860
} - 17861
} - 17862
- 17863
/// Park an unparseable schedule's marker far in the future: validation - 17864
/// should have rejected it, so this only contains legacy/corrupt entries. - 17865
fn park_marker() -> chrono::DateTime<chrono::Local> { - 17866
chrono::Local::now() + chrono::Duration::days(366) - 17867
} - 17868
- 17869
/// One scheduler pass over enabled tasks for this cwd: interval tasks use - 17870
/// their `last_run_at`; scheduled tasks consult their in-memory next-fire - 17871
/// marker, initializing it to the first future slot when absent (so newly - 17872
/// loaded/created tasks do not stampede on startup). - 17873
async fn scheduler_tick(state: &AppState) { - 17874
let now_local = chrono::Local::now(); - 17875
feeds::scheduled_ingestion(state).await; - 17876
// Reload from disk every tick: tasks.json is shared with the CLI and - 17877
// desktop, so definitions added while the server runs must fire too - 17878
// (found by the v0.6 live battery — CLI-added cron tasks never fired). - 17879
load_tasks(state); - 17880
let due: Vec<String> = { - 17881
let tasks = state - 17882
.tasks - 17883
.lock() - 17884
.unwrap_or_else(std::sync::PoisonError::into_inner); - 17885
let mut markers = state - 17886
.next_fire - 17887
.lock() - 17888
.unwrap_or_else(std::sync::PoisonError::into_inner); - 17889
tasks - 17890
.values() - 17891
.filter(|t| t.enabled && t.cwd.as_path() == state.core.cwd().as_path()) - 17892
.filter(|t| match t.due_at { - 17893
Some(due) => chrono::Utc::now() >= due, - 17894
None => match t.schedule.as_deref() { - 17895
Some(expr) => { - 17896
if let Some(zone) = t.timezone.as_deref() { - 17897
let anchor = t.last_run_at.unwrap_or(t.created_at); - 17898
vak_core::tasks::cron_next_after_timezone(expr, anchor, zone) - 17899
.map(|next| chrono::Utc::now() >= next) - 17900
.unwrap_or(false) - 17901
} else { - 17902
let marker = markers.entry(t.id.clone()).or_insert_with(|| { - 17903
vak_core::tasks::cron_next_after(expr, now_local) - 17904
.unwrap_or_else(|_| park_marker()) - 17905
}); - 17906
now_local >= *marker - 17907
} - 17908
} - 17909
None => t - 17910
.last_run_at - 17911
.map(|l| { - 17912
(now_local.with_timezone(&Utc) - l).num_seconds() - 17913
>= t.interval_secs as i64 - 17914
}) - 17915
.unwrap_or(true), - 17916
}, - 17917
}) - 17918
.map(|t| t.id.clone()) - 17919
.collect() - 17920
}; - 17921
// A cron slot is spent only by a run that started: a refused or busy - 17922
// task keeps its slot and is tried again next tick. - 17923
for id in due { - 17924
if fire_task(state, &id).await.is_ok() { - 17925
advance_marker(state, &id); - 17926
} - 17927
} - 17928
} - 17929
- 17930
fn advance_marker(state: &AppState, id: &str) { - 17931
let expr = state - 17932
.tasks - 17933
.lock() - 17934
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17935
.get(id) - 17936
.and_then(|t| t.schedule.clone()); - 17937
if let Some(expr) = expr { - 17938
let next = vak_core::tasks::cron_next_after(&expr, chrono::Local::now()) - 17939
.unwrap_or_else(|_| park_marker()); - 17940
state - 17941
.next_fire - 17942
.lock() - 17943
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17944
.insert(id.to_string(), next); - 17945
} - 17946
} - 17947
- 17948
/// Startup catch-up (docs/design/29-personal-os.md P2): when enabled and a - 17949
/// scheduled task's most recent slot happened after its last run — a slot - 17950
/// missed while the process was down — fire it once immediately. Interval - 17951
/// tasks keep their self-healing `>= interval` behavior and need nothing. - 17952
async fn catch_up_missed_tasks(state: &AppState) { - 17953
if !state.core.config().automation.catch_up_missed { - 17954
return; - 17955
} - 17956
let now_local = chrono::Local::now(); - 17957
let due: Vec<String> = { - 17958
let tasks = state - 17959
.tasks - 17960
.lock() - 17961
.unwrap_or_else(std::sync::PoisonError::into_inner); - 17962
tasks - 17963
.values() - 17964
.filter(|t| t.enabled && t.cwd.as_path() == state.core.cwd().as_path()) - 17965
.filter_map(|t| { - 17966
let expr = t.schedule.as_deref()?; - 17967
t.last_run_at - 17968
.is_some_and(|l| cron_slot_missed(expr, l, now_local)) - 17969
.then(|| t.id.clone()) - 17970
}) - 17971
.collect() - 17972
}; - 17973
for id in due { - 17974
if fire_task(state, &id).await.is_ok() { - 17975
advance_marker(state, &id); - 17976
} - 17977
} - 17978
} - 17979
- 17980
/// Background loop: evaluates due tasks every 20 seconds. Holds only weak - 17981
/// state via `state` clones living inside the router — when the server - 17982
/// shuts down the loop dies with the runtime. - 17983
pub fn start_scheduler(state: &AppState) { - 17984
load_tasks(state); - 17985
let st = state.clone(); - 17986
tokio::spawn(async move { catch_up_missed_tasks(&st).await }); - 17987
let st = state.clone(); - 17988
tokio::spawn(async move { - 17989
let mut tick = tokio::time::interval(std::time::Duration::from_secs(20)); - 17990
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - 17991
loop { - 17992
tick.tick().await; - 17993
scheduler_tick(&st).await; - 17994
} - 17995
}); - 17996
// Proactive heartbeat (docs/design/29-personal-os.md P7): its own - 17997
// per-process timer alongside the task tick; the pass itself re-checks - 17998
// the enabled flag every beat. - 17999
if state.core.config().heartbeat.enabled { - 18000
let st = state.clone(); - 18001
tokio::spawn(async move { - 18002
let mut tick = - 18003
tokio::time::interval(std::time::Duration::from_secs(heartbeat::TICK_SECS)); - 18004
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - 18005
loop { - 18006
tick.tick().await; - 18007
heartbeat::heartbeat_tick(&st).await; - 18008
} - 18009
}); - 18010
} - 18011
- 18012
// Commitment upkeep runs on its own timer, deliberately NOT gated on - 18013
// `heartbeat.enabled` (docs/design/47-commitment-kernel.md). Heartbeat is - 18014
// an opt-in model pass that costs tokens; this is clock and filesystem - 18015
// work that costs none. Tying durable work's upkeep to an opt-in prober - 18016
// would mean a commitment stopped being durable the moment somebody - 18017
// switched the prober off — and a suspended commitment nobody wakes is - 18018
// indistinguishable from lost work. - 18019
if state.core.config().commitment.enabled { - 18020
let st = state.clone(); - 18021
tokio::spawn(async move { - 18022
let mut tick = tokio::time::interval(COMMITMENT_TICK); - 18023
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - 18024
loop { - 18025
tick.tick().await; - 18026
let report = vak_core::commitments::maintain_all(&st.core.shared_data_home()).await; - 18027
if report.is_empty() { - 18028
continue; - 18029
} - 18030
// A commitment waking, lapsing, or being abandoned by policy - 18031
// is a thing that happened without anybody asking for it, so - 18032
// it lands in the attention layer rather than only in a log. - 18033
for id in report.expired.iter().chain(report.escalated.iter()) { - 18034
let _ = vak_core::inbox::record( - 18035
&st.core.shared_data_home(), - 18036
vak_core::inbox::Kind::TaskSummary, - 18037
"Commitment closed without you", - 18038
&format!("{id} reached the end of its window or escalation policy."), - 18039
None, - 18040
None, - 18041
); - 18042
} - 18043
for id in report.resumed.iter().chain(report.satisfied.iter()) { - 18044
eprintln!("[commit] {id} resumed"); - 18045
} - 18046
} - 18047
}); - 18048
} - 18049
} - 18050
- 18051
/// Commitment upkeep cadence. Slower than the heartbeat tick because nothing - 18052
/// here is latency-sensitive: a scheduled wake a minute late is fine, and a - 18053
/// tighter loop would just re-read the ledger for nothing. - 18054
const COMMITMENT_TICK: std::time::Duration = std::time::Duration::from_secs(60); - 18055
- 18056
// ---- Budget alerts (docs/design/29-personal-os.md P2) ----------------------- - 18057
- 18058
/// Proactive day-cap alerting, called after every task fire and exposed for - 18059
/// surfaces to invoke wherever day spend updates land. Fires at most once - 18060
/// per (level, UTC-day window): the audit row recorded by `record_alert` - 18061
/// doubles as the once-per-window marker. Delivery reuses the gateway - 18062
/// transports (`log:` / `webhook:` / telegram bridge); with no configured - 18063
/// targets it falls back to the log surface so an approaching cap is never - 18064
/// discovered at denial time. - 18065
pub async fn check_budget_alert(state: &AppState, session_id: &str) { - 18066
let Some(cap) = state.core.effective_finops_max_day_usd() else { - 18067
return; - 18068
}; - 18069
// Read through the EFFECTIVE sessions home (an embedded server may - 18070
// have relocated it); Core::spend_day_usd pins the constructed path. - 18071
let mut day_total = vak_core::finops::FinOpsLedger::new(&state.core.shared_data_home()) - 18072
.day_total_usd(Utc::now()); - 18073
if state.core.sessions_home() != state.core.shared_data_home() { - 18074
day_total += vak_core::finops::FinOpsLedger::new(&state.core.sessions_home()) - 18075
.day_total_usd(Utc::now()); - 18076
} - 18077
let Some(level) = vak_core::finops::alert_level(day_total, cap) else { - 18078
return; - 18079
}; - 18080
let home = state.core.shared_data_home(); - 18081
if let Some(last) = vak_core::finops::last_alert(&home, level) - 18082
&& last.ts.with_timezone(&Utc).date_naive() == Utc::now().date_naive() - 18083
{ - 18084
return; // this level already alerted inside the current day window - 18085
} - 18086
if let Err(e) = vak_core::finops::record_alert(&home, level, session_id) { - 18087
eprintln!("[finops] budget-alert ledger write failed: {e}"); - 18088
} - 18089
let text = format!( - 18090
"budget alert [{}]: ${:.2} of ${:.2} daily cap", - 18091
level.as_str(), - 18092
day_total, - 18093
cap - 18094
); - 18095
let title = format!("budget alert [{}]", level.as_str()); - 18096
let mut targets = configured_delivery_targets(state); - 18097
if targets.is_empty() { - 18098
targets.push(FALLBACK_ALERT_TARGET.to_string()); - 18099
} - 18100
for target in targets { - 18101
let _ = gateway::deliver_and_record( - 18102
&state.core, - 18103
&target, - 18104
&text, - 18105
vak_core::inbox::Kind::BudgetAlert, - 18106
title.clone(), - 18107
Some(session_id), - 18108
None, - 18109
) - 18110
.await; - 18111
} - 18112
} - 18113
- 18114
/// Every distinct `deliver_to` routing target configured across all known - 18115
/// tasks — the server's vocabulary of delivery surfaces. - 18116
pub(crate) fn configured_delivery_targets(state: &AppState) -> Vec<String> { - 18117
let mut targets: Vec<String> = state - 18118
.tasks - 18119
.lock() - 18120
.unwrap_or_else(std::sync::PoisonError::into_inner) - 18121
.values() - 18122
.filter_map(|t| t.deliver_to.clone()) - 18123
.collect(); - 18124
targets.sort(); - 18125
targets.dedup(); - 18126
targets - 18127
} - 18128
- 18129
// ---- Dev-server lifecycle (preview pane) ----------------------------------- - 18130
- 18131
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] - 18132
pub struct LaunchConfig { - 18133
pub name: String, - 18134
pub cmd: String, - 18135
#[serde(default)] - 18136
pub args: Vec<String>, - 18137
#[serde(default)] - 18138
pub port: Option<u16>, - 18139
} - 18140
- 18141
struct ManagedProc { - 18142
child: tokio::process::Child, - 18143
logs: Arc<Mutex<std::collections::VecDeque<String>>>, - 18144
} - 18145
- 18146
fn parse_launch_toml(cwd: &std::path::Path) -> Result<Vec<LaunchConfig>, String> { - 18147
let path = cwd.join(".vak/launch.toml"); - 18148
if !path.exists() { - 18149
return Ok(Vec::new()); - 18150
} - 18151
let raw = std::fs::read_to_string(&path).map_err(|e| e.to_string())?; - 18152
#[derive(serde::Deserialize)] - 18153
struct File { - 18154
#[serde(rename = "server", default)] - 18155
servers: Vec<LaunchConfig>, - 18156
} - 18157
let f: File = toml::from_str(&raw).map_err(|e| format!("launch.toml: {e}"))?; - 18158
Ok(f.servers) - 18159
} - 18160
- 18161
/// Sensible fallback when no launch.toml exists across supported runtimes: - 18162
/// Node/JS/TS (Vite, Next, Astro, React, Nuxt), Python (FastAPI/Uvicorn, Flask, Streamlit, Django), - 18163
/// Rust (cargo run), Go (go run .), or static HTML (http.server). - 18164
fn detect_launch(cwd: &std::path::Path) -> Vec<LaunchConfig> { - 18165
let mut servers = Vec::new(); - 18166
- 18167
// 1. JavaScript / TypeScript projects (package.json) - 18168
let pkg = cwd.join("package.json"); - 18169
if let Ok(raw) = std::fs::read_to_string(&pkg) - 18170
&& let Ok(v) = serde_json::from_str::<serde_json::Value>(&raw) - 18171
{ - 18172
let scripts = &v["scripts"]; - 18173
let (script_name, dev_cmd) = if scripts["dev"].is_string() { - 18174
("dev", "dev") - 18175
} else if scripts["start"].is_string() { - 18176
("start", "start") - 18177
} else if scripts["serve"].is_string() { - 18178
("serve", "serve") - 18179
} else { - 18180
("", "") - 18181
}; - 18182
- 18183
if !script_name.is_empty() { - 18184
let package_manager = javascript_package_manager(&v, |name| cwd.join(name).is_file()); - 18185
- 18186
let raw_lower = raw.to_ascii_lowercase(); - 18187
let port = if raw_lower.contains("vite") { - 18188
Some(5173) - 18189
} else if raw_lower.contains("astro") { - 18190
Some(4321) - 18191
} else { - 18192
Some(3000) - 18193
}; - 18194
- 18195
servers.push(LaunchConfig { - 18196
name: script_name.into(), - 18197
cmd: package_manager.into(), - 18198
args: vec!["run".into(), dev_cmd.into()], - 18199
port, - 18200
}); - 18201
} - 18202
} - 18203
- 18204
// 2. Python projects - 18205
let manage_py = cwd.join("manage.py"); - 18206
if manage_py.exists() { - 18207
servers.push(LaunchConfig { - 18208
name: "django".into(), - 18209
cmd: "python3".into(), - 18210
args: vec!["manage.py".into(), "runserver".into(), "8000".into()], - 18211
port: Some(8000), - 18212
}); - 18213
} - 18214
- 18215
let main_py = cwd.join("main.py"); - 18216
let app_py = cwd.join("app.py"); - 18217
let py_entry = if main_py.exists() { - 18218
Some(("main", "main.py")) - 18219
} else if app_py.exists() { - 18220
Some(("app", "app.py")) - 18221
} else { - 18222
None - 18223
}; - 18224
- 18225
if let Some((mod_name, file_name)) = py_entry { - 18226
let content = std::fs::read_to_string(cwd.join(file_name)) - 18227
.unwrap_or_default() - 18228
.to_ascii_lowercase(); - 18229
if content.contains("fastapi") || content.contains("uvicorn") { - 18230
servers.push(LaunchConfig { - 18231
name: "fastapi".into(), - 18232
cmd: "python3".into(), - 18233
args: vec![ - 18234
"-m".into(), - 18235
"uvicorn".into(), - 18236
format!("{mod_name}:app"), - 18237
"--reload".into(), - 18238
"--port".into(), - 18239
"8000".into(), - 18240
], - 18241
port: Some(8000), - 18242
}); - 18243
} else if content.contains("flask") { - 18244
servers.push(LaunchConfig { - 18245
name: "flask".into(), - 18246
cmd: "python3".into(), - 18247
args: vec![file_name.into()], - 18248
port: Some(5000), - 18249
}); - 18250
} else if content.contains("streamlit") { - 18251
servers.push(LaunchConfig { - 18252
name: "streamlit".into(), - 18253
cmd: "streamlit".into(), - 18254
args: vec![ - 18255
"run".into(), - 18256
file_name.into(), - 18257
"--server.port".into(), - 18258
"8501".into(), - 18259
], - 18260
port: Some(8501), - 18261
}); - 18262
} - 18263
} - 18264
- 18265
// 3. Rust projects - 18266
let cargo_toml = cwd.join("Cargo.toml"); - 18267
if cargo_toml.exists() && (cwd.join("src/main.rs").exists() || cwd.join("src/bin").exists()) { - 18268
servers.push(LaunchConfig { - 18269
name: "cargo".into(), - 18270
cmd: "cargo".into(), - 18271
args: vec!["run".into()], - 18272
port: Some(8080), - 18273
}); - 18274
} - 18275
- 18276
// 4. Go projects - 18277
let go_mod = cwd.join("go.mod"); - 18278
let main_go = cwd.join("main.go"); - 18279
if go_mod.exists() || main_go.exists() { - 18280
servers.push(LaunchConfig { - 18281
name: "go".into(), - 18282
cmd: "go".into(), - 18283
args: vec!["run".into(), ".".into()], - 18284
port: Some(8080), - 18285
}); - 18286
} - 18287
- 18288
// 5. Static HTML fallback - 18289
let has_top_level_html = cwd.join("index.html").is_file() - 18290
|| std::fs::read_dir(cwd).is_ok_and(|entries| { - 18291
entries.filter_map(Result::ok).any(|entry| { - 18292
entry.path().is_file() - 18293
&& entry - 18294
.path() - 18295
.extension() - 18296
.and_then(|value| value.to_str()) - 18297
.is_some_and(|extension| { - 18298
extension.eq_ignore_ascii_case("html") - 18299
|| extension.eq_ignore_ascii_case("htm") - 18300
}) - 18301
}) - 18302
}); - 18303
if servers.is_empty() && has_top_level_html { - 18304
servers.push(LaunchConfig { - 18305
name: "static".into(), - 18306
cmd: "python3".into(), - 18307
args: vec!["-m".into(), "http.server".into(), "8080".into()], - 18308
port: Some(8080), - 18309
}); - 18310
} - 18311
- 18312
servers - 18313
} - 18314
- 18315
fn javascript_dependencies_missing(cwd: &std::path::Path) -> bool { - 18316
let Ok(bytes) = std::fs::read(cwd.join("package.json")) else { - 18317
return false; - 18318
}; - 18319
let Ok(package) = serde_json::from_slice::<serde_json::Value>(&bytes) else { - 18320
return false; - 18321
}; - 18322
let declared = ["dependencies", "devDependencies", "peerDependencies"] - 18323
.into_iter() - 18324
.any(|key| { - 18325
package - 18326
.get(key) - 18327
.and_then(serde_json::Value::as_object) - 18328
.is_some_and(|values| !values.is_empty()) - 18329
}); - 18330
declared && !cwd.join("node_modules").is_dir() - 18331
} - 18332
- 18333
async fn get_launch( - 18334
State(state): State<AppState>, - 18335
Path(id): Path<String>, - 18336
axum::extract::Query(scope): axum::extract::Query<LaunchScope>, - 18337
) -> Json<serde_json::Value> { - 18338
let Some(workspace) = sandbox_session_workspace(&state, &id) else { - 18339
return Json(serde_json::json!({ "error": "unknown session" })); - 18340
}; - 18341
let launch_root = match launch_root(&state, &id, scope.candidate_id.as_deref(), &workspace) { - 18342
Ok(root) => root, - 18343
Err(error) => return Json(serde_json::json!({ "error": error })), - 18344
}; - 18345
let mut servers = match parse_launch_toml(&launch_root) { - 18346
Ok(s) => s, - 18347
Err(e) => return Json(serde_json::json!({ "error": e })), - 18348
}; - 18349
if servers.is_empty() { - 18350
servers = detect_launch(&launch_root); - 18351
} - 18352
let mut procs = state - 18353
.procs - 18354
.lock() - 18355
.unwrap_or_else(std::sync::PoisonError::into_inner); - 18356
procs.retain(|_, process| match process.child.try_wait() { - 18357
Ok(Some(_)) => { - 18358
vak_tools::bash::kill_process_group(&process.child.id()); - 18359
false - 18360
} - 18361
Ok(None) | Err(_) => true, - 18362
}); - 18363
let list: Vec<serde_json::Value> = servers - 18364
.into_iter() - 18365
.map(|mut s| { - 18366
let key = proc_key(&id, scope.candidate_id.as_deref(), &s.name); - 18367
let running = procs.contains_key(&key); - 18368
let executable_available = vak_tools::bash::executable_available(&s.cmd, &launch_root); - 18369
let preparation_required = matches!(s.cmd.as_str(), "npm" | "pnpm" | "yarn" | "bun") - 18370
&& javascript_dependencies_missing(&launch_root); - 18371
let port_available = running || s.port.is_none_or(port_is_available); - 18372
let available = executable_available && !preparation_required && port_available; - 18373
let availability = if !executable_available { - 18374
"needs_setup" - 18375
} else if preparation_required { - 18376
"needs_preparation" - 18377
} else if !port_available { - 18378
"port_in_use" - 18379
} else { - 18380
"ready" - 18381
}; - 18382
let unavailable_reason = if !executable_available { - 18383
Some(format!( - 18384
"{} is not available in the preview environment", - 18385
s.cmd - 18386
)) - 18387
} else if preparation_required { - 18388
Some("project dependencies have not been prepared for this saved version".into()) - 18389
} else if !port_available { - 18390
s.port.map(|port| format!("port {port} is already in use")) - 18391
} else { - 18392
None - 18393
}; - 18394
if running && s.port.is_none() { - 18395
s.port = None; - 18396
} - 18397
serde_json::json!({ - 18398
"name": s.name, - 18399
"cmd": s.cmd, - 18400
"args": s.args, - 18401
"port": s.port, - 18402
"running": running, - 18403
"available": available, - 18404
"availability": availability, - 18405
"unavailable_reason": unavailable_reason, - 18406
}) - 18407
}) - 18408
.collect(); - 18409
Json(serde_json::json!({ "servers": list })) - 18410
} - 18411
- 18412
fn proc_key(session: &str, candidate_id: Option<&str>, name: &str) -> String { - 18413
format!("{session}::{}::{name}", candidate_id.unwrap_or("workspace")) - 18414
} - 18415
- 18416
#[derive(Default, serde::Deserialize)] - 18417
struct LaunchScope { - 18418
candidate_id: Option<String>, - 18419
} - 18420
- 18421
fn launch_root( - 18422
state: &AppState, - 18423
session_id: &str, - 18424
candidate_id: Option<&str>, - 18425
workspace: &std::path::Path, - 18426
) -> Result<std::path::PathBuf, String> { - 18427
let Some(candidate_id) = candidate_id else { - 18428
return Ok(workspace.to_path_buf()); - 18429
}; - 18430
let saved = saved_launch_candidate(state, session_id, candidate_id)?; - 18431
let expected = sandbox_candidates_root(state).join(candidate_id); - 18432
if saved.candidate.source_root != expected { - 18433
return Err("saved draft root is invalid".into()); - 18434
} - 18435
verify_launch_tree(&saved.candidate, &expected)?; - 18436
let prepared = sandbox_previews_root(state).join(candidate_id); - 18437
if std::fs::read_to_string(prepared.join(".vak-candidate-digest")) - 18438
.is_ok_and(|digest| digest == saved.candidate_digest) - 18439
{ - 18440
verify_launch_tree(&saved.candidate, &prepared)?; - 18441
return Ok(prepared); - 18442
} - 18443
Ok(expected) - 18444
} - 18445
- 18446
fn saved_launch_candidate( - 18447
state: &AppState, - 18448
session_id: &str, - 18449
candidate_id: &str, - 18450
) -> Result<vak_sandbox::CandidateRecord, String> { - 18451
vak_sandbox::load_records(&sandbox_records_path(state)) - 18452
.map_err(|error| error.to_string())? - 18453
.into_iter() - 18454
.rev() - 18455
.find_map(|record| match record { - 18456
vak_sandbox::DurableRecord::Candidate(saved) - 18457
if saved.session_id == session_id - 18458
&& saved.candidate.candidate_id == candidate_id => - 18459
{ - 18460
Some(saved) - 18461
} - 18462
_ => None, - 18463
}) - 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 (
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.