1use std::collections::{BTreeSet, HashMap};
28use std::fmt;
29use std::path::{Path, PathBuf};
30use std::process::{Command, Stdio};
31
32use serde::Deserialize;
33
34use crate::land::PrLifecycle;
35use crate::proc::Quiet as _;
36use crate::queue::{Queue, Task, TaskStatus};
37use crate::run::RunStatus;
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum Signal {
42 Branch,
44 Sha,
46 Pr,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq)]
52pub enum Owner {
53 Task,
55 Run,
57 Pr,
59}
60
61#[derive(Debug, Clone, PartialEq, Eq)]
63pub struct Hit {
64 pub owner: Owner,
66 pub id: String,
68 pub status: String,
70 pub signal: Signal,
72 pub token: String,
74 pub via: String,
76}
77
78impl fmt::Display for Hit {
79 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
80 let kind = match self.owner {
81 Owner::Task => "task",
82 Owner::Run => "run",
83 Owner::Pr => "pull request",
84 };
85 let what = match self.signal {
86 Signal::Branch => "names branch",
87 Signal::Sha => "names commit",
88 Signal::Pr => "names pull request",
89 };
90 write!(
91 f,
92 "{kind} {} ({}): this work {what} {}, {}",
93 crate::queue::short(&self.id),
94 self.status,
95 self.token,
96 self.via
97 )
98 }
99}
100
101#[derive(Debug, Clone)]
103pub struct Duplicate(pub Vec<Hit>);
104
105impl Duplicate {
106 pub fn render(&self, override_hint: &str) -> String {
108 let mut out = String::from("this looks like work that is already in flight:");
109 for h in &self.0 {
110 out.push_str("\n - ");
111 out.push_str(&h.to_string());
112 }
113 out.push('\n');
114 out.push_str(override_hint);
115 out
116 }
117}
118
119impl fmt::Display for Duplicate {
120 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
121 f.write_str(&self.render(
122 "If it is not a duplicate, pass --force to file it anyway \
123 (an agent should report this to the operator instead).",
124 ))
125 }
126}
127
128impl std::error::Error for Duplicate {}
129
130#[derive(Debug, Default, Deserialize)]
133struct RunView {
134 #[serde(default)]
135 id: String,
136 #[serde(default)]
137 repo: PathBuf,
138 #[serde(default)]
139 status: String,
140 #[serde(default)]
141 base_commit: String,
142 #[serde(default)]
143 candidates: Vec<CandView>,
144 #[serde(default)]
145 pr: Option<PrView>,
146 #[serde(default)]
148 released_to: Option<String>,
149}
150
151#[derive(Debug, Default, Deserialize)]
152struct CandView {
153 #[serde(default)]
154 branch: String,
155}
156
157#[derive(Debug, Default, Deserialize)]
158struct PrView {
159 #[serde(default)]
160 url: String,
161 #[serde(default)]
162 number: u64,
163 #[serde(default)]
164 state: String,
165}
166
167impl RunView {
168 fn read(runs_root: &Path, id: &str) -> Option<Self> {
169 let raw = std::fs::read_to_string(runs_root.join(id).join("run.json")).ok()?;
170 serde_json::from_str(&raw).ok()
171 }
172
173 fn terminal(&self) -> bool {
176 serde_json::from_value::<RunStatus>(serde_json::Value::String(self.status.clone()))
177 .map(RunStatus::done)
178 .unwrap_or(false)
179 }
180
181 fn pr_open(&self) -> bool {
182 self.pr.as_ref().is_some_and(|p| p.state == "open")
183 }
184}
185
186struct Staleness<'a> {
203 runs_root: &'a Path,
204 lookup: &'a dyn Fn(&Path, u64) -> Option<PrLifecycle>,
205 cache: HashMap<u64, bool>,
206 asked: usize,
207 failed: bool,
208}
209
210const MAX_FORGE_LOOKUPS: usize = 5;
213
214impl<'a> Staleness<'a> {
215 fn new(runs_root: &'a Path, lookup: &'a dyn Fn(&Path, u64) -> Option<PrLifecycle>) -> Self {
216 Self {
217 runs_root,
218 lookup,
219 cache: HashMap::new(),
220 asked: 0,
221 failed: false,
222 }
223 }
224
225 fn pr_released(&mut self, repo: &Path, view: &RunView) -> bool {
228 let Some(pr) = view.pr.as_ref().filter(|p| p.state == "open") else {
229 return false;
230 };
231 if let Some(next) = &view.released_to
232 && next != &view.id
233 && RunView::read(self.runs_root, next).is_some()
234 {
235 return true;
236 }
237 if pr.number == 0 {
238 return false;
239 }
240 if let Some(known) = self.cache.get(&pr.number) {
241 return *known;
242 }
243 if self.failed || self.asked >= MAX_FORGE_LOOKUPS {
244 return false;
245 }
246 self.asked += 1;
247 let settled = match (self.lookup)(repo, pr.number) {
248 Some(PrLifecycle::Merged | PrLifecycle::Closed) => true,
249 Some(PrLifecycle::Open) => false,
250 None => {
251 self.failed = true;
252 false
253 }
254 };
255 self.cache.insert(pr.number, settled);
256 settled
257 }
258}
259
260#[derive(Debug, Clone)]
262struct Claim {
263 owner: Owner,
264 id: String,
265 status: String,
266 via: String,
267 branch: Option<String>,
268 base: Option<String>,
270 pr: Option<(u64, String)>,
272}
273
274pub fn check(
278 queue: &Queue,
279 runs_root: &Path,
280 repo: &Path,
281 text: &str,
282 review_branch: Option<&str>,
283 ignore_task: Option<&str>,
284) -> Vec<Hit> {
285 check_with(
286 queue,
287 runs_root,
288 repo,
289 text,
290 review_branch,
291 ignore_task,
292 &gh_open_pr,
293 &gh_pr_state,
294 )
295}
296
297#[allow(clippy::too_many_arguments)]
304pub fn check_with(
305 queue: &Queue,
306 runs_root: &Path,
307 repo: &Path,
308 text: &str,
309 review_branch: Option<&str>,
310 ignore_task: Option<&str>,
311 open_pr: &dyn Fn(&Path, u64) -> Option<String>,
312 pr_state: &dyn Fn(&Path, u64) -> Option<PrLifecycle>,
313) -> Vec<Hit> {
314 let mut stale = Staleness::new(runs_root, pr_state);
315 let mut idents = Idents::default();
316 let here = idents.of(repo);
317 let tasks = queue.list();
318 let own_runs: BTreeSet<String> = tasks
319 .iter()
320 .filter(|t| Some(t.id.as_str()) == ignore_task)
321 .flat_map(|t| t.runs.iter().cloned())
322 .collect();
323 let own_prs: BTreeSet<u64> = own_runs
325 .iter()
326 .filter_map(|id| RunView::read(runs_root, id))
327 .filter_map(|v| v.pr.map(|p| p.number))
328 .collect();
329
330 let mut claims: Vec<Claim> = Vec::new();
331 let mut from_task: BTreeSet<String> = BTreeSet::new();
332 for t in tasks
333 .iter()
334 .filter(|t| t.status != TaskStatus::Done && Some(t.id.as_str()) != ignore_task)
335 .filter(|t| idents.of(&t.repo) == here)
336 {
337 claims.extend(task_claims(t, runs_root, &mut from_task, repo, &mut stale));
338 }
339 for id in crate::run::list_ids_in(runs_root) {
340 if own_runs.contains(&id) {
341 continue;
342 }
343 let Some(view) = RunView::read(runs_root, &id) else {
344 continue;
345 };
346 if (view.terminal() && !view.pr_open()) || idents.of(&view.repo) != here {
347 continue;
348 }
349 let released = stale.pr_released(repo, &view);
350 if view.terminal() && released {
353 continue;
354 }
355 claims.extend(run_claims(
356 &view,
357 Owner::Run,
358 None,
359 "its own run",
360 !released,
361 ));
362 }
363
364 let mut hits: Vec<Hit> = Vec::new();
365 let mut push = |c: &Claim, signal: Signal, token: String| {
366 let hit = Hit {
367 owner: c.owner.clone(),
368 id: c.id.clone(),
369 status: c.status.clone(),
370 signal,
371 token,
372 via: c.via.clone(),
373 };
374 if !hits.contains(&hit) {
375 hits.push(hit);
376 }
377 };
378
379 let prs = pr_numbers(text);
380 let shas = sha_candidates(repo, text);
381 for c in &claims {
382 if let Some(b) = &c.branch {
383 if names_branch(text, b) || review_branch == Some(b.as_str()) {
384 push(c, Signal::Branch, b.clone());
385 }
386 if let Some(base) = &c.base {
387 for sha in &shas {
388 if on_branch_only(repo, sha, b, base) {
389 push(c, Signal::Sha, short_sha(sha));
390 }
391 }
392 }
393 }
394 if let Some((n, url)) = &c.pr {
395 if prs.contains(&Mention::Number(*n))
396 || prs.iter().any(|p| matches!(p, Mention::Url(u) if u == url))
397 {
398 push(c, Signal::Pr, format!("#{n}"));
399 }
400 }
401 }
402 let mut asked = BTreeSet::new();
403 for n in prs.iter().filter_map(|m| match m {
404 Mention::Number(n) => Some(*n),
405 Mention::Url(u) => u.rsplit('/').next().and_then(|d| d.parse().ok()),
406 }) {
407 let token = format!("#{n}");
408 if hits
409 .iter()
410 .any(|h| h.signal == Signal::Pr && h.token == token)
411 || own_prs.contains(&n)
412 || !asked.insert(n)
413 {
414 continue;
415 }
416 if let Some(url) = open_pr(repo, n) {
417 hits.push(Hit {
418 owner: Owner::Pr,
419 id: token.clone(),
420 status: "open".into(),
421 signal: Signal::Pr,
422 token,
423 via: format!("an open pull request with no run record here ({url})"),
424 });
425 }
426 }
427 hits
428}
429
430fn gh_open_pr(repo: &Path, n: u64) -> Option<String> {
435 let v = gh_pr_view(repo, n)?;
436 (v["state"] == "OPEN")
437 .then(|| v["url"].as_str().map(str::to_owned))
438 .flatten()
439}
440
441fn gh_pr_state(repo: &Path, n: u64) -> Option<PrLifecycle> {
444 match gh_pr_view(repo, n)?["state"].as_str()? {
445 "OPEN" => Some(PrLifecycle::Open),
446 "MERGED" => Some(PrLifecycle::Merged),
447 "CLOSED" => Some(PrLifecycle::Closed),
448 _ => None,
449 }
450}
451
452fn gh_pr_view(repo: &Path, n: u64) -> Option<serde_json::Value> {
453 let mut child = Command::new("gh")
454 .quiet()
455 .args(["pr", "view", &n.to_string(), "--json", "state,url"])
456 .current_dir(repo)
457 .env_remove("GH_REPO")
458 .env("GH_PROMPT_DISABLED", "1")
459 .stdin(Stdio::null())
460 .stdout(Stdio::piped())
461 .stderr(Stdio::null())
462 .spawn()
463 .ok()?;
464 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
465 loop {
466 match child.try_wait().ok()? {
467 Some(status) if status.success() => break,
468 Some(_) => return None,
469 None if std::time::Instant::now() >= deadline => {
470 let _ = child.kill();
471 let _ = child.wait();
472 return None;
473 }
474 None => std::thread::sleep(std::time::Duration::from_millis(50)),
475 }
476 }
477 let mut raw = String::new();
478 std::io::Read::read_to_string(&mut child.stdout.take()?, &mut raw).ok()?;
479 serde_json::from_str(&raw).ok()
480}
481
482fn task_claims(
483 t: &Task,
484 runs_root: &Path,
485 seen: &mut BTreeSet<String>,
486 repo: &Path,
487 stale: &mut Staleness<'_>,
488) -> Vec<Claim> {
489 let status = t.status.as_str().to_owned();
490 let mut out = Vec::new();
491 if let Some(b) = &t.review_branch {
492 out.push(Claim {
493 owner: Owner::Task,
494 id: t.id.clone(),
495 status: status.clone(),
496 via: "its review branch".into(),
497 branch: Some(b.clone()),
498 base: None,
499 pr: None,
500 });
501 }
502 for rid in &t.runs {
503 if let Some(view) = RunView::read(runs_root, rid) {
504 seen.insert(rid.clone());
505 let live_pr = !(view.terminal() && stale.pr_released(repo, &view));
508 out.extend(run_claims(
509 &view,
510 Owner::Task,
511 Some((&t.id, &status)),
512 &format!("produced by its run {}", crate::queue::short(rid)),
513 live_pr,
514 ));
515 }
516 }
517 out
518}
519
520fn run_claims(
525 view: &RunView,
526 owner: Owner,
527 task: Option<(&str, &str)>,
528 via: &str,
529 live_pr: bool,
530) -> Vec<Claim> {
531 let (id, status) = match task {
532 Some((id, status)) => (id.to_owned(), status.to_owned()),
533 None => (view.id.clone(), view.status.clone()),
534 };
535 let pr = view
536 .pr
537 .as_ref()
538 .filter(|p| live_pr && p.state == "open" && p.number > 0)
539 .map(|p| (p.number, p.url.clone()));
540 let via_pr = |extra: &str| match &pr {
541 Some((n, _)) => format!("{via} (PR #{n} open){extra}"),
542 None => format!("{via}{extra}"),
543 };
544 let mut out: Vec<Claim> = view
545 .candidates
546 .iter()
547 .filter(|c| !c.branch.is_empty())
548 .map(|c| Claim {
549 owner: owner.clone(),
550 id: id.clone(),
551 status: status.clone(),
552 via: via_pr(""),
553 branch: Some(c.branch.clone()),
554 base: (!view.base_commit.is_empty()).then(|| view.base_commit.clone()),
555 pr: None,
556 })
557 .collect();
558 if pr.is_some() {
559 out.push(Claim {
560 owner,
561 id,
562 status,
563 via: via_pr(""),
564 branch: None,
565 base: None,
566 pr,
567 });
568 }
569 out
570}
571
572#[derive(Default)]
575struct Idents(HashMap<PathBuf, PathBuf>);
576
577impl Idents {
578 fn of(&mut self, path: &Path) -> PathBuf {
579 self.0
580 .entry(path.to_path_buf())
581 .or_insert_with(|| {
582 git(
583 path,
584 &["rev-parse", "--path-format=absolute", "--git-common-dir"],
585 )
586 .map(PathBuf::from)
587 .and_then(|p| p.canonicalize().ok())
588 .or_else(|| path.canonicalize().ok())
589 .unwrap_or_else(|| path.to_path_buf())
590 })
591 .clone()
592 }
593}
594
595fn git(cwd: &Path, args: &[&str]) -> Option<String> {
596 let out = Command::new("git")
597 .quiet()
598 .args(args)
599 .current_dir(cwd)
600 .stdin(Stdio::null())
601 .stderr(Stdio::null())
602 .env("GIT_TERMINAL_PROMPT", "0")
603 .output()
604 .ok()?;
605 out.status
606 .success()
607 .then(|| String::from_utf8_lossy(&out.stdout).trim().to_owned())
608}
609
610fn git_ok(cwd: &Path, args: &[&str]) -> bool {
611 git(cwd, args).is_some()
612}
613
614fn is_ref_char(c: char) -> bool {
615 c.is_alphanumeric() || matches!(c, '_' | '-')
616}
617
618fn names_branch(text: &str, branch: &str) -> bool {
622 if branch.len() < 3 {
623 return false;
624 }
625 text.match_indices(branch).any(|(i, _)| {
626 let before = text[..i].chars().next_back();
627 let after = text[i + branch.len()..].chars().next();
628 let after_ok = match after {
629 None => true,
630 Some('/') => false,
631 Some('.') => !text[i + branch.len() + 1..]
632 .chars()
633 .next()
634 .is_some_and(is_ref_char),
635 Some(c) => !is_ref_char(c),
636 };
637 before.is_none_or(|c| !is_ref_char(c) && c != '.') && after_ok
638 })
639}
640
641#[derive(Debug, PartialEq, Eq)]
642enum Mention {
643 Number(u64),
644 Url(String),
645}
646
647fn pr_numbers(text: &str) -> Vec<Mention> {
649 let mut out = Vec::new();
650 let bytes = text.as_bytes();
651 let digits = |from: usize| -> Option<(u64, usize)> {
652 let n = text[from..].bytes().take_while(u8::is_ascii_digit).count();
653 (n > 0 && n < 10)
654 .then(|| text[from..from + n].parse().ok().map(|v| (v, from + n)))
655 .flatten()
656 };
657 for (i, _) in text.match_indices('#') {
658 if let Some((n, _)) = digits(i + 1) {
659 let word_before = i > 0 && is_ref_char(bytes[i - 1] as char);
660 if !word_before {
661 out.push(Mention::Number(n));
662 }
663 }
664 }
665 let lower = text.to_ascii_lowercase();
666 for key in ["pull request ", "pr "] {
667 for (i, _) in lower.match_indices(key) {
668 if i > 0 && is_ref_char(bytes[i - 1] as char) {
669 continue;
670 }
671 let from = i + key.len();
672 let from = if text[from..].starts_with('#') {
673 from + 1
674 } else {
675 from
676 };
677 if let Some((n, _)) = digits(from) {
678 out.push(Mention::Number(n));
679 }
680 }
681 }
682 for (i, _) in text.match_indices("/pull/") {
683 if let Some((_, end)) = digits(i + 6) {
684 let start = text[..i]
685 .rfind(|c: char| c.is_whitespace() || matches!(c, '(' | '<' | '"' | '\''))
686 .map_or(0, |p| p + 1);
687 out.push(Mention::Url(text[start..end].to_owned()));
688 }
689 }
690 out
691}
692
693fn sha_candidates(repo: &Path, text: &str) -> Vec<String> {
695 let mut seen = BTreeSet::new();
696 let mut out = Vec::new();
697 for word in text.split(|c: char| !c.is_ascii_alphanumeric()) {
698 if !(7..=40).contains(&word.len()) || !word.bytes().all(|b| b.is_ascii_hexdigit()) {
699 continue;
700 }
701 if seen.len() >= 16 || !seen.insert(word.to_ascii_lowercase()) {
702 continue;
703 }
704 if let Some(full) = git(
705 repo,
706 &[
707 "rev-parse",
708 "--verify",
709 "--quiet",
710 &format!("{word}^{{commit}}"),
711 ],
712 ) {
713 out.push(full);
714 }
715 }
716 out
717}
718
719fn on_branch_only(repo: &Path, sha: &str, branch: &str, base: &str) -> bool {
722 let tip = format!("{branch}^{{commit}}");
723 git_ok(repo, &["rev-parse", "--verify", "--quiet", &tip])
724 && git_ok(repo, &["merge-base", "--is-ancestor", sha, branch])
725 && !git_ok(repo, &["merge-base", "--is-ancestor", sha, base])
726}
727
728fn short_sha(sha: &str) -> String {
729 sha.chars().take(7).collect()
730}
731
732#[cfg(test)]
733mod tests {
734 use super::*;
735 use crate::queue::{Source, Task};
736
737 struct Fx {
738 _tmp: tempfile::TempDir,
739 repo: PathBuf,
740 runs: PathBuf,
741 q: Queue,
742 }
743
744 fn sh(cwd: &Path, args: &[&str]) -> String {
745 let out = Command::new("git")
746 .quiet()
747 .args(["-c", "user.name=t", "-c", "user.email=t@t"])
748 .args(args)
749 .current_dir(cwd)
750 .output()
751 .unwrap();
752 assert!(out.status.success(), "git {args:?}: {out:?}");
753 String::from_utf8_lossy(&out.stdout).trim().to_owned()
754 }
755
756 fn fx() -> (Fx, String, String) {
759 let tmp = tempfile::tempdir().unwrap();
760 let repo = tmp.path().join("repo");
761 std::fs::create_dir_all(&repo).unwrap();
762 sh(&repo, &["init", "-q", "-b", "main"]);
763 std::fs::write(repo.join("a"), "1").unwrap();
764 sh(&repo, &["add", "."]);
765 sh(&repo, &["commit", "-q", "-m", "base"]);
766 let base = sh(&repo, &["rev-parse", "HEAD"]);
767 sh(&repo, &["checkout", "-q", "-b", "magi/aaaa/A"]);
768 std::fs::write(repo.join("a"), "2").unwrap();
769 sh(&repo, &["commit", "-q", "-am", "work"]);
770 let tip = sh(&repo, &["rev-parse", "HEAD"]);
771 sh(&repo, &["checkout", "-q", "main"]);
772 let runs = tmp.path().join("runs");
773 std::fs::create_dir_all(&runs).unwrap();
774 let q = Queue::at(tmp.path().join("queue"));
775 (
776 Fx {
777 _tmp: tmp,
778 repo,
779 runs,
780 q,
781 },
782 base,
783 tip,
784 )
785 }
786
787 fn write_run(f: &Fx, id: &str, status: &str, base: &str, pr: Option<(u64, &str)>) {
788 let dir = f.runs.join(id);
789 std::fs::create_dir_all(&dir).unwrap();
790 let pr = pr.map(|(n, s)| {
791 serde_json::json!({"url": format!("https://github.com/o/r/pull/{n}"), "number": n, "state": s})
792 });
793 let v = serde_json::json!({
794 "schema": 999, "id": id, "repo": f.repo, "status": status,
795 "base_commit": base, "candidates": [{"branch": "magi/aaaa/A"}], "pr": pr,
796 });
797 std::fs::write(dir.join("run.json"), v.to_string()).unwrap();
798 }
799
800 fn file_task(f: &Fx, status: TaskStatus, runs: &[&str]) -> Task {
801 let mut t = Task::new("t".into(), "x".into(), f.repo.clone(), Source::Human);
802 t.status = status;
803 t.runs = runs.iter().map(|s| (*s).to_owned()).collect();
804 f.q.put(&mut t).unwrap();
805 t
806 }
807
808 const RID: &str = "20260901-100000-aaaa";
809
810 fn run(f: &Fx, text: &str, review: Option<&str>) -> Vec<Hit> {
811 check_with(
812 &f.q,
813 &f.runs,
814 &f.repo,
815 text,
816 review,
817 None,
818 &|_, _| None,
819 &|_, _| None,
820 )
821 }
822
823 #[test]
824 fn branch_matches_a_live_run_and_its_task() {
825 let (f, base, _) = fx();
826 write_run(&f, RID, "reviewing", &base, None);
827 let t = file_task(&f, TaskStatus::Running, &[RID]);
828 let hits = run(&f, "land magi/aaaa/A onto a fresh branch.", None);
829 assert!(hits.iter().any(|h| h.owner == Owner::Task
830 && h.id == t.id
831 && h.signal == Signal::Branch
832 && h.token == "magi/aaaa/A"));
833 assert!(hits.iter().any(|h| h.owner == Owner::Run && h.id == RID));
834 let msg = Duplicate(hits).to_string();
835 assert!(
836 msg.contains("--force") && msg.contains("magi/aaaa/A"),
837 "{msg}"
838 );
839 }
840
841 #[test]
842 fn branch_must_match_whole_word() {
843 let (f, base, _) = fx();
844 write_run(&f, RID, "reviewing", &base, None);
845 assert!(run(&f, "see magi/aaaa/AB and magi/aaaa/A/x", None).is_empty());
846 }
847
848 #[test]
849 fn sha_on_the_branch_matches_but_one_in_base_does_not() {
850 let (f, base, tip) = fx();
851 write_run(&f, RID, "reviewing", &base, None);
852 let hits = run(&f, &format!("land commit {} please", &tip[..8]), None);
853 assert!(
854 hits.iter().any(|h| h.signal == Signal::Sha && h.id == RID),
855 "{hits:?}"
856 );
857 assert!(run(&f, &format!("see {}", &base[..9]), None).is_empty());
858 assert!(run(&f, "deadbeef and 1234567", None).is_empty());
860 }
861
862 #[test]
863 fn pr_number_matches_in_every_spelling() {
864 let (f, base, _) = fx();
865 write_run(&f, RID, "ready", &base, Some((48, "open")));
866 for text in [
867 "finish #48",
868 "PR 48 is stale",
869 "pr #48",
870 "pull request 48",
871 "https://github.com/o/r/pull/48",
872 ] {
873 let hits = run(&f, text, None);
874 assert!(
875 hits.iter().any(|h| h.signal == Signal::Pr),
876 "{text}: {hits:?}"
877 );
878 }
879 assert!(run(&f, "see #480 and PR 4 and issue48", None).is_empty());
880 }
881
882 #[test]
883 fn terminal_runs_and_done_tasks_do_not_match() {
884 let (f, base, tip) = fx();
885 write_run(&f, RID, "merged", &base, Some((48, "merged")));
886 file_task(&f, TaskStatus::Done, &[RID]);
887 let text = format!("magi/aaaa/A {} #48", &tip[..8]);
888 assert!(run(&f, &text, None).is_empty());
889 }
890
891 #[test]
892 fn terminal_run_with_open_pr_or_open_task_still_claims() {
893 let (f, base, _) = fx();
898 write_run(&f, RID, "ready", &base, Some((48, "open")));
899 assert!(!run(&f, "magi/aaaa/A", None).is_empty());
900 let (g, base, _) = fx();
901 write_run(&g, RID, "ready", &base, None);
902 assert!(run(&g, "magi/aaaa/A", None).is_empty());
903 let t = file_task(&g, TaskStatus::Held, &[RID]);
904 let hits = run(&g, "magi/aaaa/A", None);
905 assert!(hits.iter().any(|h| h.id == t.id), "{hits:?}");
906 }
907
908 #[test]
909 fn review_only_matches_a_branch_a_live_task_owns() {
910 let (f, base, _) = fx();
911 write_run(&f, RID, "ready", &base, None);
912 assert!(run(&f, "", Some("magi/aaaa/A")).is_empty());
913 let mut t = file_task(&f, TaskStatus::Queued, &[]);
914 t.review_branch = Some("magi/aaaa/A".into());
915 f.q.put(&mut t).unwrap();
916 let hits = run(&f, "", Some("magi/aaaa/A"));
917 assert!(
918 hits.iter()
919 .any(|h| h.id == t.id && h.signal == Signal::Branch)
920 );
921 }
922
923 #[test]
924 fn other_repository_and_edited_task_do_not_match() {
925 let (f, base, _) = fx();
926 write_run(&f, RID, "reviewing", &base, None);
927 let t = file_task(&f, TaskStatus::Running, &[RID]);
928 let other = f._tmp.path().join("other");
929 std::fs::create_dir_all(&other).unwrap();
930 sh(&other, &["init", "-q"]);
931 assert!(
932 check_with(
933 &f.q,
934 &f.runs,
935 &other,
936 "magi/aaaa/A",
937 None,
938 None,
939 &|_, _| None,
940 &|_, _| None,
941 )
942 .is_empty()
943 );
944 assert!(
946 check_with(
947 &f.q,
948 &f.runs,
949 &f.repo,
950 "magi/aaaa/A",
951 None,
952 Some(&t.id),
953 &|_, _| None,
954 &|_, _| None,
955 )
956 .is_empty()
957 );
958 }
959
960 #[test]
961 fn a_worktree_is_the_same_repository() {
962 let (f, base, _) = fx();
963 write_run(&f, RID, "reviewing", &base, None);
964 let wt = f._tmp.path().join("wt");
965 sh(
966 &f.repo,
967 &["worktree", "add", "-q", wt.to_str().unwrap(), "-b", "other"],
968 );
969 assert!(
970 !check_with(
971 &f.q,
972 &f.runs,
973 &wt,
974 "magi/aaaa/A",
975 None,
976 None,
977 &|_, _| None,
978 &|_, _| None,
979 )
980 .is_empty()
981 );
982 }
983
984 #[test]
985 fn an_open_pr_without_a_run_record_matches_through_the_forge() {
986 let (f, _, _) = fx();
987 let open = |_: &Path, n: u64| (n == 48).then(|| "https://example.test/pull/48".to_owned());
988 let hit = |text: &str| {
989 check_with(&f.q, &f.runs, &f.repo, text, None, None, &open, &|_, _| {
990 None
991 })
992 };
993 let hits = hit("finish PR #48");
994 assert_eq!(hits.len(), 1, "{hits:?}");
995 assert_eq!(hits[0].owner, Owner::Pr);
996 assert!(hits[0].to_string().contains("#48"));
997 assert!(hit("finish PR #49").is_empty());
999 assert!(hit("finish the work").is_empty());
1000 }
1001
1002 #[test]
1003 fn a_forge_hit_does_not_repeat_a_pr_a_run_already_explains() {
1004 let (f, base, _) = fx();
1005 write_run(&f, RID, "ready", &base, Some((48, "open")));
1006 let open = |_: &Path, _: u64| Some("u".to_owned());
1007 let hits = check_with(&f.q, &f.runs, &f.repo, "#48", None, None, &open, &|_, _| {
1008 None
1009 });
1010 assert!(hits.iter().all(|h| h.owner != Owner::Pr), "{hits:?}");
1011 assert!(!hits.is_empty());
1012 }
1013
1014 #[test]
1015 fn an_edited_tasks_own_open_pr_is_not_a_forge_hit() {
1016 let (f, base, _) = fx();
1017 write_run(&f, RID, "ready", &base, Some((48, "open")));
1018 let t = file_task(&f, TaskStatus::Running, &[RID]);
1019 let open = |_: &Path, _: u64| Some("u".to_owned());
1020 let hits = check_with(
1021 &f.q,
1022 &f.runs,
1023 &f.repo,
1024 "#48",
1025 None,
1026 Some(&t.id),
1027 &open,
1028 &|_, _| None,
1029 );
1030 assert!(hits.is_empty(), "{hits:?}");
1031 }
1032
1033 fn with_forge(f: &Fx, text: &str, state: Option<PrLifecycle>) -> Vec<Hit> {
1034 check_with(
1035 &f.q,
1036 &f.runs,
1037 &f.repo,
1038 text,
1039 None,
1040 None,
1041 &|_, _| None,
1042 &move |_, _| state,
1043 )
1044 }
1045
1046 fn release(f: &Fx, id: &str, to: &str) {
1047 let path = f.runs.join(id).join("run.json");
1048 let mut v: serde_json::Value =
1049 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
1050 v["released_to"] = serde_json::json!(to);
1051 std::fs::write(path, v.to_string()).unwrap();
1052 }
1053
1054 #[test]
1055 fn stale_open_pr_that_the_forge_says_is_merged_or_closed_does_not_claim() {
1056 for state in [PrLifecycle::Merged, PrLifecycle::Closed] {
1057 let (f, base, _) = fx();
1058 write_run(&f, RID, "superseded", &base, Some((48, "open")));
1059 assert!(with_forge(&f, "follow up on #48", Some(state)).is_empty());
1060 assert!(with_forge(&f, "magi/aaaa/A", Some(state)).is_empty());
1061 file_task(&f, TaskStatus::Held, &[RID]);
1063 let hits = with_forge(&f, "follow up on #48", Some(state));
1064 assert!(hits.iter().all(|h| h.signal != Signal::Pr), "{hits:?}");
1065 }
1066 }
1067
1068 #[test]
1069 fn stale_open_pr_with_an_unreadable_forge_still_claims() {
1070 let (f, base, _) = fx();
1071 write_run(&f, RID, "superseded", &base, Some((48, "open")));
1072 let hits = with_forge(&f, "follow up on #48", None);
1073 assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1074 }
1075
1076 #[test]
1077 fn a_genuinely_open_pr_still_claims() {
1078 let (f, base, _) = fx();
1079 write_run(&f, RID, "blocked", &base, Some((48, "open")));
1080 let hits = with_forge(&f, "follow up on #48", Some(PrLifecycle::Open));
1081 assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1082 }
1083
1084 #[test]
1085 fn a_released_run_defers_to_its_successor_without_asking_the_forge() {
1086 let (f, base, _) = fx();
1087 let next = "20260901-110000-bbbb";
1088 write_run(&f, RID, "superseded", &base, Some((48, "open")));
1089 write_run(&f, next, "merged", &base, Some((48, "merged")));
1090 release(&f, RID, next);
1091 let asked = std::cell::Cell::new(0);
1092 let hits = check_with(
1093 &f.q,
1094 &f.runs,
1095 &f.repo,
1096 "follow up on #48",
1097 None,
1098 None,
1099 &|_, _| None,
1100 &|_, _| {
1101 asked.set(asked.get() + 1);
1102 None
1103 },
1104 );
1105 assert!(hits.is_empty(), "{hits:?}");
1106 assert_eq!(asked.get(), 0);
1107 std::fs::remove_dir_all(f.runs.join(next)).unwrap();
1109 let hits = with_forge(&f, "follow up on #48", None);
1110 assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1111 }
1112
1113 #[test]
1114 fn forge_lookups_are_cached_and_stop_after_a_failure() {
1115 let (f, base, _) = fx();
1116 for (i, id) in ["20260901-100000-aaa1", "20260901-100000-aaa2"]
1117 .iter()
1118 .enumerate()
1119 {
1120 write_run(&f, id, "blocked", &base, Some((48 + i as u64, "open")));
1121 }
1122 let asked = std::cell::Cell::new(0);
1123 check_with(
1124 &f.q,
1125 &f.runs,
1126 &f.repo,
1127 "x",
1128 None,
1129 None,
1130 &|_, _| None,
1131 &|_, _| {
1132 asked.set(asked.get() + 1);
1133 None
1134 },
1135 );
1136 assert_eq!(
1137 asked.get(),
1138 1,
1139 "an unreadable forge is asked once, not per PR"
1140 );
1141 }
1142}