1use std::collections::BTreeSet;
25use std::path::{Path, PathBuf};
26
27use anyhow::{Context as _, Result, bail};
28use jiff::{SignedDuration, Timestamp};
29use serde::Deserialize;
30
31use crate::ask::Questions;
32use crate::config::Disk;
33use crate::run::{RunState, RunStatus, SCHEMA, short_of};
34
35use crate::disk::{Prune, dir_size};
36
37#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
39pub struct Housekeeping {
40 pub folded: usize,
42 pub unreadable: usize,
52 pub orphaned_worktrees: usize,
55 pub cache_files: usize,
57 pub cache_freed: u64,
59 pub questions_abandoned: usize,
62 pub external_merges_recorded: usize,
67 pub stale_pr_states_repaired: usize,
71}
72
73pub async fn housekeep(
81 cfg: &crate::config::Config,
82 home: &Path,
83 worktrees_root: &Path,
84 repo: &Path,
85 now: Timestamp,
86) -> Housekeeping {
87 let mut out = Housekeeping::default();
88 if cfg.disk.auto_fold {
89 let runs = home.join("runs");
90 match fold_due(&runs, home, worktrees_root, &cfg.disk, now).await {
91 Ok((folded, unreadable)) => {
92 out.folded = folded;
93 out.unreadable = unreadable;
94 }
95 Err(e) => {
96 tracing::warn!("housekeep: fold due runs: {e:#}");
97 crate::notices::raise_in(
100 home,
101 crate::notices::Notice::warn(
102 "housekeep:fold",
103 "Automatic cleanup of finished runs failed; disk usage may keep growing.",
104 ),
105 );
106 }
107 }
108 out.orphaned_worktrees =
109 fold_orphaned_worktrees(&runs, worktrees_root, home, cfg.disk.fold_grace_secs, now)
110 .await;
111 if let Err(e) = crate::git::worktree_prune(repo).await {
115 tracing::warn!("housekeep: prune worktree registrations: {e:#}");
116 }
117 out.external_merges_recorded = reconcile_external_merges(&runs, home, &cfg.disk, now).await;
118 out.stale_pr_states_repaired = crate::land::repair_stale_pr_states(home, 5).await.0;
121 }
122 match prune_cache_if_over_limit(cfg, home) {
123 Ok(Some(pruned)) => {
124 out.cache_files = pruned.files;
125 out.cache_freed = pruned.freed;
126 }
127 Ok(None) => {}
128 Err(e) => tracing::warn!("housekeep: prune cache: {e:#}"),
129 }
130 out.questions_abandoned =
138 abandon_settled_questions(&Questions::at(home.join("questions")), &home.join("runs"));
139 out
140}
141
142pub fn abandon_settled_questions(store: &Questions, runs: &Path) -> usize {
153 let waiting_on: BTreeSet<String> = store
154 .list()
155 .into_iter()
156 .filter(|q| q.status.open())
157 .map(|q| q.run)
158 .collect();
159 let mut abandoned = 0;
160 for run in waiting_on {
161 let Ok(meta) = read_meta(runs, &run) else {
162 continue;
163 };
164 match store.settle_run(&run, meta.status) {
165 Ok(n) => abandoned += n,
166 Err(e) => tracing::warn!("housekeep: abandon questions for {run}: {e:#}"),
167 }
168 }
169 abandoned
170}
171
172pub async fn fold_due(
197 runs: &Path,
198 home: &Path,
199 _worktrees_root: &Path,
200 disk: &Disk,
201 now: Timestamp,
202) -> Result<(usize, usize)> {
203 let mut folded = 0usize;
204 let mut unreadable = 0usize;
205 let mut ids: Vec<String> = std::fs::read_dir(runs)
206 .into_iter()
207 .flatten()
208 .flatten()
209 .filter(|e| e.path().join("run.json").is_file())
210 .map(|e| e.file_name().to_string_lossy().into_owned())
211 .collect();
212 ids.sort_unstable();
213 for id in ids {
214 if crate::daemon::is_working_on(home, &id, now) {
215 continue;
216 }
217 let meta = match read_meta(runs, &id) {
218 Ok(meta) => meta,
219 Err(e) => {
220 unreadable += 1;
221 tracing::warn!("housekeep: run {id} unreadable, left alone: {e:#}");
222 continue;
223 }
224 };
225 if meta.status.resumable() || !due(now, meta.updated_at, disk.fold_grace_secs) {
226 continue;
227 }
228 let mut state = match read_state(runs, &id) {
237 Ok(state) => state,
238 Err(e) => {
239 unreadable += 1;
240 tracing::warn!("housekeep: run {id} unreadable, left alone: {e:#}");
241 continue;
242 }
243 };
244 if state.schema != SCHEMA {
245 tracing::info!(
246 "housekeep: run {id} was written by schema {} (this build speaks {SCHEMA}); \
247 folding it anyway",
248 state.schema
249 );
250 }
251 let drop_winner = matches!(
261 state.status,
262 RunStatus::Merged | RunStatus::Superseded | RunStatus::AlreadyInBase
263 );
264 match crate::graph::fold_run(&mut state, drop_winner, home).await {
271 Ok(_) => folded += 1,
272 Err(e) => tracing::warn!("housekeep: fold {id}: {e:#}"),
273 }
274 }
275 Ok((folded, unreadable))
276}
277
278#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
280pub struct FetchReport {
281 pub fetched: usize,
283 pub no_origin: usize,
285 pub failed: usize,
287}
288
289pub async fn fetch_origins(
297 roots: &[PathBuf],
298 per_repo: std::time::Duration,
299 stop: impl Fn() -> bool,
300) -> FetchReport {
301 let mut report = FetchReport::default();
302 for repo in crate::repos::scan(roots) {
303 if stop() {
304 break;
305 }
306 match crate::git::git_raw(&repo.path, &["remote", "get-url", "origin"]).await {
307 Ok(out) if out.ok() => {}
308 _ => {
309 report.no_origin += 1;
310 continue;
311 }
312 }
313 match crate::git::fetch_origin(&repo.path, per_repo).await {
314 Ok(()) => report.fetched += 1,
315 Err(e) => {
316 tracing::warn!("fetch origin in {}: {e:#}", repo.name);
317 report.failed += 1;
318 }
319 }
320 }
321 report
322}
323
324pub fn due(now: Timestamp, updated: Timestamp, grace_secs: u64) -> bool {
330 now.duration_since(updated) > SignedDuration::new(grace_secs as i64, 0)
331}
332
333#[derive(Deserialize)]
336struct Meta {
337 status: RunStatus,
338 updated_at: Timestamp,
339}
340
341fn read_meta(runs: &Path, id: &str) -> Result<Meta> {
345 let path = runs.join(id).join("run.json");
346 let body =
347 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
348 let meta: Meta =
349 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
350 Ok(meta)
351}
352
353const MAX_EXTERNAL_MERGE_CHECKS_PER_PASS: usize = 5;
359
360#[derive(Deserialize)]
363struct ExternalMergeMeta {
364 status: RunStatus,
365 updated_at: Timestamp,
366 #[serde(default)]
367 merge: Option<crate::run::MergeOutcome>,
368}
369
370fn read_external_merge_meta(runs: &Path, id: &str) -> Result<ExternalMergeMeta> {
371 let path = runs.join(id).join("run.json");
372 let body =
373 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
374 let meta: ExternalMergeMeta =
375 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
376 Ok(meta)
377}
378
379pub fn eligible_for_external_merge_check(
388 status: RunStatus,
389 merge_recorded: bool,
390 updated_at: Timestamp,
391 now: Timestamp,
392 grace_secs: u64,
393) -> bool {
394 status == RunStatus::Blocked && !merge_recorded && due(now, updated_at, grace_secs)
395}
396
397fn rotate_from(ids: &[String], start: usize) -> Vec<String> {
415 if ids.is_empty() {
416 return Vec::new();
417 }
418 let start = start % ids.len();
419 ids[start..]
420 .iter()
421 .chain(ids[..start].iter())
422 .cloned()
423 .collect()
424}
425
426const EXTERNAL_MERGE_CURSOR_FILE: &str = "external-merge-cursor";
430
431fn read_external_merge_cursor(home: &Path) -> usize {
436 std::fs::read_to_string(home.join(EXTERNAL_MERGE_CURSOR_FILE))
437 .ok()
438 .and_then(|s| s.trim().parse().ok())
439 .unwrap_or(0)
440}
441
442fn write_external_merge_cursor(home: &Path, cursor: usize) {
447 if let Err(e) = std::fs::write(home.join(EXTERNAL_MERGE_CURSOR_FILE), cursor.to_string()) {
448 tracing::warn!("housekeep: persist the external-merge sweep's cursor: {e:#}");
449 }
450}
451
452async fn reconcile_external_merges(runs: &Path, home: &Path, disk: &Disk, now: Timestamp) -> usize {
473 let mut ids: Vec<String> = std::fs::read_dir(runs)
474 .into_iter()
475 .flatten()
476 .flatten()
477 .filter(|e| e.path().join("run.json").is_file())
478 .map(|e| e.file_name().to_string_lossy().into_owned())
479 .collect();
480 ids.sort_unstable();
481
482 let mut eligible: Vec<String> = Vec::new();
488 for id in ids {
489 if crate::daemon::is_working_on(home, &id, now) {
490 continue;
491 }
492 let meta = match read_external_merge_meta(runs, &id) {
493 Ok(meta) => meta,
494 Err(_) => continue,
497 };
498 if eligible_for_external_merge_check(
499 meta.status,
500 meta.merge.is_some(),
501 meta.updated_at,
502 now,
503 disk.fold_grace_secs,
504 ) {
505 eligible.push(id);
506 }
507 }
508
509 if eligible.is_empty() {
510 return 0;
511 }
512 let cursor = read_external_merge_cursor(home);
513 let window: Vec<String> = rotate_from(&eligible, cursor)
514 .into_iter()
515 .take(MAX_EXTERNAL_MERGE_CHECKS_PER_PASS)
516 .collect();
517 write_external_merge_cursor(home, (cursor + window.len()) % eligible.len());
522
523 let mut reconciled = 0usize;
524 for id in window {
525 let mut state = match read_state(runs, &id) {
526 Ok(state) => state,
527 Err(_) => continue,
528 };
529 match crate::land::find_external_merge(&state).await {
530 Ok(Some(found)) => {
531 match crate::land::correct_confirmed_external_merge(&mut state, &found.url).await {
532 Ok(_) => {
533 if let Err(e) = crate::graph::fold_run(&mut state, true, home).await {
534 tracing::warn!(
535 "housekeep: fold {id} after recording its external merge: {e:#}"
536 );
537 }
538 reconciled += 1;
539 }
540 Err(e) => tracing::warn!(
541 "housekeep: record external merge {} for {id}: {e:#}",
542 found.url
543 ),
544 }
545 }
546 Ok(None) => {}
547 Err(e) => {
548 tracing::warn!("housekeep: check external merge for {id}: {e:#}");
549 crate::notices::raise_in(
550 home,
551 crate::notices::Notice::warn(
552 &format!("merged-unrecorded:{id}"),
553 format!(
554 "Run {id} is blocked with no recorded merge, and checking GitHub \
555 for a merge failed; check by hand."
556 ),
557 )
558 .link(crate::notices::Link::Run { id: id.clone() }),
559 );
560 }
561 }
562 }
563 reconciled
564}
565
566fn read_state(runs: &Path, id: &str) -> Result<RunState> {
581 let path = runs.join(id).join("run.json");
582 let body =
583 std::fs::read_to_string(&path).with_context(|| format!("read {}", path.display()))?;
584 let state: RunState =
585 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
586 Ok(state)
587}
588
589pub async fn fold_unreadable(runs: &Path, worktrees_root: &Path, id: &str) -> Result<Vec<String>> {
605 let resolved = resolve_id_path(runs, id)?;
606 let mut removed = Vec::new();
607 let run_dir = runs.join(&resolved);
608 if run_dir.exists() {
609 std::fs::remove_dir_all(&run_dir)
610 .with_context(|| format!("remove {}", run_dir.display()))?;
611 removed.push(format!("runs/{resolved}"));
612 }
613 let wt = worktrees_root.join(short_of(&resolved));
614 if wt.exists() {
615 crate::git::remove_worktree_from_linked(&wt).await;
616 for e in std::fs::read_dir(&wt).into_iter().flatten().flatten() {
617 crate::git::remove_worktree_from_linked(&e.path()).await;
618 }
619 std::fs::remove_dir_all(&wt).with_context(|| format!("remove {}", wt.display()))?;
620 removed.push(wt.to_string_lossy().into_owned());
621 }
622 Ok(removed)
623}
624
625pub async fn fold_orphaned_worktrees(
664 runs: &Path,
665 worktrees_root: &Path,
666 home: &Path,
667 grace_secs: u64,
668 now: Timestamp,
669) -> usize {
670 let known: std::collections::HashSet<String> = std::fs::read_dir(runs)
671 .into_iter()
672 .flatten()
673 .flatten()
674 .map(|e| e.file_name().to_string_lossy().into_owned())
675 .filter(|name| crate::run::is_run_id(name))
676 .map(|id| short_of(&id).to_owned())
677 .collect();
678
679 let mut folded = 0usize;
680 for entry in std::fs::read_dir(worktrees_root)
681 .into_iter()
682 .flatten()
683 .flatten()
684 {
685 if !entry.path().is_dir() {
686 continue;
687 }
688 let short = entry.file_name().to_string_lossy().into_owned();
689 if !looks_like_a_worktree_bay(&short) {
690 continue;
691 }
692 if known.contains(&short) || crate::daemon::is_working_on_short(home, &short, now) {
693 continue;
694 }
695 let wt = entry.path();
696 if !stale_enough(&wt, grace_secs, now) {
697 continue;
698 }
699 crate::git::remove_worktree_from_linked(&wt).await;
700 for e in std::fs::read_dir(&wt).into_iter().flatten().flatten() {
701 crate::git::remove_worktree_from_linked(&e.path()).await;
702 }
703 match std::fs::remove_dir_all(&wt) {
704 Ok(()) => folded += 1,
705 Err(e) => tracing::warn!(
706 "housekeep: remove orphaned worktree {}: {e:#}",
707 wt.display()
708 ),
709 }
710 }
711 folded
712}
713
714fn looks_like_a_worktree_bay(name: &str) -> bool {
723 name.len() == 4 && name.bytes().all(|b| b.is_ascii_alphanumeric())
724}
725
726fn stale_enough(dir: &Path, grace_secs: u64, now: Timestamp) -> bool {
745 let Ok(modified) = std::fs::metadata(dir).and_then(|m| m.modified()) else {
746 return false;
747 };
748 let Ok(ts) = Timestamp::try_from(modified) else {
749 return false;
750 };
751 due(now, ts, grace_secs.max(MIN_ORPHAN_AGE_SECS))
752}
753
754const MIN_ORPHAN_AGE_SECS: u64 = 5 * 60;
765
766fn resolve_id_path(runs: &Path, prefix: &str) -> Result<String> {
769 if runs.join(prefix).is_dir() && crate::run::is_run_id(prefix) {
774 return Ok(prefix.to_owned());
775 }
776 let mut hits: Vec<String> = Vec::new();
777 for e in std::fs::read_dir(runs).into_iter().flatten().flatten() {
778 if !e.path().is_dir() {
779 continue;
780 }
781 let id = e.file_name().to_string_lossy().into_owned();
782 if crate::run::is_run_id(&id) && (id.starts_with(prefix) || id.ends_with(prefix)) {
783 hits.push(id);
784 }
785 }
786 match hits.len() {
787 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
788 0 => bail!("no run matches `{prefix}`"),
789 _ => bail!(
790 "`{prefix}` matches {} runs: {}",
791 hits.len(),
792 hits.join(", ")
793 ),
794 }
795}
796
797pub fn clear_abandoned_active(state: &mut RunState, home: &Path, now: Timestamp) -> Result<bool> {
812 if crate::daemon::is_working_on(home, &state.id, now) || !state.active_all_overrun(now) {
813 return Ok(false);
814 }
815 state.abandon("fold");
816 state.save_under(home)?;
817 if let Err(e) = Questions::at(home.join("questions")).settle_run(&state.id, state.status) {
825 tracing::warn!("abandon questions for {}: {e:#}", state.id);
826 }
827 Ok(true)
828}
829
830pub fn prune_cache(home: &Path, cache: &Path, limit_bytes: u64) -> Result<Option<Prune>> {
839 crate::cache::maintenance_prune(home, cache, limit_bytes)
840}
841
842pub fn prune_cache_if_over_limit(
852 cfg: &crate::config::Config,
853 home: &Path,
854) -> Result<Option<Prune>> {
855 if cfg.disk.cache_limit_bytes == 0 {
856 return Ok(None);
857 }
858 let Some(cache) = cfg.cache_dir() else {
859 return Ok(None);
860 };
861 prune_cache(home, &cache, cfg.disk.cache_limit_bytes)
862}
863
864pub fn cache_report(cfg: &crate::config::Config) -> Option<(PathBuf, u64, u64)> {
870 let cache = cfg.cache_dir()?;
871 Some((
872 cache.clone(),
873 cache_size(&cache),
874 cfg.disk.cache_limit_bytes,
875 ))
876}
877
878pub fn cache_size(cache: &Path) -> u64 {
880 dir_size(cache)
881}
882
883#[cfg(test)]
884mod tests {
885
886 async fn sh(cwd: &Path, args: &[&str]) {
887 let out = crate::git::git_raw(cwd, args).await.expect("spawn git");
888 assert!(out.ok(), "git {args:?}: {}", out.stderr);
889 }
890
891 async fn commit(cwd: &Path, file: &str, body: &str) {
892 std::fs::write(cwd.join(file), body).unwrap();
893 sh(cwd, &["add", file]).await;
894 sh(
895 cwd,
896 &[
897 "-c",
898 "user.name=t",
899 "-c",
900 "user.email=t@example.com",
901 "commit",
902 "-q",
903 "-m",
904 body,
905 ],
906 )
907 .await;
908 }
909
910 const LONG: std::time::Duration = std::time::Duration::from_secs(30);
911
912 #[tokio::test]
913 async fn fetch_advances_origin_main_and_leaves_the_checkout_alone() {
914 let tmp = tempfile::tempdir().unwrap();
915 let t = tmp.path();
916 let bare = t.join("remote.git");
917 sh(
918 t,
919 &["init", "-q", "--bare", "-b", "main", bare.to_str().unwrap()],
920 )
921 .await;
922 let dir = t.join("root/h/o/r");
923 std::fs::create_dir_all(dir.parent().unwrap()).unwrap();
924 sh(
925 t,
926 &["clone", "-q", bare.to_str().unwrap(), dir.to_str().unwrap()],
927 )
928 .await;
929 sh(&dir, &["checkout", "-q", "-b", "main"]).await;
930 commit(&dir, "a.txt", "one").await;
931 sh(&dir, &["push", "-q", "origin", "main"]).await;
932 sh(&dir, &["checkout", "-q", "--detach"]).await;
933 std::fs::write(dir.join("a.txt"), "dirty").unwrap();
934
935 let other = t.join("other");
936 sh(
937 t,
938 &[
939 "clone",
940 "-q",
941 bare.to_str().unwrap(),
942 other.to_str().unwrap(),
943 ],
944 )
945 .await;
946 commit(&other, "b.txt", "two").await;
947 sh(&other, &["push", "-q", "origin", "HEAD:main"]).await;
948
949 let rev = |d: PathBuf, r: &'static str| async move {
950 crate::git::git(&d, &["rev-parse", r]).await.unwrap()
951 };
952 let head = rev(dir.clone(), "HEAD").await;
953 let before = rev(dir.clone(), "origin/main").await;
954 let status = crate::git::git(&dir, &["status", "--porcelain"])
955 .await
956 .unwrap();
957
958 let r = fetch_origins(&[t.join("root")], LONG, || false).await;
959 assert_eq!(
960 r,
961 FetchReport {
962 fetched: 1,
963 no_origin: 0,
964 failed: 0
965 }
966 );
967
968 assert_ne!(rev(dir.clone(), "origin/main").await, before);
969 assert_eq!(rev(dir.clone(), "HEAD").await, head);
970 assert_eq!(
971 crate::git::git(&dir, &["status", "--porcelain"])
972 .await
973 .unwrap(),
974 status
975 );
976 assert_eq!(std::fs::read_to_string(dir.join("a.txt")).unwrap(), "dirty");
977 }
978
979 #[tokio::test]
980 async fn fetch_never_writes_a_local_branch_whatever_the_configured_refspec() {
981 let tmp = tempfile::tempdir().unwrap();
982 let t = tmp.path();
983 let bare = t.join("remote.git");
984 sh(
985 t,
986 &["init", "-q", "--bare", "-b", "main", bare.to_str().unwrap()],
987 )
988 .await;
989 let dir = t.join("root/h/o/r");
990 std::fs::create_dir_all(dir.parent().unwrap()).unwrap();
991 sh(
992 t,
993 &["clone", "-q", bare.to_str().unwrap(), dir.to_str().unwrap()],
994 )
995 .await;
996 sh(&dir, &["checkout", "-q", "-b", "main"]).await;
997 commit(&dir, "a.txt", "one").await;
998 sh(&dir, &["push", "-q", "origin", "main"]).await;
999 sh(&dir, &["checkout", "-q", "--detach"]).await;
1000 sh(
1002 &dir,
1003 &[
1004 "config",
1005 "remote.origin.fetch",
1006 "+refs/heads/main:refs/heads/main",
1007 ],
1008 )
1009 .await;
1010 let _ = crate::git::git_raw(&dir, &["update-ref", "-d", "refs/remotes/origin/main"]).await;
1011
1012 let other = t.join("other");
1013 sh(
1014 t,
1015 &[
1016 "clone",
1017 "-q",
1018 bare.to_str().unwrap(),
1019 other.to_str().unwrap(),
1020 ],
1021 )
1022 .await;
1023 commit(&other, "b.txt", "two").await;
1024 sh(&other, &["push", "-q", "origin", "HEAD:main"]).await;
1025 let remote_tip = crate::git::git(&other, &["rev-parse", "HEAD"])
1026 .await
1027 .unwrap();
1028 let local_main = crate::git::git(&dir, &["rev-parse", "refs/heads/main"])
1029 .await
1030 .unwrap();
1031
1032 let r = fetch_origins(&[t.join("root")], LONG, || false).await;
1033 assert_eq!(r.fetched, 1);
1034 assert_eq!(
1035 crate::git::git(&dir, &["rev-parse", "refs/remotes/origin/main"])
1036 .await
1037 .unwrap(),
1038 remote_tip
1039 );
1040 assert_eq!(
1041 crate::git::git(&dir, &["rev-parse", "refs/heads/main"])
1042 .await
1043 .unwrap(),
1044 local_main
1045 );
1046 }
1047
1048 #[tokio::test]
1049 async fn fetch_skips_no_origin_and_survives_an_unreachable_origin() {
1050 let tmp = tempfile::tempdir().unwrap();
1051 let t = tmp.path();
1052 let root = t.join("root");
1053 let bare = t.join("remote.git");
1054 sh(t, &["init", "-q", "--bare", bare.to_str().unwrap()]).await;
1055
1056 let none = root.join("h/o/none");
1057 let dead = root.join("h/o/dead");
1058 let good = root.join("h/o/good");
1059 for d in [&none, &dead, &good] {
1060 std::fs::create_dir_all(d).unwrap();
1061 sh(d, &["init", "-q"]).await;
1062 }
1063 sh(
1064 &dead,
1065 &[
1066 "remote",
1067 "add",
1068 "origin",
1069 t.join("missing").to_str().unwrap(),
1070 ],
1071 )
1072 .await;
1073 sh(&good, &["remote", "add", "origin", bare.to_str().unwrap()]).await;
1074
1075 let r = fetch_origins(&[root], LONG, || false).await;
1076 assert_eq!(
1077 r,
1078 FetchReport {
1079 fetched: 1,
1080 no_origin: 1,
1081 failed: 1
1082 }
1083 );
1084 }
1085 use super::*;
1086 use crate::config::Disk;
1087 use std::fs;
1088
1089 fn ts(s: &str) -> Timestamp {
1090 s.parse().expect("rfc3339")
1091 }
1092
1093 fn block_on<F: std::future::Future>(f: F) -> F::Output {
1094 tokio::runtime::Runtime::new().expect("runtime").block_on(f)
1095 }
1096
1097 #[test]
1098 fn a_run_is_due_after_its_grace_and_not_before() {
1099 let now = ts("2026-09-05T00:00:00Z");
1100 let grace = 600;
1101 let old = now - SignedDuration::new(601, 0);
1102 let fresh = now - SignedDuration::new(599, 0);
1103 assert!(due(now, old, grace));
1104 assert!(!due(now, fresh, grace));
1105 let edge = now - SignedDuration::new(600, 0);
1107 assert!(!due(now, edge, grace));
1108 assert!(due(now, old, 0));
1110 }
1111
1112 #[test]
1113 fn rotate_from_moves_the_starting_point_as_the_cursor_advances() {
1114 let ids: Vec<String> = ["a", "b", "c", "d", "e"]
1115 .iter()
1116 .map(|s| s.to_string())
1117 .collect();
1118
1119 assert_eq!(rotate_from(&ids, 0), ids);
1121
1122 let at_2 = rotate_from(&ids, 2);
1125 assert_eq!(at_2, vec!["c", "d", "e", "a", "b"]);
1126 for id in &ids {
1127 assert!(at_2.contains(id));
1128 }
1129
1130 assert_eq!(rotate_from(&ids, 7), rotate_from(&ids, 2));
1133
1134 assert_eq!(rotate_from(&[], 0), Vec::<String>::new());
1136 }
1137
1138 #[test]
1145 fn advancing_the_cursor_by_the_window_size_gives_every_id_a_turn() {
1146 let ids: Vec<String> = (0..12).map(|n| format!("run-{n}")).collect();
1147 let cap = 5usize;
1148 let mut cursor = 0usize;
1149 let mut ever_seen: std::collections::HashSet<String> = std::collections::HashSet::new();
1150 for _ in 0..ids.len() {
1151 let window: Vec<String> = rotate_from(&ids, cursor).into_iter().take(cap).collect();
1152 ever_seen.extend(window.iter().cloned());
1153 cursor = (cursor + window.len()) % ids.len();
1154 }
1155 assert_eq!(
1156 ever_seen.len(),
1157 ids.len(),
1158 "every id must be checked at least once across a full rotation, whatever \
1159 the cadence between passes"
1160 );
1161 }
1162
1163 #[test]
1164 fn the_external_merge_cursor_round_trips_through_a_files_absence() {
1165 let dir = tempfile::tempdir().unwrap();
1166 let home = dir.path();
1167
1168 assert_eq!(read_external_merge_cursor(home), 0);
1170
1171 write_external_merge_cursor(home, 7);
1172 assert_eq!(read_external_merge_cursor(home), 7);
1173
1174 std::fs::write(home.join(EXTERNAL_MERGE_CURSOR_FILE), "not a number").unwrap();
1177 assert_eq!(read_external_merge_cursor(home), 0);
1178 }
1179
1180 #[test]
1181 fn a_run_qualifies_for_an_external_merge_check_only_when_blocked_unmerged_and_due() {
1182 let now = ts("2026-09-05T00:00:00Z");
1183 let grace = 600;
1184 let old = now - SignedDuration::new(601, 0);
1185 let fresh = now - SignedDuration::new(599, 0);
1186
1187 assert!(
1188 eligible_for_external_merge_check(RunStatus::Blocked, false, old, now, grace),
1189 "blocked, unmerged, and past its grace period is exactly the run this exists for"
1190 );
1191 assert!(
1192 !eligible_for_external_merge_check(RunStatus::Blocked, false, fresh, now, grace),
1193 "an operator mid-fix deserves the same grace window `fold_due` gives before \
1194 the janitor starts asking GitHub about it"
1195 );
1196 assert!(
1197 !eligible_for_external_merge_check(RunStatus::Blocked, true, old, now, grace),
1198 "a run `land::land` already recorded a merge outcome for has its own answer \
1199 already; this check is only for a run with nothing recorded at all"
1200 );
1201 assert!(
1202 !eligible_for_external_merge_check(RunStatus::Stalled, false, old, now, grace),
1203 "stalled is not blocked - it means the tally never reached quorum, which a \
1204 pull request cannot fix"
1205 );
1206 assert!(
1207 !eligible_for_external_merge_check(RunStatus::Ready, false, old, now, grace),
1208 "ready has nothing to correct - it was never landed by design"
1209 );
1210 assert!(
1211 !eligible_for_external_merge_check(RunStatus::Superseded, false, old, now, grace),
1212 "a superseded run's task already has its answer from a later attempt; \
1213 there is nothing left for GitHub to confirm here"
1214 );
1215 }
1216
1217 fn write_blocked_run(runs: &Path, id: &str, updated_at: Timestamp) {
1223 use crate::run::{Candidate, Tally};
1224 use std::collections::BTreeMap;
1225
1226 let mut state = RunState::new(
1227 PathBuf::from("/nonexistent/repo"),
1228 "main".to_owned(),
1229 "0000000000000000000000000000000000000000".to_owned(),
1230 String::new(),
1231 crate::config::Config::default(),
1232 );
1233 state.id = id.to_owned();
1234 state.status = RunStatus::Blocked;
1235 state.updated_at = updated_at;
1236 state.candidates.push(Candidate {
1237 index: 0,
1238 label: 'A',
1239 agent: "agent".to_owned(),
1240 branch: format!("magi/{id}/A"),
1241 worktree: PathBuf::from("/nonexistent/repo"),
1242 summary: String::new(),
1243 stat: String::new(),
1244 files: 0,
1245 commits: 0,
1246 empty: false,
1247 failed: None,
1248 verified_noop: None,
1249 duration_ms: 0,
1250 folded: false,
1251 });
1252 state.tally = Some(Tally {
1253 first_choice: BTreeMap::from([('A', 1)]),
1254 borda: BTreeMap::new(),
1255 winner: 'A',
1256 rankings: 1,
1257 unanimous_initial: true,
1258 deliberated: false,
1259 changed_votes: 0,
1260 unanimous_final: true,
1261 tie_break: None,
1262 judges: 1,
1263 present: 1,
1264 quorum: 1,
1265 met_quorum: true,
1266 uncontested: None,
1267 });
1268 std::fs::create_dir_all(runs.join(id)).unwrap();
1269 std::fs::write(
1270 runs.join(id).join("run.json"),
1271 serde_json::to_string_pretty(&state).unwrap(),
1272 )
1273 .unwrap();
1274 }
1275
1276 #[tokio::test]
1281 async fn reconcile_external_merges_leaves_ineligible_and_unconfirmable_runs_alone() {
1282 let dir = tempfile::tempdir().unwrap();
1283 let runs = dir.path().join("runs");
1284 let home = dir.path().to_path_buf();
1285 std::fs::create_dir_all(&runs).unwrap();
1286
1287 let now = ts("2026-09-05T00:00:00Z");
1288 let disk = Disk {
1289 fold_grace_secs: 600,
1290 ..Disk::default()
1291 };
1292
1293 let due_id = "20260905-000000-blkd";
1294 write_blocked_run(&runs, due_id, ts("2026-08-01T00:00:00Z"));
1295
1296 let fresh_id = "20260905-000000-fres";
1297 write_blocked_run(&runs, fresh_id, now);
1298
1299 let live_id = "20260905-000000-live";
1300 write_blocked_run(&runs, live_id, ts("2026-08-01T00:00:00Z"));
1301 let mut status = crate::daemon::Status::new();
1302 status.current = vec![crate::daemon::Current {
1303 task: "20260905-000000-task".to_owned(),
1304 run: live_id.to_owned(),
1305 }];
1306 status.updated_at = now;
1307 crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
1308
1309 let reconciled = reconcile_external_merges(&runs, &home, &disk, now).await;
1310 assert_eq!(
1311 reconciled, 0,
1312 "an unreachable repository can never be confirmed merged"
1313 );
1314 for id in [due_id, fresh_id, live_id] {
1315 assert_eq!(
1316 read_meta(&runs, id).unwrap().status,
1317 RunStatus::Blocked,
1318 "{id} must be left exactly as it was found"
1319 );
1320 }
1321 let notices = crate::notices::Notices::at(home.join("notifications")).list();
1322 let due_notice = notices
1323 .iter()
1324 .find(|n| n.key == format!("merged-unrecorded:{due_id}"))
1325 .expect("notice for due_id must be raised");
1326 assert_eq!(
1327 due_notice.link,
1328 Some(crate::notices::Link::Run {
1329 id: due_id.to_owned(),
1330 })
1331 );
1332 }
1333
1334 #[tokio::test]
1344 async fn reconcile_external_merges_covers_a_larger_fleet_across_repeated_passes_at_one_instant()
1345 {
1346 let dir = tempfile::tempdir().unwrap();
1347 let runs = dir.path().join("runs");
1348 let home = dir.path().to_path_buf();
1349 std::fs::create_dir_all(&runs).unwrap();
1350
1351 let now = ts("2026-09-05T00:00:00Z");
1352 let disk = Disk {
1353 fold_grace_secs: 600,
1354 ..Disk::default()
1355 };
1356 let old = ts("2026-08-01T00:00:00Z");
1357
1358 let ids: Vec<String> = (0..8).map(|n| format!("20260905-000000-r{n:03}")).collect();
1359 for id in &ids {
1360 write_blocked_run(&runs, id, old);
1361 }
1362
1363 let notices = crate::notices::Notices::at(home.join("notifications"));
1364
1365 reconcile_external_merges(&runs, &home, &disk, now).await;
1371 let after_first = notices.list().len();
1372 assert_eq!(
1373 after_first, MAX_EXTERNAL_MERGE_CHECKS_PER_PASS,
1374 "the first pass checks exactly one cap's worth"
1375 );
1376
1377 reconcile_external_merges(&runs, &home, &disk, now).await;
1378 let after_second = notices.list().len();
1379 assert_eq!(
1380 after_second,
1381 ids.len(),
1382 "a second pass at the same instant must still reach every id the \
1383 first pass had no room for, not repeat the same cap's worth"
1384 );
1385 }
1386
1387 #[test]
1388 fn the_meta_reader_is_tolerant_of_everything_except_the_deciders() {
1389 let dir = tempfile::tempdir().unwrap();
1390 let runs = dir.path().join("runs");
1391 let id = "20260905-000000-abcd";
1392 std::fs::create_dir_all(runs.join(id)).unwrap();
1393 std::fs::write(
1394 runs.join(id).join("run.json"),
1395 r#"{"schema": 99, "id": "20260905-000000-abcd", "updated_at": "2026-09-05T00:00:00Z", "status": "ready", "junk_from_another_build": [1, 2, 3]}"#,
1396 )
1397 .unwrap();
1398 let meta = read_meta(&runs, id).expect("readable");
1399 assert_eq!(meta.status, RunStatus::Ready);
1400 assert_eq!(meta.updated_at, ts("2026-09-05T00:00:00Z"));
1401 assert!(read_meta(&runs, "nope").is_err(), "missing file unreadable");
1402 std::fs::write(runs.join(id).join("run.json"), "not json at all").unwrap();
1403 assert!(read_meta(&runs, id).is_err(), "garbage unreadable");
1404 }
1405
1406 #[test]
1407 fn fold_unreadable_releases_run_dir_and_worktrees() {
1408 let dir = tempfile::tempdir().unwrap();
1409 let runs = dir.path().join("runs");
1410 let wt = dir.path().join("wt");
1411 let id = "20260905-000000-abcd";
1412 std::fs::create_dir_all(runs.join(id)).unwrap();
1413 std::fs::write(runs.join(id).join("run.json"), "garbage").unwrap();
1414 std::fs::create_dir_all(wt.join("abcd")).unwrap();
1415 std::fs::write(wt.join("abcd").join("leftover"), b"x").unwrap();
1416
1417 let removed = block_on(fold_unreadable(&runs, &wt, id)).expect("fold");
1418 assert_eq!(removed.len(), 2);
1419 assert!(!runs.join(id).exists(), "run dir gone");
1420 assert!(!wt.join("abcd").exists(), "worktrees gone");
1421
1422 std::fs::create_dir_all(runs.join(id)).unwrap();
1424 std::fs::write(runs.join(id).join("run.json"), "garbage").unwrap();
1425 std::fs::create_dir_all(wt.join("abcd")).unwrap();
1426 std::fs::write(wt.join("abcd").join("leftover"), b"x").unwrap();
1427 let removed = block_on(fold_unreadable(&runs, &wt, "20260905")).expect("by prefix");
1428 assert_eq!(removed.len(), 2);
1429 assert!(
1433 block_on(fold_unreadable(&runs, &wt, id)).is_err(),
1434 "a run already gone cannot be resolved again"
1435 );
1436 }
1437
1438 #[test]
1439 fn prune_cache_sheds_the_oldest_generation_until_it_fits() {
1440 let home = tempfile::tempdir().unwrap();
1441 let dir = tempfile::tempdir().unwrap();
1442 fs::write(dir.path().join("old"), b"xx").unwrap();
1445 fs::write(dir.path().join("new"), b"yy").unwrap();
1446 touch(&dir.path().join("old"), 1_000_000);
1447 touch(&dir.path().join("new"), 2_000_000);
1448
1449 let out = prune_cache(home.path(), dir.path(), 2)
1450 .expect("prune")
1451 .expect("the cache is free");
1452 assert_eq!(out.files, 1, "one deletion is enough to reach the cap");
1453 assert_eq!(out.remaining, 2);
1454 assert!(!dir.path().join("old").exists(), "the older file went");
1455 assert!(dir.path().join("new").exists(), "the newer one stayed");
1456
1457 let tied = tempfile::tempdir().unwrap();
1462 fs::write(tied.path().join("big"), b"xxxx").unwrap();
1463 fs::write(tied.path().join("small"), b"yy").unwrap();
1464 touch(&tied.path().join("big"), 1_000_000);
1465 touch(&tied.path().join("small"), 1_000_000);
1466 let out = prune_cache(home.path(), tied.path(), 2)
1467 .expect("prune")
1468 .expect("the cache is free");
1469 assert_eq!(out.files, 1, "the big one alone gets under the cap");
1470 assert_eq!(out.remaining, 2);
1471 assert!(tied.path().join("small").exists());
1472 }
1473
1474 #[test]
1475 fn prune_cache_if_over_limit_resolves_the_opt_outs_before_ever_measuring() {
1476 let home = tempfile::tempdir().unwrap();
1477 let dir = tempfile::tempdir().unwrap();
1478 fs::write(dir.path().join("big"), vec![0u8; 10]).unwrap();
1479
1480 let mut cfg = crate::config::Config::default();
1481 cfg.verify.gate = vec![format!(
1482 "CARGO_TARGET_DIR={} cargo make check",
1483 dir.path().display()
1484 )];
1485
1486 cfg.disk.cache_limit_bytes = 0;
1489 assert_eq!(
1490 prune_cache_if_over_limit(&cfg, home.path()).unwrap(),
1491 None,
1492 "a zero cap must not even look at the directory"
1493 );
1494 assert!(dir.path().join("big").exists());
1495
1496 let mut no_cache = crate::config::Config::default();
1499 no_cache.disk.cache_limit_bytes = 1;
1500 assert_eq!(
1501 prune_cache_if_over_limit(&no_cache, home.path()).unwrap(),
1502 None
1503 );
1504
1505 cfg.disk.cache_limit_bytes = 1;
1508 let pruned = prune_cache_if_over_limit(&cfg, home.path())
1509 .unwrap()
1510 .expect("a real cache dir over its cap prunes");
1511 assert_eq!(pruned.files, 1);
1512 assert!(!dir.path().join("big").exists());
1513 }
1514
1515 fn touch(path: &Path, secs: u64) {
1518 let f = fs::File::options().write(true).open(path).unwrap();
1519 f.set_times(fs::FileTimes::new().set_modified(
1520 std::time::SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(secs),
1521 ))
1522 .unwrap();
1523 }
1524
1525 #[test]
1528 fn fold_unreadable_clears_a_run_whose_state_never_landed() {
1529 let dir = tempfile::tempdir().unwrap();
1530 let runs = dir.path().join("runs");
1531 let wt = dir.path().join("wt");
1532 let id = "20260904-014540-88c0";
1533 std::fs::create_dir_all(runs.join(id)).unwrap();
1534 std::fs::write(runs.join(id).join("run.json.tmp"), b"").unwrap();
1535
1536 let removed = block_on(fold_unreadable(&runs, &wt, id)).expect("fold by id");
1537 assert_eq!(removed, vec![format!("runs/{id}")]);
1538 assert!(!runs.join(id).exists(), "record gone");
1539
1540 std::fs::create_dir_all(runs.join(id)).unwrap();
1542 std::fs::write(runs.join(id).join("run.json.tmp"), b"").unwrap();
1543 assert!(
1544 block_on(fold_unreadable(&runs, &wt, "88c0")).is_ok(),
1545 "by prefix"
1546 );
1547
1548 std::fs::create_dir_all(runs.join("scratch")).unwrap();
1550 assert!(
1551 block_on(fold_unreadable(&runs, &wt, "scratch")).is_err(),
1552 "a stray directory is not a run"
1553 );
1554 }
1555
1556 #[test]
1557 fn fold_due_folds_terminal_runs_of_any_schema_but_leaves_genuinely_unreadable_ones() {
1558 let dir = tempfile::tempdir().unwrap();
1559 let runs = dir.path().join("runs");
1560 let wt = dir.path().join("wt");
1561 let home = dir.path().to_path_buf();
1562 let disk = Disk::default();
1563 let now = ts("2026-09-05T00:00:00Z");
1564
1565 let judging = "20260801-000000-0001";
1567 write_meta(&runs, judging, "judging", "2026-08-01T00:00:00Z");
1568
1569 let ready_fresh = "20260904-220000-0002";
1576 write_meta(&runs, ready_fresh, "ready", "2026-09-04T22:00:00Z");
1577
1578 let garbage = "20260901-000000-0004";
1584 std::fs::create_dir_all(runs.join(garbage)).unwrap();
1585 std::fs::write(runs.join(garbage).join("run.json"), "not json").unwrap();
1586 std::fs::create_dir_all(wt.join("0004")).unwrap();
1587
1588 let due_ready = due_run(&runs, &wt, "20260801-000000-ffff", SCHEMA);
1591
1592 let due_old_schema = due_run(&runs, &wt, "20260801-000000-eeee", SCHEMA - 1);
1597
1598 let (folded, unreadable) =
1599 block_on(fold_due(&runs, &home, &wt, &disk, now)).expect("fold_due");
1600 assert_eq!(
1601 folded, 2,
1602 "both due, parseable runs fold regardless of their schema number"
1603 );
1604 assert_eq!(
1605 unreadable, 1,
1606 "only the run with broken JSON counts as unreadable"
1607 );
1608 assert!(runs.join(judging).exists(), "runnable never folded");
1609 assert!(runs.join(ready_fresh).exists(), "fresh never folded");
1610 assert!(runs.join(garbage).exists(), "unreadable record kept");
1611 assert!(wt.join("0004").exists(), "unreadable worktree kept");
1612 assert!(
1613 runs.join(&due_ready).exists(),
1614 "folding drops worktrees, not the record"
1615 );
1616 assert!(
1617 runs.join(&due_old_schema).exists(),
1618 "an old-schema record survives its fold exactly like a current one"
1619 );
1620 for id in [&due_ready, &due_old_schema] {
1626 let saved = read_meta(&runs, id).expect("folded run still parses");
1627 assert_ne!(
1628 saved.updated_at,
1629 ts("2026-08-01T00:00:00Z"),
1630 "fold_run must have saved the updated state back through the \
1631 `runs` directory this test passed to fold_due"
1632 );
1633 }
1634 }
1635
1636 #[tokio::test]
1642 async fn housekeep_leaves_everything_alone_when_auto_fold_is_disabled() {
1643 let dir = tempfile::tempdir().unwrap();
1644 let runs = dir.path().join("runs");
1645 let wt = dir.path().join("wt");
1646 let home = dir.path().to_path_buf();
1647 crate::run::set_home(dir.path().to_path_buf());
1648
1649 let due_id = due_run(&runs, &wt, "20260801-000000-abcd", SCHEMA);
1650 std::fs::create_dir_all(wt.join("orphan").join("cand-A")).unwrap();
1651
1652 let mut cfg = crate::config::Config::default();
1653 cfg.disk.auto_fold = false;
1654 cfg.disk.cache_limit_bytes = 0;
1655
1656 let out = housekeep(&cfg, &home, &wt, &dir.path().join("repo"), Timestamp::now()).await;
1657
1658 assert_eq!(out.folded, 0);
1659 assert_eq!(out.unreadable, 0);
1660 assert_eq!(out.orphaned_worktrees, 0);
1661 assert!(
1662 runs.join(&due_id).exists(),
1663 "a due run's record survives untouched"
1664 );
1665 assert!(
1666 wt.join("orphan").exists(),
1667 "an orphaned worktree survives untouched: the reclaim pass never ran"
1668 );
1669 }
1670
1671 #[test]
1672 fn fold_orphaned_worktrees_removes_only_worktrees_no_run_claims_and_none_in_flight() {
1673 let dir = tempfile::tempdir().unwrap();
1674 let runs = dir.path().join("runs");
1675 let wt = dir.path().join("wt");
1676 let home = dir.path().to_path_buf();
1677
1678 write_meta(
1681 &runs,
1682 "20260801-000000-aaaa",
1683 "ready",
1684 "2026-08-01T00:00:00Z",
1685 );
1686 std::fs::create_dir_all(wt.join("aaaa").join("cand-A")).unwrap();
1687
1688 std::fs::create_dir_all(wt.join("bbbb").join("cand-A")).unwrap();
1692
1693 std::fs::create_dir_all(wt.join("cccc")).unwrap();
1698
1699 std::fs::create_dir_all(wt.join("scratch")).unwrap();
1703
1704 let now = Timestamp::now() + SignedDuration::new((MIN_ORPHAN_AGE_SECS + 1) as i64, 0);
1710 let mut status = crate::daemon::Status::new();
1711 status.current = vec![crate::daemon::Current {
1712 task: "20260905-000000-t111".to_owned(),
1713 run: "20260905-000000-cccc".to_owned(),
1714 }];
1715 status.updated_at = now;
1716 crate::daemon::write_status_to(&home.join("daemon.json"), &status).unwrap();
1717
1718 let folded = block_on(fold_orphaned_worktrees(&runs, &wt, &home, 0, now));
1719 assert_eq!(
1720 folded, 1,
1721 "only the truly orphaned, idle, bay-shaped worktree is removed"
1722 );
1723 assert!(wt.join("aaaa").exists(), "claimed by a run record");
1724 assert!(!wt.join("bbbb").exists(), "orphaned and idle: reclaimed");
1725 assert!(wt.join("cccc").exists(), "a run in flight is never touched");
1726 assert!(
1727 wt.join("scratch").exists(),
1728 "not shaped like a worktree bay, so never a reclaim target"
1729 );
1730 }
1731
1732 #[test]
1740 fn fold_orphaned_worktrees_leaves_a_freshly_created_bay_alone() {
1741 let dir = tempfile::tempdir().unwrap();
1742 let runs = dir.path().join("runs");
1743 let wt = dir.path().join("wt");
1744 let home = dir.path().to_path_buf();
1745
1746 std::fs::create_dir_all(wt.join("dddd").join("under-review")).unwrap();
1747
1748 let now = Timestamp::now();
1749 let folded = block_on(fold_orphaned_worktrees(&runs, &wt, &home, 6 * 60 * 60, now));
1750 assert_eq!(
1751 folded, 0,
1752 "too fresh to tell apart from a run still being set up"
1753 );
1754 assert!(wt.join("dddd").exists());
1755 }
1756
1757 fn open_question(store: &Questions, run: &str) -> crate::ask::Question {
1759 let mut q = crate::ask::Question::new(
1760 run.to_owned(),
1761 "implement".to_owned(),
1762 "impl-A".to_owned(),
1763 "Which storage backend should the cache use?".to_owned(),
1764 String::new(),
1765 vec!["SQLite".to_owned(), "Redis".to_owned()],
1766 );
1767 store.put(&mut q).unwrap();
1768 q
1769 }
1770
1771 #[test]
1777 fn a_finished_runs_open_question_is_swept_up() {
1778 let dir = tempfile::tempdir().unwrap();
1779 let runs = dir.path().join("runs");
1780 let store = Questions::at(dir.path().join("questions"));
1781
1782 let failed = "20260908-205802-c9eb";
1783 write_meta(&runs, failed, "failed", "2026-09-08T20:58:02Z");
1784 let failed_q = open_question(&store, failed);
1785
1786 let merged = "20260908-205501-ca67";
1787 write_meta(&runs, merged, "merged", "2026-09-08T20:55:01Z");
1788 let merged_q = open_question(&store, merged);
1789
1790 let n = abandon_settled_questions(&store, &runs);
1791 assert_eq!(n, 2, "both dead runs' questions are swept in one pass");
1792
1793 for (id, run) in [(&failed_q.id, failed), (&merged_q.id, merged)] {
1794 let back = store.get(id).unwrap();
1795 assert!(!back.status.open(), "{run} is done; nobody reads an answer");
1796 assert!(back.detail.contains(run), "{}", back.detail);
1797 }
1798 }
1799
1800 #[test]
1801 fn a_still_alive_runs_open_question_survives_the_sweep() {
1802 let dir = tempfile::tempdir().unwrap();
1803 let runs = dir.path().join("runs");
1804 let store = Questions::at(dir.path().join("questions"));
1805
1806 for (id, status) in [
1810 ("20260908-000000-b10c", "blocked"),
1811 ("20260908-000000-5ta1", "stalled"),
1812 ("20260908-000000-jud6", "judging"),
1813 ] {
1814 write_meta(&runs, id, status, "2026-09-08T00:00:00Z");
1815 let q = open_question(&store, id);
1816
1817 let n = abandon_settled_questions(&store, &runs);
1818 assert_eq!(n, 0, "{status} run is not done; nothing to sweep");
1819 assert!(
1820 store.get(&q.id).unwrap().status.open(),
1821 "{status} run's question must still be waiting"
1822 );
1823 }
1824 }
1825
1826 #[test]
1827 fn the_sweep_leaves_an_answered_question_and_an_unreadable_run_alone() {
1828 let dir = tempfile::tempdir().unwrap();
1829 let runs = dir.path().join("runs");
1830 let store = Questions::at(dir.path().join("questions"));
1831
1832 let done = "20260908-000000-answ";
1835 write_meta(&runs, done, "failed", "2026-09-08T00:00:00Z");
1836 let mut answered = open_question(&store, done);
1837 answered
1838 .answer(crate::ask::Answer::Choice("SQLite".to_owned()))
1839 .unwrap();
1840 store.put(&mut answered).unwrap();
1841
1842 let gone = "20260908-000000-gone";
1844 let orphan = open_question(&store, gone);
1845
1846 assert_eq!(abandon_settled_questions(&store, &runs), 0);
1847 assert_eq!(
1848 store.get(&answered.id).unwrap().status,
1849 crate::ask::QuestionStatus::Answered,
1850 "a real answer is never overwritten by a sweep"
1851 );
1852 assert!(
1853 store.get(&orphan.id).unwrap().status.open(),
1854 "a run this sweep cannot read is left exactly as it was, not guessed at"
1855 );
1856 }
1857
1858 #[test]
1867 fn fold_orphaned_worktrees_floors_a_zero_grace_at_the_race_safe_minimum() {
1868 let dir = tempfile::tempdir().unwrap();
1869 let runs = dir.path().join("runs");
1870 let wt = dir.path().join("wt");
1871 let home = dir.path().to_path_buf();
1872
1873 std::fs::create_dir_all(wt.join("eeee").join("under-review")).unwrap();
1874
1875 let now = Timestamp::now();
1877 let folded = block_on(fold_orphaned_worktrees(&runs, &wt, &home, 0, now));
1878 assert_eq!(
1879 folded, 0,
1880 "a zero grace must not defeat the race-safety floor"
1881 );
1882 assert!(wt.join("eeee").exists());
1883
1884 let later = now + SignedDuration::new((MIN_ORPHAN_AGE_SECS + 1) as i64, 0);
1887 let folded = block_on(fold_orphaned_worktrees(&runs, &wt, &home, 0, later));
1888 assert_eq!(folded, 1, "old enough now, regardless of the zero grace");
1889 assert!(!wt.join("eeee").exists());
1890 }
1891
1892 #[test]
1893 fn clear_abandoned_active_only_acts_once_dead_and_overrun() {
1894 let dir = tempfile::tempdir().unwrap();
1895 crate::run::set_home(dir.path().to_path_buf());
1900 let home = dir.path().to_path_buf();
1901 let now = ts("2026-09-14T12:00:00Z");
1902 let overrun_seat = || crate::run::ActiveSeat {
1903 node: "implement".to_owned(),
1904 started_at: now - SignedDuration::new(21_000, 0),
1905 timeout_secs: 3_600,
1906 attempt: 0,
1907 task: None,
1908 command: None,
1909 index: None,
1910 total: None,
1911 };
1912
1913 let mut state = RunState::new(
1914 PathBuf::from("/repo"),
1915 "main".to_owned(),
1916 "abc1234".to_owned(),
1917 "fixture".to_owned(),
1918 crate::config::Config::default(),
1919 );
1920 state.status = RunStatus::Implementing;
1921 state.active.insert("impl-A".to_owned(), overrun_seat());
1922
1923 let mut fresh = state.clone();
1926 fresh.active.insert(
1927 "impl-B".to_owned(),
1928 crate::run::ActiveSeat {
1929 node: "implement".to_owned(),
1930 started_at: now,
1931 timeout_secs: 3_600,
1932 attempt: 0,
1933 task: None,
1934 command: None,
1935 index: None,
1936 total: None,
1937 },
1938 );
1939 assert!(!clear_abandoned_active(&mut fresh, &home, now).unwrap());
1940 assert!(!fresh.active.is_empty());
1941 assert_eq!(fresh.status, RunStatus::Implementing);
1942
1943 let store = Questions::at(home.join("questions"));
1944 let q = open_question(&store, &state.id);
1945
1946 assert!(clear_abandoned_active(&mut state, &home, now).unwrap());
1947 assert!(state.active.is_empty());
1948 assert_eq!(state.status, RunStatus::Failed);
1949 assert!(
1950 !store.get(&q.id).unwrap().status.open(),
1951 "the abandoned seat's own open question must not keep badging the \
1952 operator until some later daemon startup notices it"
1953 );
1954 }
1955
1956 fn write_meta(runs: &Path, id: &str, status: &str, updated_at: &str) {
1958 let day = &updated_at[..10];
1959 std::fs::create_dir_all(runs.join(id)).unwrap();
1960 let body = format!(
1961 r#"{{"schema": {SCHEMA}, "id": "{id}", "repo": "/nonexistent/repo", "base_branch": "main", "base_commit": "0000000000000000000000000000000000000000", "instruction": "", "created_at": "{day}T00:00:00Z", "updated_at": "{updated_at}", "status": "{status}", "seed": 1}}"#
1962 );
1963 std::fs::write(runs.join(id).join("run.json"), body).unwrap();
1964 }
1965
1966 fn due_run(runs: &Path, wt: &Path, id: &str, schema: u32) -> String {
1976 due_run_with_status(runs, wt, id, schema, RunStatus::Ready)
1977 }
1978
1979 fn due_run_with_status(
1982 runs: &Path,
1983 wt: &Path,
1984 id: &str,
1985 schema: u32,
1986 status: RunStatus,
1987 ) -> String {
1988 let mut config = crate::config::Config::default();
1989 config.graph.worktree_root = Some(wt.to_path_buf());
1990 let mut state = RunState::new(
1991 PathBuf::from("/nonexistent/repo"),
1992 "main".to_owned(),
1993 "0000000000000000000000000000000000000000".to_owned(),
1994 String::new(),
1995 config,
1996 );
1997 state.id = id.to_owned();
1998 state.status = status;
1999 state.updated_at = ts("2026-08-01T00:00:00Z");
2000 let mut value = serde_json::to_value(&state).unwrap();
2001 value["schema"] = serde_json::json!(schema);
2002 std::fs::create_dir_all(runs.join(id)).unwrap();
2003 std::fs::write(
2004 runs.join(id).join("run.json"),
2005 serde_json::to_string_pretty(&value).unwrap(),
2006 )
2007 .unwrap();
2008 id.to_owned()
2009 }
2010
2011 #[test]
2012 fn fold_due_folds_a_superseded_run_same_as_any_other_terminal_one() {
2013 let dir = tempfile::tempdir().unwrap();
2014 let runs = dir.path().join("runs");
2015 let wt = dir.path().join("wt");
2016 let home = dir.path().to_path_buf();
2017 let disk = Disk::default();
2018 let now = ts("2026-09-05T00:00:00Z");
2019
2020 let superseded = due_run_with_status(
2021 &runs,
2022 &wt,
2023 "20260801-000000-cccc",
2024 SCHEMA,
2025 RunStatus::Superseded,
2026 );
2027
2028 let (folded, unreadable) =
2029 block_on(fold_due(&runs, &home, &wt, &disk, now)).expect("fold_due");
2030 assert_eq!(
2031 folded, 1,
2032 "a superseded run has nothing left for a human to check, so it folds \
2033 exactly like a merged one"
2034 );
2035 assert_eq!(unreadable, 0);
2036 assert!(
2037 read_meta(&runs, &superseded).is_ok(),
2038 "folding drops the worktree, not the record"
2039 );
2040 }
2041
2042 fn init_repo(dir: &Path) {
2045 use crate::proc::Quiet as _;
2046 let run = |args: &[&str]| {
2047 let out = std::process::Command::new("git")
2048 .args(args)
2049 .current_dir(dir)
2050 .quiet()
2051 .output()
2052 .expect("spawn git");
2053 assert!(
2054 out.status.success(),
2055 "git {args:?} failed: {}",
2056 String::from_utf8_lossy(&out.stderr)
2057 );
2058 };
2059 run(&["init", "-b", "main"]);
2060 run(&["config", "user.name", "magi test"]);
2061 run(&["config", "user.email", "magi@example.com"]);
2062 std::fs::write(dir.join("README.md"), "# fixture\n").unwrap();
2063 run(&["add", "-A"]);
2064 run(&["commit", "-m", "init"]);
2065 }
2066
2067 async fn due_run_with_winner_worktree(
2071 runs: &Path,
2072 wt: &Path,
2073 repo: &Path,
2074 id: &str,
2075 status: RunStatus,
2076 ) -> (String, PathBuf) {
2077 use crate::run::{Candidate, Tally};
2078 use std::collections::BTreeMap;
2079
2080 let mut config = crate::config::Config::default();
2081 config.graph.worktree_root = Some(wt.to_path_buf());
2082 let mut state = RunState::new(
2083 repo.to_path_buf(),
2084 "main".to_owned(),
2085 "0000000000000000000000000000000000000000".to_owned(),
2086 String::new(),
2087 config,
2088 );
2089 state.id = id.to_owned();
2090 state.status = status;
2091 state.updated_at = ts("2026-08-01T00:00:00Z");
2092
2093 let winner_wt = state.worktree_root().join("cand-A");
2094 crate::git::worktree_add_branch(repo, &winner_wt, &format!("magi/{id}/A"), "main")
2095 .await
2096 .expect("winner worktree");
2097
2098 state.candidates.push(Candidate {
2099 index: 0,
2100 label: 'A',
2101 agent: "agent".to_owned(),
2102 branch: format!("magi/{id}/A"),
2103 worktree: winner_wt.clone(),
2104 summary: String::new(),
2105 stat: String::new(),
2106 files: 0,
2107 commits: 0,
2108 empty: false,
2109 failed: None,
2110 verified_noop: None,
2111 duration_ms: 0,
2112 folded: false,
2113 });
2114 state.tally = Some(Tally {
2115 first_choice: BTreeMap::from([('A', 1)]),
2116 borda: BTreeMap::new(),
2117 winner: 'A',
2118 rankings: 1,
2119 unanimous_initial: true,
2120 deliberated: false,
2121 changed_votes: 0,
2122 unanimous_final: true,
2123 tie_break: None,
2124 judges: 1,
2125 present: 1,
2126 quorum: 1,
2127 met_quorum: true,
2128 uncontested: None,
2129 });
2130 std::fs::create_dir_all(runs.join(id)).unwrap();
2131 std::fs::write(
2132 runs.join(id).join("run.json"),
2133 serde_json::to_string_pretty(&state).unwrap(),
2134 )
2135 .unwrap();
2136 (id.to_owned(), winner_wt)
2137 }
2138
2139 #[tokio::test]
2140 async fn fold_due_drops_a_superseded_runs_own_winner_worktree_but_keeps_a_readys() {
2141 let dir = tempfile::tempdir().unwrap();
2142 let runs = dir.path().join("runs");
2143 let wt = dir.path().join("wt");
2144 let repo = dir.path().join("repo");
2145 let home = dir.path().to_path_buf();
2146 let disk = Disk::default();
2147 let now = ts("2026-09-05T00:00:00Z");
2148 std::fs::create_dir_all(&repo).unwrap();
2149 init_repo(&repo);
2150
2151 let (superseded_id, superseded_wt) = due_run_with_winner_worktree(
2152 &runs,
2153 &wt,
2154 &repo,
2155 "20260801-000000-supw",
2156 RunStatus::Superseded,
2157 )
2158 .await;
2159 let (ready_id, ready_wt) = due_run_with_winner_worktree(
2166 &runs,
2167 &wt,
2168 &repo,
2169 "20260801-000000-rdyw",
2170 RunStatus::Ready,
2171 )
2172 .await;
2173
2174 let (folded, unreadable) = fold_due(&runs, &home, &wt, &disk, now)
2175 .await
2176 .expect("fold_due");
2177 assert_eq!(folded, 2);
2178 assert_eq!(unreadable, 0);
2179
2180 assert!(
2181 !superseded_wt.exists(),
2182 "a superseded run's own winner never lands anywhere else, so its worktree \
2183 must be dropped exactly like a merged run's"
2184 );
2185 assert!(
2186 ready_wt.exists(),
2187 "a still-ready run's winner may yet be merged by hand - folding must not \
2188 touch it"
2189 );
2190
2191 assert!(read_meta(&runs, &superseded_id).is_ok());
2193 assert!(read_meta(&runs, &ready_id).is_ok());
2194 }
2195}