1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum Signal {
50 Branch,
52 Sha,
54 Pr,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq)]
60pub enum Owner {
61 Task,
63 Run,
65 Pr,
67}
68
69#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct Hit {
72 pub owner: Owner,
74 pub id: String,
76 pub status: String,
78 pub signal: Signal,
80 pub token: String,
82 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#[derive(Debug, Clone)]
111pub struct Duplicate {
112 pub hits: Vec<Hit>,
114 pub judge: Option<String>,
117}
118
119impl Duplicate {
120 pub fn new(hits: Vec<Hit>) -> Self {
122 Self { hits, judge: None }
123 }
124
125 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
155pub enum Ruling {
156 Owns,
158 Mentions,
160 Unsure,
162}
163
164#[derive(Debug, Clone, PartialEq, Eq)]
166pub struct Judgement {
167 pub ruling: Ruling,
169 pub reason: String,
171 pub agent: String,
173}
174
175pub type JudgeFuture = Pin<Box<dyn Future<Output = Result<Judgement>> + Send>>;
178
179const JUDGE_BUDGET: Duration = Duration::from_secs(120);
181const JUDGE_TURN: Duration = Duration::from_secs(90);
183
184fn parse_ruling(text: &str) -> Result<(Ruling, String)> {
187 #[derive(Deserialize)]
188 struct Raw {
189 ruling: String,
190 reason: String,
191 }
192 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
215pub 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
285pub 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 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
348pub 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#[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 #[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 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
426struct 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
450const 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 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#[derive(Debug, Clone)]
502struct Claim {
503 owner: Owner,
504 id: String,
505 status: String,
506 via: String,
507 branch: Option<String>,
508 base: Option<String>,
510 pr: Option<(u64, String)>,
512}
513
514pub 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#[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 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 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
670fn 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
681fn 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 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
760fn 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#[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
858fn 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
887fn 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
933fn 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
959fn 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 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 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 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 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 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 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 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 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}