- 16650
) - 16651
.into_response() - 16652
} - 16653
Err(e) => ( - 16654
StatusCode::INTERNAL_SERVER_ERROR, - 16655
Json(serde_json::json!({ "error": e.to_string() })), - 16656
) - 16657
.into_response(), - 16658
} - 16659
} - 16660
- 16661
async fn discard_best_run( - 16662
State(state): State<AppState>, - 16663
Path(id): Path<String>, - 16664
) -> axum::response::Response { - 16665
use axum::response::IntoResponse; - 16666
let meta = state - 16667
.best_runs - 16668
.lock() - 16669
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16670
.get(&id) - 16671
.cloned(); - 16672
let Some(meta) = meta else { - 16673
return StatusCode::NOT_FOUND.into_response(); - 16674
}; - 16675
cleanup_worktree(&state, &id, &meta); - 16676
(StatusCode::OK, Json(serde_json::json!({"discarded": id}))).into_response() - 16677
} - 16678
- 16679
fn cleanup_worktree(state: &AppState, child_id: &str, meta: &BestRunMeta) { - 16680
let wt = vak_core::worktree::Worktree { - 16681
path: meta.wt_path.clone(), - 16682
branch: meta.branch.clone(), - 16683
}; - 16684
let _ = vak_core::worktree::remove(&meta.repo, &wt); - 16685
state - 16686
.best_runs - 16687
.lock() - 16688
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16689
.remove(child_id); - 16690
} - 16691
- 16692
// ---- PR monitoring (gh-backed) --------------------------------------------- - 16693
- 16694
async fn gh_output(cwd: &std::path::Path, args: &[&str]) -> Result<String, String> { - 16695
let out = tokio::process::Command::new("gh") - 16696
.args(args) - 16697
.current_dir(cwd) - 16698
.output() - 16699
.await - 16700
.map_err(|e| format!("gh not available: {e}"))?; - 16701
if !out.status.success() { - 16702
return Err(String::from_utf8_lossy(&out.stderr).trim().to_string()); - 16703
} - 16704
Ok(String::from_utf8_lossy(&out.stdout).into_owned()) - 16705
} - 16706
- 16707
/// One-shot PR status for the session workspace: current branch, the PR - 16708
/// attached to it (if any), and its check rollup. Tooling absence surfaces - 16709
/// as `{error}` — never a hang, never a panic. - 16710
async fn session_pr( - 16711
State(state): State<AppState>, - 16712
Path(id): Path<String>, - 16713
) -> Json<serde_json::Value> { - 16714
let Some(handle) = state.get(&id) else { - 16715
return Json(serde_json::json!({ "error": "unknown session" })); - 16716
}; - 16717
let cwd = handle.cwd.clone(); - 16718
let Some(branch) = git_output(&cwd, &["rev-parse", "--abbrev-ref", "HEAD"]).await else { - 16719
return Json(serde_json::json!({ "error": "not a git repository" })); - 16720
}; - 16721
let branch = branch.trim().to_string(); - 16722
if branch.is_empty() || branch == "HEAD" { - 16723
return Json(serde_json::json!({ "error": "detached HEAD" })); - 16724
} - 16725
- 16726
let raw = match gh_output( - 16727
&cwd, - 16728
&[ - 16729
"pr", - 16730
"view", - 16731
&branch, - 16732
"--json", - 16733
"number,title,url,state,mergeable,statusCheckRollup", - 16734
], - 16735
) - 16736
.await - 16737
{ - 16738
Ok(r) => r, - 16739
Err(e) => { - 16740
// No PR for this branch vs no gh at all — distinguish for UX. - 16741
let msg = e.to_lowercase(); - 16742
let kind = if msg.contains("no pull requests") || msg.contains("no merges requested") { - 16743
"no_pr" - 16744
} else { - 16745
"gh_unavailable" - 16746
}; - 16747
return Json(serde_json::json!({ - 16748
"branch": branch, - 16749
"pr": serde_json::Value::Null, - 16750
"reason": kind, - 16751
"error": e, - 16752
})); - 16753
} - 16754
}; - 16755
- 16756
let Ok(view) = serde_json::from_str::<serde_json::Value>(&raw) else { - 16757
return Json(serde_json::json!({ "error": "unparsable gh output", "branch": branch })); - 16758
}; - 16759
let mut pass = 0u32; - 16760
let mut fail = 0u32; - 16761
let mut pending = 0u32; - 16762
if let Some(rollup) = view["statusCheckRollup"].as_array() { - 16763
for c in rollup { - 16764
let status = c["status"].as_str().unwrap_or(""); - 16765
let conclusion = c["conclusion"].as_str().unwrap_or(""); - 16766
match (status, conclusion) { - 16767
(_, "SUCCESS") => pass += 1, - 16768
(_, "FAILURE") | (_, "CANCELLED") | (_, "TIMED_OUT") => fail += 1, - 16769
("COMPLETED", _) => {} - 16770
_ => pending += 1, - 16771
} - 16772
} - 16773
} - 16774
Json(serde_json::json!({ - 16775
"branch": branch, - 16776
"pr": { - 16777
"number": view["number"], - 16778
"title": view["title"], - 16779
"url": view["url"], - 16780
"state": view["state"], - 16781
"mergeable": view["mergeable"], - 16782
}, - 16783
"checks": view["statusCheckRollup"], - 16784
"summary": { "pass": pass, "fail": fail, "pending": pending }, - 16785
})) - 16786
} - 16787
- 16788
#[derive(serde::Deserialize)] - 16789
struct PrMergeBody { - 16790
number: u64, - 16791
#[serde(default)] - 16792
method: Option<String>, - 16793
} - 16794
- 16795
/// Merge an open PR via gh. `--auto` honors branch protection: gh merges - 16796
/// when checks go green. - 16797
async fn pr_merge( - 16798
State(state): State<AppState>, - 16799
Path(id): Path<String>, - 16800
Json(body): Json<PrMergeBody>, - 16801
) -> axum::response::Response { - 16802
use axum::response::IntoResponse; - 16803
let Some(handle) = state.get(&id) else { - 16804
return StatusCode::NOT_FOUND.into_response(); - 16805
}; - 16806
let method = body.method.unwrap_or_else(|| "squash".to_string()); - 16807
let flag = match method.as_str() { - 16808
"merge" => "--merge", - 16809
"rebase" => "--rebase", - 16810
_ => "--squash", - 16811
}; - 16812
let mut args = vec![ - 16813
"pr".to_string(), - 16814
"merge".to_string(), - 16815
body.number.to_string(), - 16816
flag.to_string(), - 16817
"--auto".to_string(), - 16818
]; - 16819
if method == "squash" { - 16820
args.push("--delete-branch".to_string()); - 16821
} - 16822
let arg_refs: Vec<&str> = args.iter().map(String::as_str).collect(); - 16823
match gh_output(&handle.cwd, &arg_refs).await { - 16824
Ok(_) => ( - 16825
StatusCode::OK, - 16826
Json(serde_json::json!({ "merging": body.number })), - 16827
) - 16828
.into_response(), - 16829
Err(e) => ( - 16830
StatusCode::CONFLICT, - 16831
Json(serde_json::json!({ "error": e })), - 16832
) - 16833
.into_response(), - 16834
} - 16835
} - 16836
- 16837
// ---- Scheduled tasks (local routines) -------------------------------------- - 16838
- 16839
/// Final assistant text of a session's active chain, if any. Used to give - 16840
/// routine runs a real answer instead of a status word. - 16841
fn last_assistant_text(handle: &SessionHandle) -> Option<String> { - 16842
let guard = handle - 16843
.session - 16844
.lock() - 16845
.unwrap_or_else(std::sync::PoisonError::into_inner); - 16846
let log = guard.as_ref()?; - 16847
let narration = log - 16848
.derive_messages() - 16849
.into_iter() - 16850
.rev() - 16851
.find(|m| m.role == vak_llm::Role::Assistant) - 16852
.map(|m| m.text_content()) - 16853
.unwrap_or_default(); - 16854
// A run whose answer was a card has little or no narration; the cards are - 16855
// still the answer. - 16856
let text = crate::projection::text_with_run_cards( - 16857
log, - 16858
crate::projection::clean_scaffolding(&narration), - 16859
); - 16860
if text.trim().is_empty() { - 16861
None - 16862
} else { - 16863
Some(text) - 16864
} - 16865
} - 16866
- 16867
fn tasks_file(core: &Core) -> PathBuf { - 16868
let shared = vak_core::tasks::tasks_file(&core.shared_data_home()); - 16869
if shared.exists() || core.sessions_home() == core.shared_data_home() { - 16870
shared - 16871
} else { - 16872
let session = vak_core::tasks::tasks_file(&core.sessions_home()); - 16873
if session.exists() { session } else { shared } - 16874
} - 16875
} - 16876
- 16877
/// Loads `tasks.json` and makes `state.tasks` match it exactly (inserts, - 16878
/// updates *and* removals) rather than merging insert-only. The whole - 16879
/// read-and-replace runs under `state.tasks`'s lock so it can never - 16880
/// interleave with `update_tasks`'s mutate-then-persist below: either this - 16881
/// runs entirely before a concurrent create/update/delete's persist, or - 16882
/// entirely after, never in the gap between that mutation's memory write - 16883
/// and its disk write. Previously an insert-only merge meant an external - 16884
/// delete (CLI, desktop app) — or even this process's own `delete_task` - 16885
/// racing a scheduler tick — could be silently undone the next time - 16886
/// anything called `update_tasks`, since the removed id would still be on - 16887
/// disk and get merged straight back into memory. - 16888
fn load_tasks(state: &AppState) { - 16889
// The disk read itself must happen while holding the lock, not before - 16890
// it: reading first and only acquiring the lock to apply the snapshot - 16891
// leaves a gap where a concurrent `update_tasks` (create/update/delete) - 16892
// can mutate memory *and* persist in between the read and the replace. - 16893
// This function would then overwrite that fresh insert with the stale - 16894
// pre-mutation snapshot it already had in hand, silently losing it — - 16895
// exactly the kind of loss the merge-vs-replace note below was written - 16896
// to prevent, just moved one step earlier. - 16897
let mut map = state - 16898
.tasks - 16899
.lock() - 16900
.unwrap_or_else(std::sync::PoisonError::into_inner); - 16901
let mut tasks_vec = match vak_core::tasks::TaskStore::load(&state.core.shared_data_home()) { - 16902
Ok(store) => store.all(), - 16903
Err(_) => Vec::new(), - 16904
}; - 16905
if state.core.sessions_home() != state.core.shared_data_home() - 16906
&& let Ok(store) = vak_core::tasks::TaskStore::load(&state.core.sessions_home()) - 16907
{ - 16908
for t in store.all() { - 16909
if !tasks_vec.iter().any(|existing| existing.id == t.id) { - 16910
tasks_vec.push(t); - 16911
} - 16912
} - 16913
} - 16914
*map = tasks_vec.into_iter().map(|t| (t.id.clone(), t)).collect(); - 16915
if recover_interrupted_tasks(&mut map) { - 16916
write_tasks_file(state, &map); - 16917
} - 16918
} - 16919
- 16920
fn recover_interrupted_tasks(tasks: &mut HashMap<String, TaskDef>) -> bool { - 16921
let mut recovered = false; - 16922
for task in tasks.values_mut() { - 16923
if task.last_run_status.as_deref() == Some("working") { - 16924
task.last_run_status = Some("interrupted".into()); - 16925
task.last_delivery_state = Some("pending".into()); - 16926
recovered = true; - 16927
} - 16928
} - 16929
recovered - 16930
} - 16931
- 16932
/// fsyncs a directory so a prior rename into it is durable across a crash, - 16933
/// not just torn-write-free while running. No-op on non-unix, where the - 16934
/// rename itself is still atomic but directory fsync isn't a thing. - 16935
#[cfg(unix)] - 16936
fn sync_tasks_dir(path: &std::path::Path) -> std::io::Result<()> { - 16937
std::fs::File::open(path).and_then(|dir| dir.sync_all()) - 16938
} - 16939
#[cfg(not(unix))] - 16940
fn sync_tasks_dir(_path: &std::path::Path) -> std::io::Result<()> { - 16941
Ok(()) - 16942
} - 16943
- 16944
/// Serializes `map` and writes it to `tasks.json` atomically and durably: - 16945
/// write-to-temp, fsync the temp file, rename over the real path, fsync the - 16946
/// directory. Mirrors `vak_core::tasks::TaskStore::save` (and the same - 16947
/// crash-durability fix) since this is a second, independent writer of the - 16948
/// same file — kept in sync here because `AppState.tasks` lives in the - 16949
/// server, not in a `TaskStore`. - 16950
fn write_tasks_file(state: &AppState, map: &HashMap<String, TaskDef>) { - 16951
let mut list: Vec<TaskDef> = map.values().cloned().collect(); - 16952
list.sort_by_key(|t| t.created_at); - 16953
let target = tasks_file(&state.core); - 16954
let result = (|| -> std::io::Result<()> { - 16955
if let Some(parent) = target.parent() { - 16956
std::fs::create_dir_all(parent)?; - 16957
} - 16958
let json = serde_json::to_string_pretty(&list)?; - 16959
let tmp = target.with_extension("json.tmp"); - 16960
{ - 16961
let mut file = std::fs::File::create(&tmp)?; - 16962
file.write_all(json.as_bytes())?; - 16963
file.sync_all()?; - 16964
} - 16965
std::fs::rename(&tmp, &target)?; - 16966
if let Some(parent) = target.parent() { - 16967
sync_tasks_dir(parent)?; - 16968
} - 16969
Ok(()) - 16970
})(); - 16971
if let Err(e) = result { - 16972
eprintln!("[scheduler] tasks file save failed: {e}"); - 16973
} - 16974
} - 16975
- 16976
/// Mutates `state.tasks` and persists the result to disk under a single - 16977
/// hold of the lock, so no other reader/writer (in particular - 16978
/// `load_tasks`'s scheduler tick) can observe or race the gap between the - 16979
/// in-memory change and the on-disk write. - 16980
fn update_tasks<T>(state: &AppState, f: impl FnOnce(&mut HashMap<String, TaskDef>) -> T) -> T { - 16981
let mut map = state - 16982
.tasks - 16983
.lock() - 16984
.unwrap_or_else(std::sync::PoisonError::into_inner); - 16985
let result = f(&mut map); - 16986
write_tasks_file(state, &map); - 16987
result - 16988
} - 16989
- 16990
async fn list_tasks(State(state): State<AppState>) -> Json<serde_json::Value> { - 16991
let cwd = state.core.cwd().clone(); - 16992
let next_fire = state - 16993
.next_fire - 16994
.lock() - 16995
.unwrap_or_else(std::sync::PoisonError::into_inner); - 16996
let mut mine: Vec<TaskDef> = state - 16997
.tasks - 16998
.lock() - 16999
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17000
.values() - 17001
.filter(|t| t.cwd == cwd) - 17002
.cloned() - 17003
.collect(); - 17004
mine.sort_by_key(|t| t.created_at); - 17005
let now = chrono::Local::now(); - 17006
let tasks = mine - 17007
.into_iter() - 17008
.map(|task| { - 17009
let next = task - 17010
.schedule - 17011
.as_deref() - 17012
.and_then(|expr| { - 17013
next_fire - 17014
.get(&task.id) - 17015
.map(|at| at.with_timezone(&Utc)) - 17016
.or_else(|| { - 17017
vak_core::tasks::cron_next_after(expr, now) - 17018
.ok() - 17019
.map(|at| at.with_timezone(&Utc)) - 17020
}) - 17021
}) - 17022
.or_else(|| { - 17023
task.last_run_at - 17024
.map(|last| last + chrono::Duration::seconds(task.interval_secs as i64)) - 17025
}) - 17026
.or_else(|| Some(now.with_timezone(&Utc))); - 17027
let mut value = serde_json::to_value(task).unwrap_or_else(|_| serde_json::json!({})); - 17028
if let Some(object) = value.as_object_mut() { - 17029
object.insert( - 17030
"next_run_at".into(), - 17031
next.map(|at| serde_json::Value::String(at.to_rfc3339())) - 17032
.unwrap_or(serde_json::Value::Null), - 17033
); - 17034
object.insert( - 17035
"timezone".into(), - 17036
serde_json::Value::String(now.offset().to_string()), - 17037
); - 17038
} - 17039
value - 17040
}) - 17041
.collect::<Vec<_>>(); - 17042
Json(serde_json::json!({ "tasks": tasks })) - 17043
} - 17044
- 17045
#[derive(serde::Deserialize)] - 17046
struct TaskCreateBody { - 17047
name: String, - 17048
/// LLM turn instruction. Optional only for `script:` watchdog tasks. - 17049
#[serde(default)] - 17050
prompt: String, - 17051
#[serde(default = "task_default_interval")] - 17052
interval_secs: u64, - 17053
#[serde(default)] - 17054
deliver_to: Option<String>, - 17055
/// 5-field cron (`m h dom mon dow`, local time) replacing interval ticks. - 17056
#[serde(default)] - 17057
schedule: Option<String>, - 17058
#[serde(default)] - 17059
timezone: Option<String>, - 17060
#[serde(default)] - 17061
due_at: Option<chrono::DateTime<chrono::Utc>>, - 17062
/// Watchdog shell one-liner; XOR with `prompt`, never touches the LLM. - 17063
#[serde(default)] - 17064
script: Option<String>, - 17065
/// Pin dispatches to one model id (`provider/model` or bare model id). - 17066
#[serde(default)] - 17067
model_pin: Option<String>, - 17068
#[serde(default)] - 17069
agent_id: Option<String>, - 17070
#[serde(default)] - 17071
agent_revision: Option<u64>, - 17072
} - 17073
- 17074
fn task_default_interval() -> u64 { - 17075
3600 - 17076
} - 17077
- 17078
/// Structural validation shared by POST and PATCH: TaskDef::validate owns - 17079
/// the prompt-XOR-script and cron-grammar rules; the server adds its - 17080
/// transport-shape rules on top. Returns a typed 400 payload on failure. - 17081
fn validate_task_fields(task: &TaskDef) -> Result<(), (StatusCode, serde_json::Value)> { - 17082
if task.deliver_to.as_deref().is_some_and(|t| !t.contains(':')) { - 17083
return Err(( - 17084
StatusCode::BAD_REQUEST, - 17085
serde_json::json!({ - 17086
"error": "deliver_to must be '<surface>:<chat>', e.g. 'log:ops'" - 17087
}), - 17088
)); - 17089
} - 17090
task.validate().map_err(|e| { - 17091
( - 17092
StatusCode::BAD_REQUEST, - 17093
serde_json::json!({ "error": e.to_string() }), - 17094
) - 17095
}) - 17096
} - 17097
- 17098
async fn create_task( - 17099
State(state): State<AppState>, - 17100
Json(body): Json<TaskCreateBody>, - 17101
) -> axum::response::Response { - 17102
use axum::response::IntoResponse; - 17103
let scheduled = body.schedule.is_some(); - 17104
if !scheduled && body.interval_secs < 60 { - 17105
return ( - 17106
StatusCode::BAD_REQUEST, - 17107
Json(serde_json::json!({ "error": "interval must be >= 60s" })), - 17108
) - 17109
.into_response(); - 17110
} - 17111
let task = TaskDef { - 17112
id: uuid::Uuid::now_v7().to_string(), - 17113
name: body.name, - 17114
prompt: body.prompt, - 17115
interval_secs: body.interval_secs, - 17116
enabled: true, - 17117
cwd: state.core.cwd().clone(), - 17118
created_at: chrono::Utc::now(), - 17119
last_run_at: None, - 17120
last_session_id: None, - 17121
last_summary: None, - 17122
last_result_id: None, - 17123
last_run_status: None, - 17124
last_delivery_state: None, - 17125
last_wt: None, - 17126
deliver_to: body.deliver_to, - 17127
schedule: body.schedule.filter(|s| !s.trim().is_empty()), - 17128
timezone: body.timezone.filter(|s| !s.trim().is_empty()), - 17129
due_at: body.due_at, - 17130
script: body.script.filter(|s| !s.trim().is_empty()), - 17131
model_pin: body.model_pin.filter(|m| !m.trim().is_empty()), - 17132
agent_id: body.agent_id.filter(|m| !m.trim().is_empty()), - 17133
agent_revision: body.agent_revision, - 17134
}; - 17135
if let Err((status, payload)) = validate_task_fields(&task) { - 17136
return (status, Json(payload)).into_response(); - 17137
} - 17138
update_tasks(&state, |map| { - 17139
map.insert(task.id.clone(), task); - 17140
}); - 17141
(StatusCode::OK, Json(serde_json::json!({"ok": true}))).into_response() - 17142
} - 17143
- 17144
#[derive(serde::Deserialize)] - 17145
struct TaskPatchBody { - 17146
enabled: Option<bool>, - 17147
name: Option<String>, - 17148
prompt: Option<String>, - 17149
interval_secs: Option<u64>, - 17150
deliver_to: Option<Option<String>>, - 17151
/// Tri-state: absent = keep, null/empty = clear, string = set. - 17152
#[serde(default)] - 17153
schedule: OptionalStr, - 17154
#[serde(default)] - 17155
timezone: OptionalStr, - 17156
#[serde(default)] - 17157
due_at: Option<Option<chrono::DateTime<chrono::Utc>>>, - 17158
#[serde(default)] - 17159
script: OptionalStr, - 17160
#[serde(default)] - 17161
model_pin: OptionalStr, - 17162
/// Tri-state Agent selection: absent = keep, null/empty = clear, string = set. - 17163
#[serde(default)] - 17164
agent_id: OptionalStr, - 17165
/// Absent = keep, null = clear, number = set. - 17166
#[serde(default)] - 17167
agent_revision: Option<Option<u64>>, - 17168
} - 17169
- 17170
/// Distinguishes an absent JSON field from an explicit `null` (which plain - 17171
/// `Option<Option<T>>` cannot: both deserialize to outer `None`). - 17172
#[derive(Debug, Clone, Default)] - 17173
enum OptionalStr { - 17174
#[default] - 17175
Keep, - 17176
Clear, - 17177
Set(String), - 17178
} - 17179
- 17180
impl<'de> serde::Deserialize<'de> for OptionalStr { - 17181
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> - 17182
where - 17183
D: serde::Deserializer<'de>, - 17184
{ - 17185
// Only reached when the key is present; null → None here. - 17186
match Option::<String>::deserialize(deserializer)? { - 17187
None => Ok(OptionalStr::Clear), - 17188
Some(s) if s.trim().is_empty() => Ok(OptionalStr::Clear), - 17189
Some(s) => Ok(OptionalStr::Set(s)), - 17190
} - 17191
} - 17192
} - 17193
- 17194
async fn patch_task( - 17195
State(state): State<AppState>, - 17196
Path(id): Path<String>, - 17197
Json(body): Json<TaskPatchBody>, - 17198
) -> axum::response::Response { - 17199
use axum::response::IntoResponse; - 17200
if let Some(Some(t)) = &body.deliver_to - 17201
&& !t.contains(':') - 17202
{ - 17203
return ( - 17204
StatusCode::BAD_REQUEST, - 17205
Json(serde_json::json!({ - 17206
"error": "deliver_to must be '<surface>:<chat>', e.g. 'log:ops'" - 17207
})), - 17208
) - 17209
.into_response(); - 17210
} - 17211
// Mutate, validate and (on success) persist under a single hold of the - 17212
// tasks lock, so the patch can never be observed as committed in memory - 17213
// but not yet on disk (or vice versa) by a concurrent `load_tasks` tick. - 17214
// `update_tasks` can't itself carry an early `return` out of this async - 17215
// fn, so the closure reports outcome via `Result` and the response is - 17216
// built from that afterward. - 17217
let outcome = update_tasks( - 17218
&state, - 17219
|map| -> Result<TaskDef, (StatusCode, serde_json::Value)> { - 17220
let t = map.get_mut(&id).ok_or_else(|| { - 17221
( - 17222
StatusCode::NOT_FOUND, - 17223
serde_json::json!({ "error": format!("no task '{id}'") }), - 17224
) - 17225
})?; - 17226
// Apply to a candidate and validate BEFORE committing so a - 17227
// rejected patch never leaves half-mutated state behind. - 17228
let mut candidate = t.clone(); - 17229
if let Some(v) = body.enabled { - 17230
candidate.enabled = v; - 17231
} - 17232
if let Some(v) = body.name { - 17233
candidate.name = v; - 17234
} - 17235
if let Some(v) = body.prompt { - 17236
candidate.prompt = v; - 17237
} - 17238
if let Some(v) = body.interval_secs - 17239
&& v >= 60 - 17240
{ - 17241
candidate.interval_secs = v; - 17242
} - 17243
if let Some(v) = body.deliver_to { - 17244
candidate.deliver_to = v; - 17245
} - 17246
match body.schedule { - 17247
OptionalStr::Keep => {} - 17248
OptionalStr::Clear => candidate.schedule = None, - 17249
OptionalStr::Set(ref s) => candidate.schedule = Some(s.clone()), - 17250
} - 17251
match body.timezone { - 17252
OptionalStr::Keep => {} - 17253
OptionalStr::Clear => candidate.timezone = None, - 17254
OptionalStr::Set(ref s) => candidate.timezone = Some(s.clone()), - 17255
} - 17256
if let Some(v) = body.due_at { - 17257
candidate.due_at = v; - 17258
} - 17259
match body.script { - 17260
OptionalStr::Keep => {} - 17261
OptionalStr::Clear => candidate.script = None, - 17262
OptionalStr::Set(ref s) => candidate.script = Some(s.clone()), - 17263
} - 17264
match body.model_pin { - 17265
OptionalStr::Keep => {} - 17266
OptionalStr::Clear => candidate.model_pin = None, - 17267
OptionalStr::Set(ref s) => candidate.model_pin = Some(s.clone()), - 17268
} - 17269
match body.agent_id { - 17270
OptionalStr::Keep => {} - 17271
OptionalStr::Clear => { - 17272
candidate.agent_id = None; - 17273
candidate.agent_revision = None; - 17274
} - 17275
OptionalStr::Set(ref s) => candidate.agent_id = Some(s.clone()), - 17276
} - 17277
if let Some(revision) = body.agent_revision { - 17278
candidate.agent_revision = revision; - 17279
} - 17280
if let Err((status, payload)) = validate_task_fields(&candidate) { - 17281
return Err((status, payload)); - 17282
} - 17283
*t = candidate.clone(); - 17284
// Re-enabling reschedules interval tasks from now; dropping the - 17285
// cron marker makes the next tick recompute the schedule from - 17286
// scratch. - 17287
if body.enabled == Some(true) && t.schedule.is_none() { - 17288
t.last_run_at = None; - 17289
} - 17290
Ok(candidate) - 17291
}, - 17292
); - 17293
let updated = match outcome { - 17294
Ok(updated) => updated, - 17295
// `update_tasks` still writes tasks.json on the Err path (it can't - 17296
// see into the Result), but the write reproduces the same - 17297
// unmodified map, so a rejected/missing patch persists nothing new. - 17298
Err((status, payload)) => return (status, Json(payload)).into_response(), - 17299
}; - 17300
state - 17301
.next_fire - 17302
.lock() - 17303
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17304
.remove(&id); - 17305
(StatusCode::OK, Json(serde_json::json!(updated))).into_response() - 17306
} - 17307
- 17308
async fn delete_task(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 17309
// Remove-and-persist under one lock hold, closing the window where a - 17310
// concurrent scheduler `load_tasks` tick could otherwise re-read the - 17311
// not-yet-updated disk file and resurrect the task right after this - 17312
// handler releases the lock but before it writes tasks.json. - 17313
let removed = update_tasks(&state, |map| map.remove(&id)); - 17314
if let Some(task) = removed { - 17315
let idle = task.last_session_id.as_deref().is_none_or(|sid| { - 17316
state.get(sid).is_none_or(|h| { - 17317
h.session - 17318
.lock() - 17319
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17320
.is_some() - 17321
}) - 17322
}); - 17323
if idle && let Some(wt) = &task.last_wt { - 17324
let old = vak_core::worktree::Worktree { - 17325
path: wt.path.clone(), - 17326
branch: wt.branch.clone(), - 17327
}; - 17328
let _ = vak_core::worktree::remove(&task.cwd, &old); - 17329
} - 17330
StatusCode::OK - 17331
} else { - 17332
StatusCode::NOT_FOUND - 17333
} - 17334
} - 17335
- 17336
/// Fire a task immediately (also resets its schedule). - 17337
async fn run_task_now(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 17338
match fire_task(&state, &id).await { - 17339
Ok(_) => return StatusCode::ACCEPTED, - 17340
Err(NotFired::Gone) => return StatusCode::NOT_FOUND, - 17341
Err(NotFired::Refused) => return StatusCode::UNPROCESSABLE_ENTITY, - 17342
Err(NotFired::Busy) => {} - 17343
} - 17344
// A scheduler tick may hold the one-shot inflight slot for this script - 17345
// task — the requested execution is happening at this very moment, so - 17346
// report accepted rather than conflict (found by the 0.7 suite: the - 17347
// tick raced run-now on freshly created interval tasks). - 17348
let busy_elsewhere = state - 17349
.tasks - 17350
.lock() - 17351
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17352
.get(&id) - 17353
.map(|t| { - 17354
t.script - 17355
.as_deref() - 17356
.map(str::trim) - 17357
.is_some_and(|s| !s.is_empty()) - 17358
}) - 17359
.unwrap_or(false) - 17360
&& state - 17361
.script_inflight - 17362
.lock() - 17363
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17364
.contains(&id); - 17365
if busy_elsewhere { - 17366
StatusCode::ACCEPTED - 17367
} else { - 17368
StatusCode::CONFLICT - 17369
} - 17370
} - 17371
- 17372
/// Replay pending deliveries for one task without executing the task again. - 17373
async fn retry_task_delivery( - 17374
State(state): State<AppState>, - 17375
Path(id): Path<String>, - 17376
) -> axum::response::Response { - 17377
use axum::response::IntoResponse; - 17378
let records = match delivery::outbox_records(&state.core) { - 17379
Ok(records) => records, - 17380
Err(error) => { - 17381
return ( - 17382
StatusCode::INTERNAL_SERVER_ERROR, - 17383
Json(serde_json::json!({ "error": error })), - 17384
) - 17385
.into_response(); - 17386
} - 17387
}; - 17388
let jobs = records - 17389
.into_iter() - 17390
.filter(|record| record.state == vak_delivery::outbox::OutboxState::Pending) - 17391
.filter(|record| match &record.job.content { - 17392
vak_delivery::DeliveryContent::Answer(answer) => { - 17393
answer.metadata.get("vak_task_id").map(String::as_str) == Some(id.as_str()) - 17394
} - 17395
_ => false, - 17396
}) - 17397
.map(|record| record.job.job_id) - 17398
.collect::<Vec<_>>(); - 17399
let mut replayed = 0usize; - 17400
let mut failed = 0usize; - 17401
for job_id in jobs { - 17402
match delivery::replay_outbox_job(&state.core, &job_id).await { - 17403
Ok(()) => replayed += 1, - 17404
Err(_) => failed += 1, - 17405
} - 17406
} - 17407
if replayed > 0 && failed == 0 { - 17408
update_tasks(&state, |tasks| { - 17409
if let Some(task) = tasks.get_mut(&id) { - 17410
task.last_delivery_state = Some("delivered".into()); - 17411
} - 17412
}); - 17413
} - 17414
Json(serde_json::json!({ "replayed": replayed, "failed": failed })).into_response() - 17415
} - 17416
- 17417
/// Split a `model_pin` into (provider, model). A bare model id pins only - 17418
/// the model and keeps this server's active provider. - 17419
pub(crate) fn split_model_pin(pin: &str, current_provider: &str) -> (String, String) { - 17420
match pin.split_once('/') { - 17421
Some((provider, model)) if !provider.trim().is_empty() && !model.trim().is_empty() => { - 17422
(provider.trim().to_string(), model.trim().to_string()) - 17423
} - 17424
_ => (current_provider.to_string(), pin.trim().to_string()), - 17425
} - 17426
} - 17427
- 17428
/// Spawn one isolated run for `task` if its previous run is idle. Returns - 17429
/// the child session id on success. Script tasks take the brokered-bash - 17430
/// branch instead: no provider dispatch, no worktree, no child session. - 17431
/// Why a due task did not start a run. - 17432
enum NotFired { - 17433
/// The task no longer exists. - 17434
Gone, - 17435
/// Its previous run is still going; the slot waits for the next tick. - 17436
Busy, - 17437
/// It cannot run as it stands; the reason is in the inbox. - 17438
Refused, - 17439
} - 17440
- 17441
/// Tells the person why a due task could not start, once per missed slot: - 17442
/// the scheduler retries the slot every tick, and the key changes only when - 17443
/// the task has run since. - 17444
fn refuse_task(state: &AppState, task: &TaskDef, reason: String) -> NotFired { - 17445
let slot = task.last_run_at.unwrap_or(task.created_at).to_rfc3339(); - 17446
let key = format!("routine-failed|{}|{slot}", task.id); - 17447
let _ = vak_core::inbox::record_with_result_and_key( - 17448
&state.core.shared_data_home(), - 17449
vak_core::inbox::Kind::RoutineFailed, - 17450
&format!("routine '{}' could not run", task.name), - 17451
&reason, - 17452
None, - 17453
Some(&task.id), - 17454
None, - 17455
Some(&key), - 17456
); - 17457
update_tasks(state, |map| { - 17458
if let Some(t) = map.get_mut(&task.id) { - 17459
t.last_run_status = Some("refused".into()); - 17460
t.last_summary = Some(reason.clone()); - 17461
} - 17462
}); - 17463
NotFired::Refused - 17464
} - 17465
- 17466
async fn fire_task(state: &AppState, id: &str) -> Result<String, NotFired> { - 17467
let Some(snapshot) = state - 17468
.tasks - 17469
.lock() - 17470
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17471
.get(id) - 17472
.cloned() - 17473
else { - 17474
return Err(NotFired::Gone); - 17475
}; - 17476
// Previous run still going? - 17477
if let Some(prev) = snapshot.last_session_id.as_deref() - 17478
&& state.get(prev).is_some_and(|h| { - 17479
h.session - 17480
.lock() - 17481
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17482
.is_none() - 17483
}) - 17484
{ - 17485
return Err(NotFired::Busy); - 17486
} - 17487
let script = snapshot - 17488
.script - 17489
.as_deref() - 17490
.map(str::trim) - 17491
.filter(|s| !s.is_empty()); - 17492
if let Some(script) = script { - 17493
return fire_script_task(state, &snapshot, script) - 17494
.await - 17495
.ok_or(NotFired::Busy); - 17496
} - 17497
let provider = match state.core.provider() { - 17498
Ok(provider) => provider, - 17499
Err(error) => { - 17500
return Err(refuse_task( - 17501
state, - 17502
&snapshot, - 17503
format!( - 17504
"No model is available to run it ({error}). Connect a provider in Settings; the routine runs at its next check." - 17505
), - 17506
)); - 17507
} - 17508
}; - 17509
// A scheduled run works in its own git worktree of the space, so it - 17510
// cannot run in a folder that is not a repository. - 17511
if !vak_core::worktree::is_git_repo(&snapshot.cwd) { - 17512
return Err(refuse_task( - 17513
state, - 17514
&snapshot, - 17515
format!( - 17516
"{} is not a git repository, and a scheduled run works in its own copy of one. Run `git init` there and commit, or move the routine to a folder that is a repository.", - 17517
snapshot.cwd.display() - 17518
), - 17519
)); - 17520
} - 17521
// Drop the previous worktree (latest-only retention). - 17522
if let Some(wt) = &snapshot.last_wt { - 17523
let old = vak_core::worktree::Worktree { - 17524
path: wt.path.clone(), - 17525
branch: wt.branch.clone(), - 17526
}; - 17527
let _ = vak_core::worktree::remove(&snapshot.cwd, &old); - 17528
} - 17529
let rid = format!("task-{}", uuid::Uuid::now_v7().simple()); - 17530
let wt = match vak_core::worktree::create(&snapshot.cwd, &rid) { - 17531
Ok(wt) => wt, - 17532
Err(error) => { - 17533
return Err(refuse_task( - 17534
state, - 17535
&snapshot, - 17536
format!("Its working copy could not be made: {error}."), - 17537
)); - 17538
} - 17539
}; - 17540
let fired_at_utc = chrono::Utc::now(); - 17541
let scheduled_prompt = format!( - 17542
"{}\n\n[Scheduled-run context: fired at UTC {}; local system time {}. Re-evaluate relative dates against this run time unless the request explicitly established a specific date.]", - 17543
snapshot.prompt, - 17544
fired_at_utc.to_rfc3339(), - 17545
fired_at_utc.with_timezone(&chrono::Local).to_rfc3339(), - 17546
); - 17547
let child_id = spawn_isolated_run( - 17548
state, - 17549
provider.clone(), - 17550
&wt, - 17551
&scheduled_prompt, - 17552
snapshot.model_pin.as_deref(), - 17553
snapshot.agent_id.as_deref(), - 17554
snapshot.agent_revision, - 17555
false, - 17556
) - 17557
.await - 17558
.map_err(|error| { - 17559
let _ = vak_core::worktree::remove(&snapshot.cwd, &wt); - 17560
refuse_task(state, &snapshot, format!("It could not start: {error}.")) - 17561
})?; - 17562
- 17563
update_tasks(state, |map| { - 17564
if let Some(t) = map.get_mut(id) { - 17565
if t.due_at.is_some() { - 17566
t.enabled = false; - 17567
} - 17568
t.last_run_at = Some(chrono::Utc::now()); - 17569
t.last_session_id = Some(child_id.clone()); - 17570
t.last_summary = None; - 17571
t.last_result_id = None; - 17572
t.last_run_status = Some("working".into()); - 17573
t.last_delivery_state = Some("pending".into()); - 17574
t.last_wt = Some(WtMeta { - 17575
path: wt.path, - 17576
branch: wt.branch, - 17577
}); - 17578
} - 17579
}); - 17580
- 17581
// Watcher: record the run's final assistant text on the task when it - 17582
// finishes, and push it out through the gateway when a deliver target - 17583
// is set. The event's `summary` is a status word; the transcript holds - 17584
// the actual answer a phone user should receive. - 17585
if let Some(h) = state.get(&child_id) { - 17586
let st = state.clone(); - 17587
let tid = id.to_string(); - 17588
let child_handle = h.clone(); - 17589
let child_session = child_id.clone(); - 17590
let task_name = snapshot.name.clone(); - 17591
let deliver_to = snapshot.deliver_to.clone(); - 17592
let rx = h.events_tx.subscribe(); - 17593
tokio::spawn(async move { - 17594
use tokio_stream::StreamExt; - 17595
use tokio_stream::wrappers::BroadcastStream; - 17596
let mut stream = BroadcastStream::new(rx); - 17597
while let Some(Ok(ev)) = stream.next().await { - 17598
if let AgentEvent::RunFinished { summary, is_error } = ev.event { - 17599
let text = - 17600
last_assistant_text(&child_handle).unwrap_or_else(|| summary.clone()); - 17601
let result_id = child_handle - 17602
.presentation - 17603
.lock() - 17604
.unwrap_or_else(std::sync::PoisonError::into_inner) - 17605
.items - 17606
.iter() - 17607
.rev() - 17608
.find_map(|item| { - 17609
item.outcome - 17610
.as_ref() - 17611
.map(|outcome| outcome.result_id.clone()) - 17612
}); - 17613
update_tasks(&st, |map| { - 17614
if let Some(t) = map.get_mut(&tid) { - 17615
t.last_summary = Some(text.clone()); - 17616
t.last_result_id = result_id.clone(); - 17617
} - 17618
}); - 17619
let delivery_state = if let Some(target) = &deliver_to { - 17620
// Delivery failure must not lose the recorded summary; - 17621
// it only means this transport could not be reached. - 17622
gateway::deliver_and_record_with_result( - 17623
&st.core, - 17624
target, - 17625
&format!("routine '{task_name}' finished:\n{text}"), - 17626
vak_core::inbox::Kind::TaskSummary, - 17627
format!("routine '{task_name}' finished"), - 17628
Some(&child_session), - 17629
Some(&tid), - 17630
result_id.as_deref(), - 17631
) - 17632
.await - 17633
.unwrap_or("pending") - 17634
} else { - 17635
let dedupe_key = Some(format!("inbox|{child_session}")); - 17636
let _ = vak_core::inbox::record_with_result_and_key( - 17637
&st.core.shared_data_home(), - 17638
vak_core::inbox::Kind::TaskSummary, - 17639
&format!("routine '{task_name}' finished"), - 17640
&format!("routine '{task_name}' finished:\n{text}"), - 17641
Some(&child_session), - 17642
Some(&tid), - 17643
result_id.as_deref(), - 17644
dedupe_key.as_deref(), - 17645
); - 17646
"inbox" - 17647
}; - 17648
check_budget_alert(&st, &tid).await; - 17649
update_tasks(&st, |map| {
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.