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::future::Future;
30use std::path::{Path, PathBuf};
31use std::pin::Pin;
32use std::process::{Command, Stdio};
33use std::time::{Duration, Instant};
34
35use anyhow::{Context as _, Result, bail};
36use serde::Deserialize;
37
38use crate::agent;
39use crate::config::Config;
40use crate::prompt;
41
42use crate::land::PrLifecycle;
43use crate::proc::Quiet as _;
44use crate::queue::{Queue, Task, TaskStatus};
45use crate::run::RunStatus;
46
47/// What kind of identifier matched.
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum Signal {
50    /// A branch name.
51    Branch,
52    /// A commit SHA.
53    Sha,
54    /// A pull request number or URL.
55    Pr,
56}
57
58/// Who owns the thing that matched.
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub enum Owner {
61    /// A queued task.
62    Task,
63    /// A recorded run.
64    Run,
65    /// A pull request the forge reports open that no record here owns.
66    Pr,
67}
68
69/// One reason a new piece of work looks like one already under way.
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct Hit {
72    /// Task or run.
73    pub owner: Owner,
74    /// Full id of the task or run.
75    pub id: String,
76    /// The owner's status word (`queued`, `reviewing`, ...).
77    pub status: String,
78    /// What kind of identifier matched.
79    pub signal: Signal,
80    /// The identifier as it matched: a branch, a SHA, `#48`.
81    pub token: String,
82    /// How the owner is tied to it, e.g. `produced by its run c9eb`.
83    pub via: String,
84}
85
86impl fmt::Display for Hit {
87    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
88        let kind = match self.owner {
89            Owner::Task => "task",
90            Owner::Run => "run",
91            Owner::Pr => "pull request",
92        };
93        let what = match self.signal {
94            Signal::Branch => "names branch",
95            Signal::Sha => "names commit",
96            Signal::Pr => "names pull request",
97        };
98        write!(
99            f,
100            "{kind} {} ({}): this work {what} {}, {}",
101            crate::queue::short(&self.id),
102            self.status,
103            self.token,
104            self.via
105        )
106    }
107}
108
109/// The refusal: one or more [`Hit`]s and nothing filed.
110#[derive(Debug, Clone)]
111pub struct Duplicate {
112    /// What matched.
113    pub hits: Vec<Hit>,
114    /// The judge's word, when one was asked: a ready-made line for the
115    /// refusal text. `None` when no judge ran.
116    pub judge: Option<String>,
117}
118
119impl Duplicate {
120    /// A refusal on the identifier match alone.
121    pub fn new(hits: Vec<Hit>) -> Self {
122        Self { hits, judge: None }
123    }
124
125    /// The refusal text, ending in how to override it.
126    pub fn render(&self, override_hint: &str) -> String {
127        let mut out = String::from("this looks like work that is already in flight:");
128        for h in &self.hits {
129            out.push_str("\n  - ");
130            out.push_str(&h.to_string());
131        }
132        if let Some(j) = &self.judge {
133            out.push('\n');
134            out.push_str(j);
135        }
136        out.push('\n');
137        out.push_str(override_hint);
138        out
139    }
140}
141
142impl fmt::Display for Duplicate {
143    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
144        f.write_str(&self.render(
145            "If it is not a duplicate, pass --force to file it anyway \
146             (an agent should report this to the operator instead).",
147        ))
148    }
149}
150
151impl std::error::Error for Duplicate {}
152
153/// What the judge made of a match.
154#[derive(Debug, Clone, Copy, PartialEq, Eq)]
155pub enum Ruling {
156    /// The new work would do work on the same branch / pull request / commit.
157    Owns,
158    /// The new work only cites it as context.
159    Mentions,
160    /// The judge could not tell.
161    Unsure,
162}
163
164/// The judge's answer, with the agent that gave it.
165#[derive(Debug, Clone, PartialEq, Eq)]
166pub struct Judgement {
167    /// The verdict.
168    pub ruling: Ruling,
169    /// One line, as the agent wrote it.
170    pub reason: String,
171    /// Id of the agent that answered.
172    pub agent: String,
173}
174
175/// A judge: given the instruction and the matched claims, rules on them.
176/// Owned arguments, so a boxed future needs no borrow.
177pub type JudgeFuture = Pin<Box<dyn Future<Output = Result<Judgement>> + Send>>;
178
179/// Whole-chain wall-clock budget for the judge call.
180const JUDGE_BUDGET: Duration = Duration::from_secs(120);
181/// Per-agent cap inside that budget.
182const JUDGE_TURN: Duration = Duration::from_secs(90);
183
184/// Read the judge's reply. Strict: anything but one of the three rulings with a
185/// non-empty reason is an error, which the caller treats as a refusal.
186fn parse_ruling(text: &str) -> Result<(Ruling, String)> {
187    #[derive(Deserialize)]
188    struct Raw {
189        ruling: String,
190        reason: String,
191    }
192    // The whole reply must be one object (optionally in a code fence): hunting
193    // for an object inside prose or a second answer could recover an approval
194    // from a reply that is not itself a single valid ruling.
195    let mut body = text.trim();
196    if let Some(rest) = body.strip_prefix("```") {
197        let rest = rest.strip_prefix("json").unwrap_or(rest);
198        body = rest.trim().strip_suffix("```").unwrap_or(rest).trim();
199    }
200    let raw: Raw =
201        serde_json::from_str(body).context("the judge's reply is not a single JSON object")?;
202    let ruling = match raw.ruling.trim().to_ascii_lowercase().as_str() {
203        "owns" => Ruling::Owns,
204        "mentions" => Ruling::Mentions,
205        "unsure" => Ruling::Unsure,
206        other => bail!("unknown ruling `{other}`"),
207    };
208    let reason = raw.reason.split_whitespace().collect::<Vec<_>>().join(" ");
209    if reason.is_empty() {
210        bail!("the judge gave no reason");
211    }
212    Ok((ruling, reason.chars().take(300).collect()))
213}
214
215/// Ask the `[roles] chatter` chain (or the default agent) once. Each agent is
216/// tried at most once; only an error, quota or unusable answer moves on. A
217/// usable answer that does not parse, or that rules `owns` / `unsure`, ends
218/// the chain.
219pub async fn chain_judge(
220    cfg: &Config,
221    repo: &Path,
222    instruction: String,
223    hits: Vec<Hit>,
224) -> Result<Judgement> {
225    let chain = agent::pick_chain(
226        &cfg.agents,
227        cfg.roles.chatter.as_ref(),
228        &agent::installed,
229        "dupes judge",
230    )?;
231    let claims: Vec<String> = hits.iter().map(ToString::to_string).collect();
232    let body = prompt::dupes_judge(&instruction, &claims);
233    let artifacts = std::env::temp_dir().join(format!("magi-dupes-{:016x}", crate::rng::entropy()));
234    let started = Instant::now();
235    let mut last: anyhow::Error = anyhow::anyhow!("no judge agent ran");
236    let mut result = None;
237    for spec in &chain {
238        let left = JUDGE_BUDGET.saturating_sub(started.elapsed());
239        if left.is_zero() {
240            break;
241        }
242        let mut seat = agent::SeatState::new("dupes", &spec.id, crate::rng::entropy());
243        let inv = agent::Invocation {
244            cwd: repo,
245            prompt: &body,
246            timeout: left.min(JUDGE_TURN),
247            allow_write: false,
248            sessions: false,
249            artifacts: &artifacts,
250            stem: &format!("judge-{}", spec.id),
251            run: "dupes",
252            node: "dupes",
253            cache_dir: None,
254            attachments: &[],
255            writable: &[],
256        };
257        let out = agent::invoke(spec, &mut seat, &inv).await;
258        if agent::chain_advances(&out) {
259            last = match out {
260                Err(e) => e.context(format!("judge `{}` failed", spec.id)),
261                Ok(o) => anyhow::anyhow!(
262                    "judge `{}` gave no usable reply (exit {:?}, timed out {}, quota {})",
263                    spec.id,
264                    o.exit_code,
265                    o.timed_out,
266                    o.quota_exhausted()
267                ),
268            };
269            continue;
270        }
271        result = Some(
272            out.and_then(|o| parse_ruling(&o.text))
273                .map(|(ruling, reason)| Judgement {
274                    ruling,
275                    reason,
276                    agent: spec.id.clone(),
277                }),
278        );
279        break;
280    }
281    let _ = std::fs::remove_dir_all(&artifacts);
282    result.unwrap_or(Err(last))
283}
284
285/// The one decision point every caller shares. `hits` is [`check`]'s output:
286/// empty passes without the judge being asked. Otherwise the judge is asked
287/// once, and only a clean `mentions` lets the work through; `owns`, `unsure`
288/// and every failure refuse, carrying the judge's reason when there is one.
289///
290/// For a review-only request `text` is empty, so `review_branch` is put in
291/// front of the judge as the thing the work is about.
292pub async fn screen(
293    hits: Vec<Hit>,
294    text: &str,
295    review_branch: Option<&str>,
296    judge: &(dyn Fn(String, Vec<Hit>) -> JudgeFuture + Sync),
297) -> Result<(), Duplicate> {
298    if hits.is_empty() {
299        return Ok(());
300    }
301    // A judge that has not seen the whole text cannot clear it: the part it
302    // missed may be the part that does the work.
303    if text.chars().count() > prompt::DUPES_JUDGE_MAX_CHARS {
304        return Err(Duplicate {
305            hits,
306            judge: Some(format!(
307                "judge could not decide: the text is longer than {} characters, \
308                 too long to judge in full",
309                prompt::DUPES_JUDGE_MAX_CHARS
310            )),
311        });
312    }
313    let mut subject = text.to_owned();
314    if let Some(b) = review_branch {
315        if !subject.is_empty() {
316            subject.push_str("\n\n");
317        }
318        subject.push_str(&format!(
319            "(This is a review-only request for branch `{b}`: it would do work on that branch.)"
320        ));
321    }
322    let note = match judge(subject, hits.clone()).await {
323        Ok(j) if j.ruling == Ruling::Mentions => {
324            tracing::info!(
325                agent = %j.agent,
326                reason = %j.reason,
327                hits = hits.len(),
328                "duplicate check: the judge says the work only mentions what is in flight"
329            );
330            return Ok(());
331        }
332        Ok(j) => {
333            let word = if j.ruling == Ruling::Owns {
334                "owns"
335            } else {
336                "unsure"
337            };
338            format!("judge ({}): {word} - {}", j.agent, j.reason)
339        }
340        Err(e) => format!("judge could not decide: {e:#}"),
341    };
342    Err(Duplicate {
343        hits,
344        judge: Some(note),
345    })
346}
347
348/// [`screen`] with the production judge: the `[roles] chatter` chain of `cfg`.
349/// No config (`None`) means no judge, so a hit refuses as it always did.
350pub async fn screen_with_config(
351    hits: Vec<Hit>,
352    text: &str,
353    review_branch: Option<&str>,
354    repo: &Path,
355    cfg: Option<&Config>,
356) -> Result<(), Duplicate> {
357    let judge = |instruction: String, hits: Vec<Hit>| -> JudgeFuture {
358        let cfg = cfg.cloned();
359        let repo = repo.to_path_buf();
360        Box::pin(async move {
361            match cfg {
362                Some(cfg) => chain_judge(&cfg, &repo, instruction, hits).await,
363                None => bail!("no readable configuration to resolve a judge agent from"),
364            }
365        })
366    };
367    screen(hits, text, review_branch, &judge).await
368}
369
370/// The slice of `run.json` this module needs. Every field is optional so a
371/// record from any schema still yields what it has.
372#[derive(Debug, Default, Deserialize)]
373struct RunView {
374    #[serde(default)]
375    id: String,
376    #[serde(default)]
377    repo: PathBuf,
378    #[serde(default)]
379    status: String,
380    #[serde(default)]
381    base_commit: String,
382    #[serde(default)]
383    candidates: Vec<CandView>,
384    #[serde(default)]
385    pr: Option<PrView>,
386    /// The run that took this one's worktree (and pull request) over.
387    #[serde(default)]
388    released_to: Option<String>,
389}
390
391#[derive(Debug, Default, Deserialize)]
392struct CandView {
393    #[serde(default)]
394    branch: String,
395}
396
397#[derive(Debug, Default, Deserialize)]
398struct PrView {
399    #[serde(default)]
400    url: String,
401    #[serde(default)]
402    number: u64,
403    #[serde(default)]
404    state: String,
405}
406
407impl RunView {
408    fn read(runs_root: &Path, id: &str) -> Option<Self> {
409        let raw = std::fs::read_to_string(runs_root.join(id).join("run.json")).ok()?;
410        serde_json::from_str(&raw).ok()
411    }
412
413    /// An unknown status word counts as unfinished: claiming too much is the
414    /// cheap error here.
415    fn terminal(&self) -> bool {
416        serde_json::from_value::<RunStatus>(serde_json::Value::String(self.status.clone()))
417            .map(RunStatus::done)
418            .unwrap_or(false)
419    }
420
421    fn pr_open(&self) -> bool {
422        self.pr.as_ref().is_some_and(|p| p.state == "open")
423    }
424}
425
426/// Decides whether the pull request a run's record still calls open is in fact
427/// settled, so a stale record stops claiming it.
428///
429/// A record freezes the last state its land loop polled. A run that was handed
430/// over (`released_to`) or left blocked while another run landed the same pull
431/// request keeps saying `open` forever, so the record alone cannot be trusted
432/// to *keep* a claim alive. Two answers, cheapest first:
433///
434/// 1. a released run does not own the pull request any more: ownership follows
435///    `released_to`, and the successor's own record (read like any other run's)
436///    decides - provided that record is readable;
437/// 2. otherwise the forge is asked, once per number, at most
438///    [`MAX_FORGE_LOOKUPS`] numbers, and never again after one lookup failed
439///    (offline stays fast). Merged or closed means no claim; open, an error or
440///    a skipped lookup means the claim stands. Claiming too much is the cheap
441///    error.
442struct Staleness<'a> {
443    runs_root: &'a Path,
444    lookup: &'a dyn Fn(&Path, u64) -> Option<PrLifecycle>,
445    cache: HashMap<u64, bool>,
446    asked: usize,
447    failed: bool,
448}
449
450/// How many pull requests one invocation may put to the forge on behalf of
451/// stale records.
452const MAX_FORGE_LOOKUPS: usize = 5;
453
454impl<'a> Staleness<'a> {
455    fn new(runs_root: &'a Path, lookup: &'a dyn Fn(&Path, u64) -> Option<PrLifecycle>) -> Self {
456        Self {
457            runs_root,
458            lookup,
459            cache: HashMap::new(),
460            asked: 0,
461            failed: false,
462        }
463    }
464
465    /// `true` when `view`'s recorded-open pull request is known not to be this
466    /// run's to claim any more.
467    fn pr_released(&mut self, repo: &Path, view: &RunView) -> bool {
468        let Some(pr) = view.pr.as_ref().filter(|p| p.state == "open") else {
469            return false;
470        };
471        if let Some(next) = &view.released_to
472            && next != &view.id
473            && RunView::read(self.runs_root, next).is_some()
474        {
475            return true;
476        }
477        if pr.number == 0 {
478            return false;
479        }
480        if let Some(known) = self.cache.get(&pr.number) {
481            return *known;
482        }
483        if self.failed || self.asked >= MAX_FORGE_LOOKUPS {
484            return false;
485        }
486        self.asked += 1;
487        let settled = match (self.lookup)(repo, pr.number) {
488            Some(PrLifecycle::Merged | PrLifecycle::Closed) => true,
489            Some(PrLifecycle::Open) => false,
490            None => {
491                self.failed = true;
492                false
493            }
494        };
495        self.cache.insert(pr.number, settled);
496        settled
497    }
498}
499
500/// Something an unfinished piece of work owns.
501#[derive(Debug, Clone)]
502struct Claim {
503    owner: Owner,
504    id: String,
505    status: String,
506    via: String,
507    branch: Option<String>,
508    /// Commit the branch forked from; without it a SHA cannot be judged.
509    base: Option<String>,
510    /// `(number, url)` of an open pull request.
511    pr: Option<(u64, String)>,
512}
513
514/// Look for work already in flight that `text` (and `review_branch`, for a
515/// review-only run) points at. `repo` is the repository the new work is for;
516/// `ignore_task` is a task being edited, whose own claims never count.
517pub fn check(
518    queue: &Queue,
519    runs_root: &Path,
520    repo: &Path,
521    text: &str,
522    review_branch: Option<&str>,
523    ignore_task: Option<&str>,
524) -> Vec<Hit> {
525    check_with(
526        queue,
527        runs_root,
528        repo,
529        text,
530        review_branch,
531        ignore_task,
532        &gh_open_pr,
533        &gh_pr_state,
534    )
535}
536
537/// [`check`] with the forge lookup supplied: `open_pr(repo, n)` says whether
538/// pull request `n` of `repo` is open (`Some(url)`) or not / unknown (`None`).
539/// It is asked only about numbers the text names that no local record already
540/// explained, so a PR magi never produced (opened by hand) still collides.
541/// `pr_state(repo, n)` is the forge's word on a pull request a run record calls
542/// open (`None` = unreadable, which keeps the claim); see [`Staleness`].
543#[allow(clippy::too_many_arguments)]
544pub fn check_with(
545    queue: &Queue,
546    runs_root: &Path,
547    repo: &Path,
548    text: &str,
549    review_branch: Option<&str>,
550    ignore_task: Option<&str>,
551    open_pr: &dyn Fn(&Path, u64) -> Option<String>,
552    pr_state: &dyn Fn(&Path, u64) -> Option<PrLifecycle>,
553) -> Vec<Hit> {
554    let mut stale = Staleness::new(runs_root, pr_state);
555    let mut idents = Idents::default();
556    let here = idents.of(repo);
557    let tasks = queue.list();
558    let own_runs: BTreeSet<String> = tasks
559        .iter()
560        .filter(|t| Some(t.id.as_str()) == ignore_task)
561        .flat_map(|t| t.runs.iter().cloned())
562        .collect();
563    // PRs the edited task's own runs opened: never a rival, forge or not.
564    let own_prs: BTreeSet<u64> = own_runs
565        .iter()
566        .filter_map(|id| RunView::read(runs_root, id))
567        .filter_map(|v| v.pr.map(|p| p.number))
568        .collect();
569
570    let mut claims: Vec<Claim> = Vec::new();
571    let mut from_task: BTreeSet<String> = BTreeSet::new();
572    for t in tasks
573        .iter()
574        .filter(|t| t.status != TaskStatus::Done && Some(t.id.as_str()) != ignore_task)
575        .filter(|t| idents.of(&t.repo) == here)
576    {
577        claims.extend(task_claims(t, runs_root, &mut from_task, repo, &mut stale));
578    }
579    for id in crate::run::list_ids_in(runs_root) {
580        if own_runs.contains(&id) {
581            continue;
582        }
583        let Some(view) = RunView::read(runs_root, &id) else {
584            continue;
585        };
586        if (view.terminal() && !view.pr_open()) || idents.of(&view.repo) != here {
587            continue;
588        }
589        let released = stale.pr_released(repo, &view);
590        // A terminal run held up only by a pull request that is not its own
591        // (any more) claims nothing at all.
592        if view.terminal() && released {
593            continue;
594        }
595        claims.extend(run_claims(
596            &view,
597            Owner::Run,
598            None,
599            "its own run",
600            !released,
601        ));
602    }
603
604    let mut hits: Vec<Hit> = Vec::new();
605    let mut push = |c: &Claim, signal: Signal, token: String| {
606        let hit = Hit {
607            owner: c.owner.clone(),
608            id: c.id.clone(),
609            status: c.status.clone(),
610            signal,
611            token,
612            via: c.via.clone(),
613        };
614        if !hits.contains(&hit) {
615            hits.push(hit);
616        }
617    };
618
619    let prs = pr_numbers(text);
620    let shas = sha_candidates(repo, text);
621    for c in &claims {
622        if let Some(b) = &c.branch {
623            if names_branch(text, b) || review_branch == Some(b.as_str()) {
624                push(c, Signal::Branch, b.clone());
625            }
626            if let Some(base) = &c.base {
627                for sha in &shas {
628                    if on_branch_only(repo, sha, b, base) {
629                        push(c, Signal::Sha, short_sha(sha));
630                    }
631                }
632            }
633        }
634        if let Some((n, url)) = &c.pr {
635            if prs.contains(&Mention::Number(*n))
636                || prs.iter().any(|p| matches!(p, Mention::Url(u) if u == url))
637            {
638                push(c, Signal::Pr, format!("#{n}"));
639            }
640        }
641    }
642    let mut asked = BTreeSet::new();
643    for n in prs.iter().filter_map(|m| match m {
644        Mention::Number(n) => Some(*n),
645        Mention::Url(u) => u.rsplit('/').next().and_then(|d| d.parse().ok()),
646    }) {
647        let token = format!("#{n}");
648        if hits
649            .iter()
650            .any(|h| h.signal == Signal::Pr && h.token == token)
651            || own_prs.contains(&n)
652            || !asked.insert(n)
653        {
654            continue;
655        }
656        if let Some(url) = open_pr(repo, n) {
657            hits.push(Hit {
658                owner: Owner::Pr,
659                id: token.clone(),
660                status: "open".into(),
661                signal: Signal::Pr,
662                token,
663                via: format!("an open pull request with no run record here ({url})"),
664            });
665        }
666    }
667    hits
668}
669
670/// Ask the forge, via `gh`, whether PR `n` is open. Best effort and bounded:
671/// any failure (no `gh`, no remote, offline, a slow answer) is `None`, i.e.
672/// today's behaviour. `GH_REPO` is dropped so the PR is looked up in `repo`'s
673/// own remote, not whatever the environment points at.
674fn gh_open_pr(repo: &Path, n: u64) -> Option<String> {
675    let v = gh_pr_view(repo, n)?;
676    (v["state"] == "OPEN")
677        .then(|| v["url"].as_str().map(str::to_owned))
678        .flatten()
679}
680
681/// [`gh_open_pr`]'s question asked of a pull request a run record names: its
682/// lifecycle, or `None` when the forge cannot say.
683fn gh_pr_state(repo: &Path, n: u64) -> Option<PrLifecycle> {
684    match gh_pr_view(repo, n)?["state"].as_str()? {
685        "OPEN" => Some(PrLifecycle::Open),
686        "MERGED" => Some(PrLifecycle::Merged),
687        "CLOSED" => Some(PrLifecycle::Closed),
688        _ => None,
689    }
690}
691
692fn gh_pr_view(repo: &Path, n: u64) -> Option<serde_json::Value> {
693    let mut child = Command::new("gh")
694        .quiet()
695        .args(["pr", "view", &n.to_string(), "--json", "state,url"])
696        .current_dir(repo)
697        .env_remove("GH_REPO")
698        .env("GH_PROMPT_DISABLED", "1")
699        .stdin(Stdio::null())
700        .stdout(Stdio::piped())
701        .stderr(Stdio::null())
702        .spawn()
703        .ok()?;
704    let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
705    loop {
706        match child.try_wait().ok()? {
707            Some(status) if status.success() => break,
708            Some(_) => return None,
709            None if std::time::Instant::now() >= deadline => {
710                let _ = child.kill();
711                let _ = child.wait();
712                return None;
713            }
714            None => std::thread::sleep(std::time::Duration::from_millis(50)),
715        }
716    }
717    let mut raw = String::new();
718    std::io::Read::read_to_string(&mut child.stdout.take()?, &mut raw).ok()?;
719    serde_json::from_str(&raw).ok()
720}
721
722fn task_claims(
723    t: &Task,
724    runs_root: &Path,
725    seen: &mut BTreeSet<String>,
726    repo: &Path,
727    stale: &mut Staleness<'_>,
728) -> Vec<Claim> {
729    let status = t.status.as_str().to_owned();
730    let mut out = Vec::new();
731    if let Some(b) = &t.review_branch {
732        out.push(Claim {
733            owner: Owner::Task,
734            id: t.id.clone(),
735            status: status.clone(),
736            via: "its review branch".into(),
737            branch: Some(b.clone()),
738            base: None,
739            pr: None,
740        });
741    }
742    for rid in &t.runs {
743        if let Some(view) = RunView::read(runs_root, rid) {
744            seen.insert(rid.clone());
745            // A merged or closed pull request is not in flight, even when the
746            // task that owns it is.
747            let live_pr = !(view.terminal() && stale.pr_released(repo, &view));
748            out.extend(run_claims(
749                &view,
750                Owner::Task,
751                Some((&t.id, &status)),
752                &format!("produced by its run {}", crate::queue::short(rid)),
753                live_pr,
754            ));
755        }
756    }
757    out
758}
759
760/// The claims of one run, attributed to `owner` (the task, when there is one).
761///
762/// `live_pr` is false when the record's open pull request is known to be
763/// settled or somebody else's (see [`Staleness`]); it then names no PR.
764fn run_claims(
765    view: &RunView,
766    owner: Owner,
767    task: Option<(&str, &str)>,
768    via: &str,
769    live_pr: bool,
770) -> Vec<Claim> {
771    let (id, status) = match task {
772        Some((id, status)) => (id.to_owned(), status.to_owned()),
773        None => (view.id.clone(), view.status.clone()),
774    };
775    let pr = view
776        .pr
777        .as_ref()
778        .filter(|p| live_pr && p.state == "open" && p.number > 0)
779        .map(|p| (p.number, p.url.clone()));
780    let via_pr = |extra: &str| match &pr {
781        Some((n, _)) => format!("{via} (PR #{n} open){extra}"),
782        None => format!("{via}{extra}"),
783    };
784    let mut out: Vec<Claim> = view
785        .candidates
786        .iter()
787        .filter(|c| !c.branch.is_empty())
788        .map(|c| Claim {
789            owner: owner.clone(),
790            id: id.clone(),
791            status: status.clone(),
792            via: via_pr(""),
793            branch: Some(c.branch.clone()),
794            base: (!view.base_commit.is_empty()).then(|| view.base_commit.clone()),
795            pr: None,
796        })
797        .collect();
798    if pr.is_some() {
799        out.push(Claim {
800            owner,
801            id,
802            status,
803            via: via_pr(""),
804            branch: None,
805            base: None,
806            pr,
807        });
808    }
809    out
810}
811
812/// Repository identity: the canonical git common dir, so two worktrees of one
813/// repository are the same repository and a path spelled two ways is too.
814#[derive(Default)]
815struct Idents(HashMap<PathBuf, PathBuf>);
816
817impl Idents {
818    fn of(&mut self, path: &Path) -> PathBuf {
819        self.0
820            .entry(path.to_path_buf())
821            .or_insert_with(|| {
822                git(
823                    path,
824                    &["rev-parse", "--path-format=absolute", "--git-common-dir"],
825                )
826                .map(PathBuf::from)
827                .and_then(|p| p.canonicalize().ok())
828                .or_else(|| path.canonicalize().ok())
829                .unwrap_or_else(|| path.to_path_buf())
830            })
831            .clone()
832    }
833}
834
835fn git(cwd: &Path, args: &[&str]) -> Option<String> {
836    let out = Command::new("git")
837        .quiet()
838        .args(args)
839        .current_dir(cwd)
840        .stdin(Stdio::null())
841        .stderr(Stdio::null())
842        .env("GIT_TERMINAL_PROMPT", "0")
843        .output()
844        .ok()?;
845    out.status
846        .success()
847        .then(|| String::from_utf8_lossy(&out.stdout).trim().to_owned())
848}
849
850fn git_ok(cwd: &Path, args: &[&str]) -> bool {
851    git(cwd, args).is_some()
852}
853
854fn is_ref_char(c: char) -> bool {
855    c.is_alphanumeric() || matches!(c, '_' | '-')
856}
857
858/// Does `text` contain `branch` as a whole word? `.` and `/` before it are
859/// fine (`origin/magi/x/A`), a name character on either side is not, so
860/// `magi/27b2/A` does not match inside `magi/27b2/AB`.
861fn names_branch(text: &str, branch: &str) -> bool {
862    if branch.len() < 3 {
863        return false;
864    }
865    text.match_indices(branch).any(|(i, _)| {
866        let before = text[..i].chars().next_back();
867        let after = text[i + branch.len()..].chars().next();
868        let after_ok = match after {
869            None => true,
870            Some('/') => false,
871            Some('.') => !text[i + branch.len() + 1..]
872                .chars()
873                .next()
874                .is_some_and(is_ref_char),
875            Some(c) => !is_ref_char(c),
876        };
877        before.is_none_or(|c| !is_ref_char(c) && c != '.') && after_ok
878    })
879}
880
881#[derive(Debug, PartialEq, Eq)]
882enum Mention {
883    Number(u64),
884    Url(String),
885}
886
887/// `#48`, `PR 48`, `PR #48`, `pull request 48`, and `.../pull/48` URLs.
888fn pr_numbers(text: &str) -> Vec<Mention> {
889    let mut out = Vec::new();
890    let bytes = text.as_bytes();
891    let digits = |from: usize| -> Option<(u64, usize)> {
892        let n = text[from..].bytes().take_while(u8::is_ascii_digit).count();
893        (n > 0 && n < 10)
894            .then(|| text[from..from + n].parse().ok().map(|v| (v, from + n)))
895            .flatten()
896    };
897    for (i, _) in text.match_indices('#') {
898        if let Some((n, _)) = digits(i + 1) {
899            let word_before = i > 0 && is_ref_char(bytes[i - 1] as char);
900            if !word_before {
901                out.push(Mention::Number(n));
902            }
903        }
904    }
905    let lower = text.to_ascii_lowercase();
906    for key in ["pull request ", "pr "] {
907        for (i, _) in lower.match_indices(key) {
908            if i > 0 && is_ref_char(bytes[i - 1] as char) {
909                continue;
910            }
911            let from = i + key.len();
912            let from = if text[from..].starts_with('#') {
913                from + 1
914            } else {
915                from
916            };
917            if let Some((n, _)) = digits(from) {
918                out.push(Mention::Number(n));
919            }
920        }
921    }
922    for (i, _) in text.match_indices("/pull/") {
923        if let Some((_, end)) = digits(i + 6) {
924            let start = text[..i]
925                .rfind(|c: char| c.is_whitespace() || matches!(c, '(' | '<' | '"' | '\''))
926                .map_or(0, |p| p + 1);
927            out.push(Mention::Url(text[start..end].to_owned()));
928        }
929    }
930    out
931}
932
933/// Hex words of 7..=40 chars in `text` that resolve to a commit in `repo`.
934fn sha_candidates(repo: &Path, text: &str) -> Vec<String> {
935    let mut seen = BTreeSet::new();
936    let mut out = Vec::new();
937    for word in text.split(|c: char| !c.is_ascii_alphanumeric()) {
938        if !(7..=40).contains(&word.len()) || !word.bytes().all(|b| b.is_ascii_hexdigit()) {
939            continue;
940        }
941        if seen.len() >= 16 || !seen.insert(word.to_ascii_lowercase()) {
942            continue;
943        }
944        if let Some(full) = git(
945            repo,
946            &[
947                "rev-parse",
948                "--verify",
949                "--quiet",
950                &format!("{word}^{{commit}}"),
951            ],
952        ) {
953            out.push(full);
954        }
955    }
956    out
957}
958
959/// Is `sha` reachable from `branch` but not from `base`? A SHA already in the
960/// base would match every branch there is.
961fn on_branch_only(repo: &Path, sha: &str, branch: &str, base: &str) -> bool {
962    let tip = format!("{branch}^{{commit}}");
963    git_ok(repo, &["rev-parse", "--verify", "--quiet", &tip])
964        && git_ok(repo, &["merge-base", "--is-ancestor", sha, branch])
965        && !git_ok(repo, &["merge-base", "--is-ancestor", sha, base])
966}
967
968fn short_sha(sha: &str) -> String {
969    sha.chars().take(7).collect()
970}
971
972#[cfg(test)]
973mod tests {
974    use super::*;
975    use crate::queue::{Source, Task};
976
977    struct Fx {
978        _tmp: tempfile::TempDir,
979        repo: PathBuf,
980        runs: PathBuf,
981        q: Queue,
982    }
983
984    fn sh(cwd: &Path, args: &[&str]) -> String {
985        let out = Command::new("git")
986            .quiet()
987            .args(["-c", "user.name=t", "-c", "user.email=t@t"])
988            .args(args)
989            .current_dir(cwd)
990            .output()
991            .unwrap();
992        assert!(out.status.success(), "git {args:?}: {out:?}");
993        String::from_utf8_lossy(&out.stdout).trim().to_owned()
994    }
995
996    /// A repo with `main` and a branch `magi/aaaa/A` one commit ahead.
997    /// Returns the fixture plus the (main, branch) tip SHAs.
998    fn fx() -> (Fx, String, String) {
999        let tmp = tempfile::tempdir().unwrap();
1000        let repo = tmp.path().join("repo");
1001        std::fs::create_dir_all(&repo).unwrap();
1002        sh(&repo, &["init", "-q", "-b", "main"]);
1003        std::fs::write(repo.join("a"), "1").unwrap();
1004        sh(&repo, &["add", "."]);
1005        sh(&repo, &["commit", "-q", "-m", "base"]);
1006        let base = sh(&repo, &["rev-parse", "HEAD"]);
1007        sh(&repo, &["checkout", "-q", "-b", "magi/aaaa/A"]);
1008        std::fs::write(repo.join("a"), "2").unwrap();
1009        sh(&repo, &["commit", "-q", "-am", "work"]);
1010        let tip = sh(&repo, &["rev-parse", "HEAD"]);
1011        sh(&repo, &["checkout", "-q", "main"]);
1012        let runs = tmp.path().join("runs");
1013        std::fs::create_dir_all(&runs).unwrap();
1014        let q = Queue::at(tmp.path().join("queue"));
1015        (
1016            Fx {
1017                _tmp: tmp,
1018                repo,
1019                runs,
1020                q,
1021            },
1022            base,
1023            tip,
1024        )
1025    }
1026
1027    fn write_run(f: &Fx, id: &str, status: &str, base: &str, pr: Option<(u64, &str)>) {
1028        let dir = f.runs.join(id);
1029        std::fs::create_dir_all(&dir).unwrap();
1030        let pr = pr.map(|(n, s)| {
1031            serde_json::json!({"url": format!("https://github.com/o/r/pull/{n}"), "number": n, "state": s})
1032        });
1033        let v = serde_json::json!({
1034            "schema": 999, "id": id, "repo": f.repo, "status": status,
1035            "base_commit": base, "candidates": [{"branch": "magi/aaaa/A"}], "pr": pr,
1036        });
1037        std::fs::write(dir.join("run.json"), v.to_string()).unwrap();
1038    }
1039
1040    fn file_task(f: &Fx, status: TaskStatus, runs: &[&str]) -> Task {
1041        let mut t = Task::new("t".into(), "x".into(), f.repo.clone(), Source::Human);
1042        t.status = status;
1043        t.runs = runs.iter().map(|s| (*s).to_owned()).collect();
1044        f.q.put(&mut t).unwrap();
1045        t
1046    }
1047
1048    const RID: &str = "20260901-100000-aaaa";
1049
1050    fn run(f: &Fx, text: &str, review: Option<&str>) -> Vec<Hit> {
1051        check_with(
1052            &f.q,
1053            &f.runs,
1054            &f.repo,
1055            text,
1056            review,
1057            None,
1058            &|_, _| None,
1059            &|_, _| None,
1060        )
1061    }
1062
1063    #[test]
1064    fn branch_matches_a_live_run_and_its_task() {
1065        let (f, base, _) = fx();
1066        write_run(&f, RID, "reviewing", &base, None);
1067        let t = file_task(&f, TaskStatus::Running, &[RID]);
1068        let hits = run(&f, "land magi/aaaa/A onto a fresh branch.", None);
1069        assert!(hits.iter().any(|h| h.owner == Owner::Task
1070            && h.id == t.id
1071            && h.signal == Signal::Branch
1072            && h.token == "magi/aaaa/A"));
1073        assert!(hits.iter().any(|h| h.owner == Owner::Run && h.id == RID));
1074        let msg = Duplicate::new(hits).to_string();
1075        assert!(
1076            msg.contains("--force") && msg.contains("magi/aaaa/A"),
1077            "{msg}"
1078        );
1079    }
1080
1081    #[test]
1082    fn branch_must_match_whole_word() {
1083        let (f, base, _) = fx();
1084        write_run(&f, RID, "reviewing", &base, None);
1085        assert!(run(&f, "see magi/aaaa/AB and magi/aaaa/A/x", None).is_empty());
1086    }
1087
1088    #[test]
1089    fn sha_on_the_branch_matches_but_one_in_base_does_not() {
1090        let (f, base, tip) = fx();
1091        write_run(&f, RID, "reviewing", &base, None);
1092        let hits = run(&f, &format!("land commit {} please", &tip[..8]), None);
1093        assert!(
1094            hits.iter().any(|h| h.signal == Signal::Sha && h.id == RID),
1095            "{hits:?}"
1096        );
1097        assert!(run(&f, &format!("see {}", &base[..9]), None).is_empty());
1098        // Hex that resolves to nothing is ignored.
1099        assert!(run(&f, "deadbeef and 1234567", None).is_empty());
1100    }
1101
1102    #[test]
1103    fn pr_number_matches_in_every_spelling() {
1104        let (f, base, _) = fx();
1105        write_run(&f, RID, "ready", &base, Some((48, "open")));
1106        for text in [
1107            "finish #48",
1108            "PR 48 is stale",
1109            "pr #48",
1110            "pull request 48",
1111            "https://github.com/o/r/pull/48",
1112        ] {
1113            let hits = run(&f, text, None);
1114            assert!(
1115                hits.iter().any(|h| h.signal == Signal::Pr),
1116                "{text}: {hits:?}"
1117            );
1118        }
1119        assert!(run(&f, "see #480 and PR 4 and issue48", None).is_empty());
1120    }
1121
1122    #[test]
1123    fn terminal_runs_and_done_tasks_do_not_match() {
1124        let (f, base, tip) = fx();
1125        write_run(&f, RID, "merged", &base, Some((48, "merged")));
1126        file_task(&f, TaskStatus::Done, &[RID]);
1127        let text = format!("magi/aaaa/A {} #48", &tip[..8]);
1128        assert!(run(&f, &text, None).is_empty());
1129    }
1130
1131    #[test]
1132    fn terminal_run_with_open_pr_or_open_task_still_claims() {
1133        // Changed deliberately: a terminal run's recorded-open PR now claims
1134        // only while the forge says open or cannot be read (`run` passes an
1135        // unreadable forge). A merged / closed answer or a released run no
1136        // longer claims; see the `stale_open_*` tests below.
1137        let (f, base, _) = fx();
1138        write_run(&f, RID, "ready", &base, Some((48, "open")));
1139        assert!(!run(&f, "magi/aaaa/A", None).is_empty());
1140        let (g, base, _) = fx();
1141        write_run(&g, RID, "ready", &base, None);
1142        assert!(run(&g, "magi/aaaa/A", None).is_empty());
1143        let t = file_task(&g, TaskStatus::Held, &[RID]);
1144        let hits = run(&g, "magi/aaaa/A", None);
1145        assert!(hits.iter().any(|h| h.id == t.id), "{hits:?}");
1146    }
1147
1148    #[test]
1149    fn review_only_matches_a_branch_a_live_task_owns() {
1150        let (f, base, _) = fx();
1151        write_run(&f, RID, "ready", &base, None);
1152        assert!(run(&f, "", Some("magi/aaaa/A")).is_empty());
1153        let mut t = file_task(&f, TaskStatus::Queued, &[]);
1154        t.review_branch = Some("magi/aaaa/A".into());
1155        f.q.put(&mut t).unwrap();
1156        let hits = run(&f, "", Some("magi/aaaa/A"));
1157        assert!(
1158            hits.iter()
1159                .any(|h| h.id == t.id && h.signal == Signal::Branch)
1160        );
1161    }
1162
1163    #[test]
1164    fn other_repository_and_edited_task_do_not_match() {
1165        let (f, base, _) = fx();
1166        write_run(&f, RID, "reviewing", &base, None);
1167        let t = file_task(&f, TaskStatus::Running, &[RID]);
1168        let other = f._tmp.path().join("other");
1169        std::fs::create_dir_all(&other).unwrap();
1170        sh(&other, &["init", "-q"]);
1171        assert!(
1172            check_with(
1173                &f.q,
1174                &f.runs,
1175                &other,
1176                "magi/aaaa/A",
1177                None,
1178                None,
1179                &|_, _| None,
1180                &|_, _| None,
1181            )
1182            .is_empty()
1183        );
1184        // Editing the task that owns the run never collides with itself.
1185        assert!(
1186            check_with(
1187                &f.q,
1188                &f.runs,
1189                &f.repo,
1190                "magi/aaaa/A",
1191                None,
1192                Some(&t.id),
1193                &|_, _| None,
1194                &|_, _| None,
1195            )
1196            .is_empty()
1197        );
1198    }
1199
1200    #[test]
1201    fn a_worktree_is_the_same_repository() {
1202        let (f, base, _) = fx();
1203        write_run(&f, RID, "reviewing", &base, None);
1204        let wt = f._tmp.path().join("wt");
1205        sh(
1206            &f.repo,
1207            &["worktree", "add", "-q", wt.to_str().unwrap(), "-b", "other"],
1208        );
1209        assert!(
1210            !check_with(
1211                &f.q,
1212                &f.runs,
1213                &wt,
1214                "magi/aaaa/A",
1215                None,
1216                None,
1217                &|_, _| None,
1218                &|_, _| None,
1219            )
1220            .is_empty()
1221        );
1222    }
1223
1224    #[test]
1225    fn an_open_pr_without_a_run_record_matches_through_the_forge() {
1226        let (f, _, _) = fx();
1227        let open = |_: &Path, n: u64| (n == 48).then(|| "https://example.test/pull/48".to_owned());
1228        let hit = |text: &str| {
1229            check_with(&f.q, &f.runs, &f.repo, text, None, None, &open, &|_, _| {
1230                None
1231            })
1232        };
1233        let hits = hit("finish PR #48");
1234        assert_eq!(hits.len(), 1, "{hits:?}");
1235        assert_eq!(hits[0].owner, Owner::Pr);
1236        assert!(hits[0].to_string().contains("#48"));
1237        // Closed / unknown / unnamed PRs never match.
1238        assert!(hit("finish PR #49").is_empty());
1239        assert!(hit("finish the work").is_empty());
1240    }
1241
1242    #[test]
1243    fn a_forge_hit_does_not_repeat_a_pr_a_run_already_explains() {
1244        let (f, base, _) = fx();
1245        write_run(&f, RID, "ready", &base, Some((48, "open")));
1246        let open = |_: &Path, _: u64| Some("u".to_owned());
1247        let hits = check_with(&f.q, &f.runs, &f.repo, "#48", None, None, &open, &|_, _| {
1248            None
1249        });
1250        assert!(hits.iter().all(|h| h.owner != Owner::Pr), "{hits:?}");
1251        assert!(!hits.is_empty());
1252    }
1253
1254    #[test]
1255    fn an_edited_tasks_own_open_pr_is_not_a_forge_hit() {
1256        let (f, base, _) = fx();
1257        write_run(&f, RID, "ready", &base, Some((48, "open")));
1258        let t = file_task(&f, TaskStatus::Running, &[RID]);
1259        let open = |_: &Path, _: u64| Some("u".to_owned());
1260        let hits = check_with(
1261            &f.q,
1262            &f.runs,
1263            &f.repo,
1264            "#48",
1265            None,
1266            Some(&t.id),
1267            &open,
1268            &|_, _| None,
1269        );
1270        assert!(hits.is_empty(), "{hits:?}");
1271    }
1272
1273    fn with_forge(f: &Fx, text: &str, state: Option<PrLifecycle>) -> Vec<Hit> {
1274        check_with(
1275            &f.q,
1276            &f.runs,
1277            &f.repo,
1278            text,
1279            None,
1280            None,
1281            &|_, _| None,
1282            &move |_, _| state,
1283        )
1284    }
1285
1286    fn release(f: &Fx, id: &str, to: &str) {
1287        let path = f.runs.join(id).join("run.json");
1288        let mut v: serde_json::Value =
1289            serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
1290        v["released_to"] = serde_json::json!(to);
1291        std::fs::write(path, v.to_string()).unwrap();
1292    }
1293
1294    #[test]
1295    fn stale_open_pr_that_the_forge_says_is_merged_or_closed_does_not_claim() {
1296        for state in [PrLifecycle::Merged, PrLifecycle::Closed] {
1297            let (f, base, _) = fx();
1298            write_run(&f, RID, "superseded", &base, Some((48, "open")));
1299            assert!(with_forge(&f, "follow up on #48", Some(state)).is_empty());
1300            assert!(with_forge(&f, "magi/aaaa/A", Some(state)).is_empty());
1301            // Same through an owning, still-live task: the PR is not its claim.
1302            file_task(&f, TaskStatus::Held, &[RID]);
1303            let hits = with_forge(&f, "follow up on #48", Some(state));
1304            assert!(hits.iter().all(|h| h.signal != Signal::Pr), "{hits:?}");
1305        }
1306    }
1307
1308    #[test]
1309    fn stale_open_pr_with_an_unreadable_forge_still_claims() {
1310        let (f, base, _) = fx();
1311        write_run(&f, RID, "superseded", &base, Some((48, "open")));
1312        let hits = with_forge(&f, "follow up on #48", None);
1313        assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1314    }
1315
1316    #[test]
1317    fn a_genuinely_open_pr_still_claims() {
1318        let (f, base, _) = fx();
1319        write_run(&f, RID, "blocked", &base, Some((48, "open")));
1320        let hits = with_forge(&f, "follow up on #48", Some(PrLifecycle::Open));
1321        assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1322    }
1323
1324    #[test]
1325    fn a_released_run_defers_to_its_successor_without_asking_the_forge() {
1326        let (f, base, _) = fx();
1327        let next = "20260901-110000-bbbb";
1328        write_run(&f, RID, "superseded", &base, Some((48, "open")));
1329        write_run(&f, next, "merged", &base, Some((48, "merged")));
1330        release(&f, RID, next);
1331        let asked = std::cell::Cell::new(0);
1332        let hits = check_with(
1333            &f.q,
1334            &f.runs,
1335            &f.repo,
1336            "follow up on #48",
1337            None,
1338            None,
1339            &|_, _| None,
1340            &|_, _| {
1341                asked.set(asked.get() + 1);
1342                None
1343            },
1344        );
1345        assert!(hits.is_empty(), "{hits:?}");
1346        assert_eq!(asked.get(), 0);
1347        // A successor that cannot be read decides nothing: claim stands.
1348        std::fs::remove_dir_all(f.runs.join(next)).unwrap();
1349        let hits = with_forge(&f, "follow up on #48", None);
1350        assert!(hits.iter().any(|h| h.signal == Signal::Pr), "{hits:?}");
1351    }
1352
1353    #[test]
1354    fn forge_lookups_are_cached_and_stop_after_a_failure() {
1355        let (f, base, _) = fx();
1356        for (i, id) in ["20260901-100000-aaa1", "20260901-100000-aaa2"]
1357            .iter()
1358            .enumerate()
1359        {
1360            write_run(&f, id, "blocked", &base, Some((48 + i as u64, "open")));
1361        }
1362        let asked = std::cell::Cell::new(0);
1363        check_with(
1364            &f.q,
1365            &f.runs,
1366            &f.repo,
1367            "x",
1368            None,
1369            None,
1370            &|_, _| None,
1371            &|_, _| {
1372                asked.set(asked.get() + 1);
1373                None
1374            },
1375        );
1376        assert_eq!(
1377            asked.get(),
1378            1,
1379            "an unreadable forge is asked once, not per PR"
1380        );
1381    }
1382
1383    use std::sync::Arc;
1384    use std::sync::atomic::{AtomicUsize, Ordering};
1385
1386    fn verdict(ruling: Ruling) -> Result<Judgement> {
1387        Ok(Judgement {
1388            ruling,
1389            reason: "because".into(),
1390            agent: "j".into(),
1391        })
1392    }
1393
1394    /// Run `screen` over `text` with a judge answering `answer`; returns the
1395    /// outcome and how many times the judge was asked.
1396    fn judged(
1397        f: &Fx,
1398        text: &str,
1399        answer: impl Fn() -> Result<Judgement> + Send + Sync + 'static,
1400    ) -> (Result<(), Duplicate>, usize) {
1401        let calls = Arc::new(AtomicUsize::new(0));
1402        let seen = calls.clone();
1403        let judge = move |_: String, _: Vec<Hit>| -> JudgeFuture {
1404            seen.fetch_add(1, Ordering::SeqCst);
1405            let r = answer();
1406            Box::pin(async move { r })
1407        };
1408        let hits = run(f, text, None);
1409        let out = tokio::runtime::Builder::new_current_thread()
1410            .build()
1411            .unwrap()
1412            .block_on(screen(hits, text, None, &judge));
1413        (out, calls.load(Ordering::SeqCst))
1414    }
1415
1416    const NAMES: &str = "seen on magi/aaaa/A, which does not touch this file";
1417
1418    fn fx_with_run() -> Fx {
1419        let (f, base, _) = fx();
1420        write_run(&f, RID, "reviewing", &base, None);
1421        f
1422    }
1423
1424    #[test]
1425    fn a_mentions_ruling_lets_the_work_through() {
1426        let f = fx_with_run();
1427        let (out, calls) = judged(&f, NAMES, || verdict(Ruling::Mentions));
1428        assert!(out.is_ok());
1429        assert_eq!(calls, 1);
1430    }
1431
1432    #[test]
1433    fn an_owns_ruling_refuses_and_says_why() {
1434        let f = fx_with_run();
1435        let (out, _) = judged(&f, NAMES, || verdict(Ruling::Owns));
1436        let msg = out.unwrap_err().to_string();
1437        assert!(
1438            msg.contains("owns - because") && msg.contains("--force"),
1439            "{msg}"
1440        );
1441        assert!(msg.contains("magi/aaaa/A"), "{msg}");
1442    }
1443
1444    #[test]
1445    fn an_unsure_ruling_refuses() {
1446        let f = fx_with_run();
1447        let (out, _) = judged(&f, NAMES, || verdict(Ruling::Unsure));
1448        assert!(out.unwrap_err().to_string().contains("unsure - because"));
1449    }
1450
1451    #[test]
1452    fn a_failing_judge_refuses() {
1453        let f = fx_with_run();
1454        let (out, calls) = judged(&f, NAMES, || Err(anyhow::anyhow!("quota")));
1455        let msg = out.unwrap_err().to_string();
1456        assert!(msg.contains("judge could not decide: quota"), "{msg}");
1457        assert_eq!(calls, 1);
1458    }
1459
1460    #[test]
1461    fn garbage_and_unknown_rulings_do_not_parse() {
1462        assert!(parse_ruling("sure, go ahead").is_err());
1463        assert!(parse_ruling(r#"{"ruling":"maybe","reason":"x"}"#).is_err());
1464        assert!(parse_ruling(r#"{"ruling":"mentions"}"#).is_err());
1465        assert!(parse_ruling(r#"{"ruling":"mentions","reason":"  "}"#).is_err());
1466        let two = "{\"ruling\":\"mentions\",\"reason\":\"c\"}\n{\"ruling\":\"owns\"}";
1467        assert!(parse_ruling(two).is_err());
1468        assert!(parse_ruling("ok {\"ruling\":\"mentions\",\"reason\":\"c\"}").is_err());
1469        let (r, why) = parse_ruling("{\"ruling\":\"Mentions\",\"reason\":\"cites\\nit\"}").unwrap();
1470        assert_eq!((r, why.as_str()), (Ruling::Mentions, "cites it"));
1471    }
1472
1473    #[test]
1474    fn a_text_too_long_to_judge_in_full_refuses_without_asking() {
1475        let f = fx_with_run();
1476        let long = format!("{NAMES} {}", "x".repeat(prompt::DUPES_JUDGE_MAX_CHARS));
1477        let (out, calls) = judged(&f, &long, || verdict(Ruling::Mentions));
1478        assert!(out.unwrap_err().to_string().contains("too long to judge"));
1479        assert_eq!(calls, 0);
1480    }
1481
1482    #[test]
1483    fn no_hit_never_asks_the_judge() {
1484        let (f, _, _) = fx();
1485        let (out, calls) = judged(&f, "nothing named here", || verdict(Ruling::Owns));
1486        assert!(out.is_ok());
1487        assert_eq!(calls, 0);
1488    }
1489
1490    #[test]
1491    fn without_a_config_a_hit_still_refuses() {
1492        let f = fx_with_run();
1493        let hits = run(&f, NAMES, None);
1494        let out = tokio::runtime::Builder::new_current_thread()
1495            .build()
1496            .unwrap()
1497            .block_on(screen_with_config(hits, NAMES, None, &f.repo, None));
1498        assert!(out.unwrap_err().judge.is_some());
1499    }
1500}