- 16464
}; - 16465
let child_core = vak_core::Core::new_with_trust(wt.path.clone(), true) - 16466
.map(|c| { - 16467
c.with_agent_identity(identity) - 16468
.with_surface(vak_core::Surface::Background) - 16469
// Unattended, and stamped BEFORE `start_session` composes and - 16470
// freezes the prompt. Stamping afterwards would be too late: - 16471
// the prompt would already have advertised a gated capability - 16472
// that this run can only ever be refused, which is the exact - 16473
// mismatch `vak_core::reach` exists to remove. `begin_turn` - 16474
// installs the matching approver. - 16475
.with_approver_answerable(false) - 16476
}) - 16477
.map_err(|e| format!("child core failed: {e}"))?; - 16478
child_core.set_provider_instance(provider); - 16479
// The shared root: the child resolves its own Agent's home beneath it, - 16480
// as every Core does. Seeding it with this Core's (already Agent-scoped) - 16481
// home nested one Agent's home inside another's. - 16482
child_core.set_sessions_home(state.core.shared_data_home()); - 16483
if let Some(pin) = model_pin.map(str::trim).filter(|p| !p.is_empty()) { - 16484
let (pin_provider, pin_model) = split_model_pin(pin, &child_core.effective_provider()); - 16485
child_core.set_route(pin_provider, pin_model); - 16486
} - 16487
- 16488
let child_log = child_core - 16489
.start_session() - 16490
.await - 16491
.map_err(|_| "child session failed to start".to_string())?; - 16492
let Some(child_header) = child_log.header() else { - 16493
return Err("child session has no header".to_string()); - 16494
}; - 16495
// The handle is the ledger's own id, so a run recorded on a task - 16496
// (`last_session_id`) opens from disk after a restart. - 16497
let child_id = child_header.session_id.clone(); - 16498
let handle = register_handle( - 16499
state, - 16500
child_id.clone(), - 16501
child_log, - 16502
wt.path.clone(), - 16503
child_core.clone(), - 16504
); - 16505
if start_turn { - 16506
begin_turn(&handle, &child_core, prompt, false); - 16507
} - 16508
Ok(child_id) - 16509
} - 16510
- 16511
/// Fire a single-turn agent run on a (usually fresh) session handle. - 16512
/// - 16513
/// `attended` says whether anyone is watching this run's event stream. It is - 16514
/// not cosmetic: a scheduled routine and a best-of-N leg both arrive here, - 16515
/// nobody is subscribed to either, and an approval gate raised on one used - 16516
/// to emit an SSE event into the void and then block the run until the - 16517
/// process restarted. An unattended run gets an approver that says so, and - 16518
/// `Core::with_approver` carries that fact into the prompt so the model is - 16519
/// never offered a capability whose gate can only ever be refused. - 16520
fn begin_turn(handle: &Arc<SessionHandle>, core: &Core, prompt: &str, attended: bool) { - 16521
let approver: Arc<dyn Approver> = Arc::new(HttpApprover { - 16522
events_tx: handle.events_tx.clone(), - 16523
pending: handle.pending.clone(), - 16524
session_id: handle.id.clone(), - 16525
activity_buffer: handle.activity_buffer.clone(), - 16526
answerable: attended, - 16527
}); - 16528
let core = core.clone().with_approver(approver.as_ref()); - 16529
let events = mpsc_to_broadcast(handle.events_tx.clone()); - 16530
let steering = Arc::new(SteeringQueues::new()); - 16531
let cancel = handle - 16532
.cancel - 16533
.lock() - 16534
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16535
.clone(); - 16536
let Some(log) = handle - 16537
.session - 16538
.lock() - 16539
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16540
.take() - 16541
else { - 16542
return; // busy — caller should have checked - 16543
}; - 16544
let turn_session_id = log - 16545
.header() - 16546
.map(|h| h.session_id.clone()) - 16547
.unwrap_or_default(); - 16548
let core = core.clone(); - 16549
let prompt = prompt.to_string(); - 16550
let h2 = handle.clone(); - 16551
tokio::spawn(async move { - 16552
let outcome = core - 16553
.run_turn_with( - 16554
log, - 16555
&prompt, - 16556
cancel, - 16557
Some(approver), - 16558
None, - 16559
Some(steering.clone()), - 16560
events, - 16561
) - 16562
.await; - 16563
*h2.cancel - 16564
.lock() - 16565
.unwrap_or_else(std::sync::PoisonError::into_inner) = CancellationToken::new(); - 16566
match outcome { - 16567
Ok((_, mut restored)) => { - 16568
for activity in std::mem::take( - 16569
&mut *h2 - 16570
.activity_buffer - 16571
.lock() - 16572
.unwrap_or_else(std::sync::PoisonError::into_inner), - 16573
) { - 16574
let _ = restored.append_activity(activity); - 16575
} - 16576
*h2.presentation - 16577
.lock() - 16578
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 16579
crate::projection::snapshot(&turn_session_id, &restored); - 16580
*h2.session - 16581
.lock() - 16582
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(restored); - 16583
let _ = h2.events_tx.send(AgentEvent::RunFinished { - 16584
summary: "completed".into(), - 16585
is_error: false, - 16586
}); - 16587
} - 16588
Err(e) => { - 16589
// Same leak class as side chats: restore from the durable - 16590
// ledger so the handle is not wedged on "run in progress". - 16591
if let Some(log) = reopen_ledger(&core, &turn_session_id) { - 16592
*h2.session - 16593
.lock() - 16594
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(log); - 16595
} - 16596
let _ = h2.events_tx.send(AgentEvent::RunFinished { - 16597
summary: format!("error: {e}"), - 16598
is_error: true, - 16599
}); - 16600
} - 16601
} - 16602
drop(steering); - 16603
}); - 16604
} - 16605
- 16606
async fn keep_best_run( - 16607
State(state): State<AppState>, - 16608
Path(id): Path<String>, - 16609
) -> axum::response::Response { - 16610
use axum::response::IntoResponse; - 16611
let meta = state - 16612
.best_runs - 16613
.lock() - 16614
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16615
.get(&id) - 16616
.cloned(); - 16617
let Some(meta) = meta else { - 16618
return StatusCode::NOT_FOUND.into_response(); - 16619
}; - 16620
// Merge into the user's checkout. Conflicts/dirty trees surface as errors. - 16621
let out = tokio::process::Command::new("git") - 16622
.args([ - 16623
"merge", - 16624
"--no-ff", - 16625
"-m", - 16626
&format!("best-of-n: merge {}", meta.branch), - 16627
]) - 16628
.arg(&meta.branch) - 16629
.current_dir(&meta.repo) - 16630
.output() - 16631
.await; - 16632
match out { - 16633
Ok(o) if o.status.success() => { - 16634
cleanup_worktree(&state, &id, &meta); - 16635
(StatusCode::OK, Json(serde_json::json!({"kept": id}))).into_response() - 16636
} - 16637
Ok(o) => { - 16638
// Abort any conflicted merge so the tree is not left dirty. - 16639
let _ = tokio::process::Command::new("git") - 16640
.args(["merge", "--abort"]) - 16641
.current_dir(&meta.repo) - 16642
.output() - 16643
.await; - 16644
( - 16645
StatusCode::CONFLICT, - 16646
Json(serde_json::json!({ - 16647
"error": "merge failed", - 16648
"stderr": String::from_utf8_lossy(&o.stderr), - 16649
})), - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.