- 16164
text = format!("{}/{}", sub.trim_end_matches('/'), text); - 16165
} - 16166
files.push(text); - 16167
if files.len() >= limit { - 16168
break; - 16169
} - 16170
} - 16171
files.sort_unstable(); - 16172
( - 16173
StatusCode::OK, - 16174
Json(serde_json::json!({ "files": files, "truncated": files.len() >= limit })), - 16175
) - 16176
.into_response() - 16177
} - 16178
- 16179
#[derive(serde::Deserialize)] - 16180
struct SideBody { - 16181
question: String, - 16182
} - 16183
- 16184
/// `/btw`: ask a question using the session's context WITHOUT landing it on - 16185
/// the main chain. Mechanics: append the Q + run the turn as a sibling - 16186
/// branch (parent = current main tail), then restore the tail so future - 16187
/// main turns continue exactly where they were. The side entries stay in - 16188
/// the ledger — reconstructable, never deleted. - 16189
async fn side_chat( - 16190
State(state): State<AppState>, - 16191
Path(id): Path<String>, - 16192
Json(body): Json<SideBody>, - 16193
) -> axum::response::Response { - 16194
use axum::response::IntoResponse; - 16195
let Some(handle) = state.get(&id) else { - 16196
return StatusCode::NOT_FOUND.into_response(); - 16197
}; - 16198
let Some(mut taken) = handle - 16199
.session - 16200
.lock() - 16201
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16202
.take() - 16203
else { - 16204
return StatusCode::CONFLICT.into_response(); // main run active - 16205
}; - 16206
if let Err(e) = state.core.provider() { - 16207
*handle - 16208
.session - 16209
.lock() - 16210
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 16211
return provider_unavailable(e); - 16212
} - 16213
- 16214
let tail_main = taken.tail_id().cloned(); - 16215
if let Err(_e) = taken.append_message(vak_session::MessageRecord { - 16216
message: vak_llm::Message::user_text(body.question.clone()), - 16217
meta: None, - 16218
}) { - 16219
*handle - 16220
.session - 16221
.lock() - 16222
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 16223
return StatusCode::INTERNAL_SERVER_ERROR.into_response(); - 16224
} - 16225
let side_tx = handle.side_events_tx.clone(); - 16226
let side_activity = Arc::new(Mutex::new(Vec::new())); - 16227
let approver: Arc<dyn Approver> = Arc::new(HttpApprover { - 16228
events_tx: handle.events_tx.clone(), - 16229
pending: handle.pending.clone(), - 16230
session_id: handle.id.clone(), - 16231
activity_buffer: side_activity.clone(), - 16232
answerable: true, - 16233
}); - 16234
let events = mpsc_to_broadcast(side_tx.clone()); - 16235
let cancel = handle - 16236
.side_cancel - 16237
.lock() - 16238
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16239
.clone(); - 16240
// A side chat reads this session's context, so it runs under the same - 16241
// ceiling the session was created with. - 16242
let core = handle.core.clone(); - 16243
- 16244
tokio::spawn(async move { - 16245
// No steering on side chats by design: they are read-only Q&A over - 16246
// the session context, not a second control surface. - 16247
let outcome = core - 16248
.run_turn_with( - 16249
taken, - 16250
&body.question, - 16251
cancel, - 16252
Some(approver), - 16253
None, - 16254
None, - 16255
events, - 16256
) - 16257
.await; - 16258
let (summary, is_error) = match &outcome { - 16259
Ok((vak_agent::TurnOutcome::Completed { .. }, _)) => ("completed".to_string(), false), - 16260
Ok((vak_agent::TurnOutcome::Aborted { .. }, _)) => ("aborted".to_string(), false), - 16261
Ok((_, _)) => ("ended".to_string(), false), - 16262
Err(e) => (format!("error: {e}"), true), - 16263
}; - 16264
let turn_ok = matches!(&outcome, Ok((_, _))); - 16265
if let Ok((_, mut restored)) = outcome { - 16266
for activity in std::mem::take( - 16267
&mut *side_activity - 16268
.lock() - 16269
.unwrap_or_else(std::sync::PoisonError::into_inner), - 16270
) { - 16271
let _ = restored.append_activity(activity); - 16272
} - 16273
// Rewind the branch pointer to the main line: the side entries - 16274
// remain in the ledger as a sibling branch — reconstructable via - 16275
// their parent chain, invisible to derive_messages(). - 16276
if let Some(main_tail) = &tail_main { - 16277
let _ = restored.branch_at(main_tail); - 16278
} - 16279
*handle - 16280
.presentation - 16281
.lock() - 16282
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 16283
crate::projection::snapshot(&id, &restored); - 16284
*handle - 16285
.session - 16286
.lock() - 16287
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(restored); - 16288
} - 16289
let _ = side_tx.send(AgentEvent::RunFinished { summary, is_error }); - 16290
// On a failed turn the taken log is gone with the Err — reopen the - 16291
// durable ledger so the session does not stay wedged as - 16292
// "run in progress" forever (found by the v0.6 deployment gate). - 16293
if !turn_ok && let Some(log) = reopen_ledger(&core, &id) { - 16294
*handle - 16295
.session - 16296
.lock() - 16297
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(log); - 16298
} - 16299
}); - 16300
- 16301
StatusCode::ACCEPTED.into_response() - 16302
} - 16303
- 16304
async fn side_cancel_run(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 16305
let Some(handle) = state.get(&id) else { - 16306
return StatusCode::NOT_FOUND; - 16307
}; - 16308
handle - 16309
.side_cancel - 16310
.lock() - 16311
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16312
.cancel(); - 16313
let _ = handle.side_events_tx.send(AgentEvent::RunFinished { - 16314
summary: "cancelled by client".into(), - 16315
is_error: false, - 16316
}); - 16317
StatusCode::ACCEPTED - 16318
} - 16319
- 16320
#[derive(serde::Deserialize)] - 16321
struct BestBody { - 16322
prompt: String, - 16323
n: Option<usize>, - 16324
} - 16325
- 16326
/// Best-of-N: fan the same prompt across N isolated git worktrees, each with - 16327
/// its own session + event stream. Candidates are compared by diff; `keep` - 16328
/// merges a branch, `discard` drops it. Ledger-native: every run is a normal - 16329
/// session under the shared store. - 16330
async fn start_bestofn( - 16331
State(state): State<AppState>, - 16332
Path(id): Path<String>, - 16333
Json(body): Json<BestBody>, - 16334
) -> axum::response::Response { - 16335
use axum::response::IntoResponse; - 16336
- 16337
let Some(anchor) = state.get(&id) else { - 16338
return StatusCode::NOT_FOUND.into_response(); - 16339
}; - 16340
if anchor - 16341
.session - 16342
.lock() - 16343
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16344
.is_none() - 16345
{ - 16346
return StatusCode::CONFLICT.into_response(); - 16347
} - 16348
let provider = match state.core.provider() { - 16349
Ok(p) => p, - 16350
Err(e) => return provider_unavailable(e), - 16351
}; - 16352
- 16353
let n = body.n.unwrap_or(2).clamp(1, 4); - 16354
let repo = state.core.cwd().clone(); - 16355
if !vak_core::worktree::is_git_repo(&repo) { - 16356
return StatusCode::CONFLICT.into_response(); - 16357
} - 16358
- 16359
// Create worktrees first; roll back everything on partial failure. - 16360
let mut created: Vec<(String, vak_core::worktree::Worktree)> = Vec::new(); - 16361
for i in 0..n { - 16362
// v7 shares its leading chars within one millisecond; disambiguate. - 16363
let rid = format!("{}-{i}", uuid::Uuid::now_v7().simple()); - 16364
match vak_core::worktree::create(&repo, &rid) { - 16365
Ok(wt) => created.push((rid, wt)), - 16366
Err(e) => { - 16367
for (_, wt) in &created { - 16368
let _ = vak_core::worktree::remove(&repo, wt); - 16369
} - 16370
return ( - 16371
StatusCode::INTERNAL_SERVER_ERROR, - 16372
Json(serde_json::json!({ "error": format!("worktree create failed: {e}") })), - 16373
) - 16374
.into_response(); - 16375
} - 16376
} - 16377
} - 16378
- 16379
let mut runs = Vec::new(); - 16380
for (_, wt) in &created { - 16381
match spawn_isolated_run( - 16382
&state, - 16383
provider.clone(), - 16384
wt, - 16385
&body.prompt, - 16386
None, - 16387
None, - 16388
None, - 16389
true, - 16390
) - 16391
.await - 16392
{ - 16393
Ok(child_id) => { - 16394
state - 16395
.best_runs - 16396
.lock() - 16397
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16398
.insert( - 16399
child_id.clone(), - 16400
BestRunMeta { - 16401
repo: repo.clone(), - 16402
wt_path: wt.path.clone(), - 16403
branch: wt.branch.clone(), - 16404
}, - 16405
); - 16406
runs.push(serde_json::json!({ - 16407
"session_id": child_id, - 16408
"branch": wt.branch, - 16409
"path": wt.path, - 16410
})); - 16411
} - 16412
Err(e) => { - 16413
for (_, w) in &created { - 16414
let _ = vak_core::worktree::remove(&repo, w); - 16415
} - 16416
return ( - 16417
StatusCode::INTERNAL_SERVER_ERROR, - 16418
Json(serde_json::json!({ "error": e })), - 16419
) - 16420
.into_response(); - 16421
} - 16422
} - 16423
} - 16424
- 16425
(StatusCode::OK, Json(serde_json::json!({ "runs": runs }))).into_response() - 16426
} - 16427
- 16428
/// One isolated run inside `wt`: child Core + session + registered handle + - 16429
/// fired turn. Shared by best-of-N and the task scheduler. `model_pin` - 16430
/// (docs/design/29-personal-os.md P2) overrides the child's provider/model - 16431
/// so BOTH main dispatches and any receipts carry the pinned id only — a - 16432
/// pinned task never escalates to another model. - 16433
#[allow(clippy::too_many_arguments)] - 16434
async fn spawn_isolated_run( - 16435
state: &AppState, - 16436
provider: Arc<dyn Provider>, - 16437
wt: &vak_core::worktree::Worktree, - 16438
prompt: &str, - 16439
model_pin: Option<&str>, - 16440
agent_id: Option<&str>, - 16441
agent_revision: Option<u64>, - 16442
start_turn: bool, - 16443
) -> Result<String, String> { - 16444
let identity = if let Some(agent_id) = agent_id { - 16445
let profiles = agents::effective(&state.active_core())?; - 16446
let profile = profiles - 16447
.iter() - 16448
.find(|profile| profile.id == agent_id) - 16449
.ok_or_else(|| format!("Agent '{agent_id}' no longer exists"))?; - 16450
if !profile.is_admissible() { - 16451
return Err(format!("Agent '{agent_id}' is paused or archived")); - 16452
} - 16453
if let Some(expected) = agent_revision - 16454
&& expected != profile.revision - 16455
{ - 16456
return Err(format!( - 16457
"Agent '{agent_id}' changed from revision {expected} to {}", - 16458
profile.revision - 16459
)); - 16460
} - 16461
Some(profile.identity()) - 16462
} else { - 16463
None - 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)]
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.