Skip to main content

magi/
dupes.rs

1//! Duplicate-work detection at the moment work is filed.
2//!
3//! Two tasks once chased the same commit without knowing about each other: an
4//! open task and its open pull request already covered a change when a chat
5//! filed a second task that named the same branch and commit. Nothing at
6//! filing time said so. This module is that check, and it is deliberately
7//! small: **concrete identifiers only** - a branch name, a commit SHA, a pull
8//! request number - never text similarity. A false positive costs one
9//! `--force`; a false negative is the behaviour before this module existed.
10//!
11//! A *claim* is a thing an unfinished piece of work owns: the branches of its
12//! candidates and the pull request it opened. The claims come from
13//!
14//! - every task that is not `done`: its `review_branch`, and each run in
15//!   `Task::runs`;
16//! - every run that has not reached a terminal status;
17//! - every run that is terminal but whose pull request is still open (the
18//!   shape of the original accident: the run ended, the PR did not).
19//!
20//! Run records are read through a tolerant view rather than `RunState::load`:
21//! a schema bump must not turn into a silent false negative (the same lesson
22//! as `clean::fold_due`).
23//!
24//! Everything here is synchronous, read-only and takes its roots as
25//! arguments, so tests need no process-global home.
26
27use 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/// What kind of identifier matched.
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum Signal {
42    /// A branch name.
43    Branch,
44    /// A commit SHA.
45    Sha,
46    /// A pull request number or URL.
47    Pr,
48}
49
50/// Who owns the thing that matched.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub enum Owner {
53    /// A queued task.
54    Task,
55    /// A recorded run.
56    Run,
57    /// A pull request the forge reports open that no record here owns.
58    Pr,
59}
60
61/// One reason a new piece of work looks like one already under way.
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub struct Hit {
64    /// Task or run.
65    pub owner: Owner,
66    /// Full id of the task or run.
67    pub id: String,
68    /// The owner's status word (`queued`, `reviewing`, ...).
69    pub status: String,
70    /// What kind of identifier matched.
71    pub signal: Signal,
72    /// The identifier as it matched: a branch, a SHA, `#48`.
73    pub token: String,
74    /// How the owner is tied to it, e.g. `produced by its run c9eb`.
75    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/// The refusal: one or more [`Hit`]s and nothing filed.
102#[derive(Debug, Clone)]
103pub struct Duplicate(pub Vec<Hit>);
104
105impl Duplicate {
106    /// The refusal text, ending in how to override it.
107    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/// The slice of `run.json` this module needs. Every field is optional so a
131/// record from any schema still yields what it has.
132#[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    /// The run that took this one's worktree (and pull request) over.
147    #[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    /// An unknown status word counts as unfinished: claiming too much is the
174    /// cheap error here.
175    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
186/// Decides whether the pull request a run's record still calls open is in fact
187/// settled, so a stale record stops claiming it.
188///
189/// A record freezes the last state its land loop polled. A run that was handed
190/// over (`released_to`) or left blocked while another run landed the same pull
191/// request keeps saying `open` forever, so the record alone cannot be trusted
192/// to *keep* a claim alive. Two answers, cheapest first:
193///
194/// 1. a released run does not own the pull request any more: ownership follows
195///    `released_to`, and the successor's own record (read like any other run's)
196///    decides - provided that record is readable;
197/// 2. otherwise the forge is asked, once per number, at most
198///    [`MAX_FORGE_LOOKUPS`] numbers, and never again after one lookup failed
199///    (offline stays fast). Merged or closed means no claim; open, an error or
200///    a skipped lookup means the claim stands. Claiming too much is the cheap
201///    error.
202struct 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
210/// How many pull requests one invocation may put to the forge on behalf of
211/// stale records.
212const 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    /// `true` when `view`'s recorded-open pull request is known not to be this
226    /// run's to claim any more.
227    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/// Something an unfinished piece of work owns.
261#[derive(Debug, Clone)]
262struct Claim {
263    owner: Owner,
264    id: String,
265    status: String,
266    via: String,
267    branch: Option<String>,
268    /// Commit the branch forked from; without it a SHA cannot be judged.
269    base: Option<String>,
270    /// `(number, url)` of an open pull request.
271    pr: Option<(u64, String)>,
272}
273
274/// Look for work already in flight that `text` (and `review_branch`, for a
275/// review-only run) points at. `repo` is the repository the new work is for;
276/// `ignore_task` is a task being edited, whose own claims never count.
277pub 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/// [`check`] with the forge lookup supplied: `open_pr(repo, n)` says whether
298/// pull request `n` of `repo` is open (`Some(url)`) or not / unknown (`None`).
299/// It is asked only about numbers the text names that no local record already
300/// explained, so a PR magi never produced (opened by hand) still collides.
301/// `pr_state(repo, n)` is the forge's word on a pull request a run record calls
302/// open (`None` = unreadable, which keeps the claim); see [`Staleness`].
303#[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    // PRs the edited task's own runs opened: never a rival, forge or not.
324    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        // A terminal run held up only by a pull request that is not its own
351        // (any more) claims nothing at all.
352        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
430/// Ask the forge, via `gh`, whether PR `n` is open. Best effort and bounded:
431/// any failure (no `gh`, no remote, offline, a slow answer) is `None`, i.e.
432/// today's behaviour. `GH_REPO` is dropped so the PR is looked up in `repo`'s
433/// own remote, not whatever the environment points at.
434fn 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
441/// [`gh_open_pr`]'s question asked of a pull request a run record names: its
442/// lifecycle, or `None` when the forge cannot say.
443fn 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            // A merged or closed pull request is not in flight, even when the
506            // task that owns it is.
507            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
520/// The claims of one run, attributed to `owner` (the task, when there is one).
521///
522/// `live_pr` is false when the record's open pull request is known to be
523/// settled or somebody else's (see [`Staleness`]); it then names no PR.
524fn 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/// Repository identity: the canonical git common dir, so two worktrees of one
573/// repository are the same repository and a path spelled two ways is too.
574#[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
618/// Does `text` contain `branch` as a whole word? `.` and `/` before it are
619/// fine (`origin/magi/x/A`), a name character on either side is not, so
620/// `magi/27b2/A` does not match inside `magi/27b2/AB`.
621fn 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
647/// `#48`, `PR 48`, `PR #48`, `pull request 48`, and `.../pull/48` URLs.
648fn 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
693/// Hex words of 7..=40 chars in `text` that resolve to a commit in `repo`.
694fn 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
719/// Is `sha` reachable from `branch` but not from `base`? A SHA already in the
720/// base would match every branch there is.
721fn 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    /// A repo with `main` and a branch `magi/aaaa/A` one commit ahead.
757    /// Returns the fixture plus the (main, branch) tip SHAs.
758    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        // Hex that resolves to nothing is ignored.
859        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        // Changed deliberately: a terminal run's recorded-open PR now claims
894        // only while the forge says open or cannot be read (`run` passes an
895        // unreadable forge). A merged / closed answer or a released run no
896        // longer claims; see the `stale_open_*` tests below.
897        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        // Editing the task that owns the run never collides with itself.
945        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        // Closed / unknown / unnamed PRs never match.
998        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            // Same through an owning, still-live task: the PR is not its claim.
1062            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        // A successor that cannot be read decides nothing: claim stands.
1108        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}