Skip to main content

magi/
release_watch.rs

1//! Watches every open release-bump pull request until it lands.
2//!
3//! [`crate::bump`] opens a `chore/release-vX.Y.Z` pull request with auto-merge
4//! armed and then stops looking. A required check that goes red (a flaky
5//! windows or browser job) leaves auto-merge waiting forever, and nothing
6//! pages anybody because the notices there only fire when the pull request
7//! cannot be *opened*. This module is the part that keeps looking.
8//!
9//! It runs inside `magi serve` as its own task, like [`crate::waiter`], and is
10//! tied to no run: a pending bump is coalesced across runs and outlives the
11//! one that opened it. The forge is the source of truth (every open pull
12//! request whose head branch starts with `chore/release-v`, in every checkout
13//! magi knows); the pending marker is not consulted, so a missing one changes
14//! nothing.
15//!
16//! - **The decision is pure.** [`decide`] maps a checks snapshot and the
17//!   recorded [`WatchState`] to one [`Step`]. A forge answer that cannot be
18//!   read is `Unknown`: nothing moves, not even the stall clock.
19//! - **A failed run is rerun once, durably.** The run id is written to the
20//!   state file *before* `gh run rerun <id> --failed` is called, so a restart
21//!   never reruns twice; if the call then fails, the cost is one escalation to
22//!   a human, never a rerun loop.
23//! - **Escalation is a notice plus a question** ([`crate::bump::NOTICE_NODE`],
24//!   no `cwd`, no run). Nothing is ever merged, closed or pushed: the watcher
25//!   only reruns jobs. Silence is a hold. The owner's choice is applied by the
26//!   watcher itself on its next lap, recorded in `WatchState::applied` so an
27//!   answer is applied once.
28//! - **Notice wording is fixed per stage** and carries only the pull request,
29//!   so a repeated poll never relights it; everything that varies (job names,
30//!   links) rides in the question.
31
32use std::collections::BTreeMap;
33use std::future::Future;
34use std::path::{Path, PathBuf};
35use std::pin::Pin;
36use std::time::Duration;
37
38use anyhow::{Result, bail};
39use serde::{Deserialize, Serialize};
40
41use crate::ask::{Answer, Question, Questions};
42use crate::land::{self, PrLifecycle, RollupView, Verdict};
43use crate::notices::{self, Notice, Notices};
44
45/// Pause between laps.
46const LAP: Duration = Duration::from_secs(300);
47
48/// How long after a rerun the forge may still report the old failure before
49/// that is read as "still red". `gh run rerun` returns before the checks flip
50/// to pending.
51const RERUN_GRACE: i64 = 180;
52
53/// Choice: rerun the failed jobs once more. The three choices are matched
54/// verbatim when applied.
55pub const RERUN_AGAIN: &str = "rerun again";
56/// Choice: stay quiet until a check or the head changes.
57pub const HOLD: &str = "hold";
58/// Choice: stop watching this pull request.
59pub const LEAVE_IT: &str = "leave it";
60
61/// What one look at one pull request concludes.
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub(crate) enum Step {
64    /// Merged or closed: stop watching.
65    Done,
66    /// Nothing to do yet.
67    Wait,
68    /// The forge answer could not be read; change nothing.
69    Unknown,
70    /// Rerun the failed jobs of these workflow runs, once each.
71    Rerun(Vec<String>),
72    /// A human is needed.
73    Escalate(Why),
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77pub(crate) enum Why {
78    /// Still red after the one rerun, or red with nothing to rerun.
79    StillRed,
80    /// No progress for longer than `[daemon] release_stall_minutes`.
81    Stalled,
82}
83
84/// Everything remembered about one watched pull request.
85#[derive(Debug, Clone, Default, Serialize, Deserialize)]
86#[serde(default)]
87pub(crate) struct WatchState {
88    /// Checkout the pull request belongs to.
89    pub repo: String,
90    pub url: String,
91    /// Head commit the rest of this record is about.
92    pub head: String,
93    /// Workflow runs already rerun for this head, with when (unix seconds).
94    pub reruns: BTreeMap<String, i64>,
95    /// Digest of the head and every check's verdict, and when it last changed.
96    pub fingerprint: String,
97    pub progress_at: i64,
98    /// The open question asked about this pull request.
99    pub question: Option<String>,
100    /// Fingerprint the owner chose to hold: no new question until it changes.
101    pub held: Option<String>,
102    /// The owner said to leave it.
103    pub ignored: bool,
104    /// Question ids whose answer was already applied.
105    pub applied: Vec<String>,
106}
107
108fn is_failed(v: Verdict) -> bool {
109    v == Verdict::Fail
110}
111
112/// Digest of what "progress" means: the head and every check's verdict.
113pub(crate) fn fingerprint(snap: &RollupView) -> String {
114    let mut parts: Vec<String> = snap
115        .checks
116        .iter()
117        .map(|c| format!("{}={:?}", c.name, c.verdict))
118        .collect();
119    parts.sort();
120    format!("{}|{}", snap.head, parts.join(","))
121}
122
123/// Fold a fresh snapshot into the record: a new head forgets the old head's
124/// reruns, and a changed fingerprint restarts the stall clock.
125pub(crate) fn observe(st: &mut WatchState, snap: &RollupView, now: i64) {
126    if st.head != snap.head {
127        st.reruns.clear();
128        st.head = snap.head.clone();
129    }
130    let fp = fingerprint(snap);
131    if st.fingerprint != fp {
132        st.fingerprint = fp;
133        st.progress_at = now;
134    }
135}
136
137/// Decide what to do about a pull request. `snap` is `None` when the forge
138/// could not be read. `stall_secs == 0` disables the stall rule.
139pub(crate) fn decide(
140    snap: Option<&RollupView>,
141    st: &WatchState,
142    now: i64,
143    stall_secs: i64,
144) -> Step {
145    let Some(snap) = snap else {
146        return Step::Unknown;
147    };
148    if snap.state != PrLifecycle::Open {
149        return Step::Done;
150    }
151    if st.ignored || st.held.as_deref() == Some(fingerprint(snap).as_str()) {
152        return Step::Wait;
153    }
154
155    let any_pending = snap.checks.iter().any(|c| c.verdict == Verdict::Pending);
156    // A run is complete once none of its own checks is still pending; another
157    // run's pending job says nothing about it.
158    let complete = |run: &str| {
159        !snap
160            .checks
161            .iter()
162            .any(|c| c.verdict == Verdict::Pending && c.run.as_deref() == Some(run))
163    };
164    let mut fresh = Vec::new();
165    let mut spent_recently = false;
166    let mut spent = false;
167    let mut external = false;
168    for c in snap.checks.iter().filter(|c| is_failed(c.verdict)) {
169        match c.run.as_deref() {
170            None => external = true,
171            Some(run) if !complete(run) => {}
172            Some(run) => match st.reruns.get(run) {
173                None => {
174                    if !fresh.iter().any(|r: &String| r == run) {
175                        fresh.push(run.to_owned());
176                    }
177                }
178                Some(&at) if now - at < RERUN_GRACE => spent_recently = true,
179                Some(_) => spent = true,
180            },
181        }
182    }
183    if !fresh.is_empty() {
184        return Step::Rerun(fresh);
185    }
186    if spent_recently {
187        return Step::Wait;
188    }
189    if spent || (external && !any_pending) {
190        return Step::Escalate(Why::StillRed);
191    }
192    if stall_secs > 0 && now - st.progress_at >= stall_secs {
193        return Step::Escalate(Why::Stalled);
194    }
195    Step::Wait
196}
197
198/// `owner/repo#N` out of a pull request url, the stable name for notices.
199pub(crate) fn pr_key(url: &str) -> Option<String> {
200    let rest = url.split("://").nth(1)?;
201    let mut it = rest.split('/');
202    let _host = it.next()?;
203    let owner = it.next()?;
204    let repo = it.next()?;
205    if it.next()? != "pull" {
206        return None;
207    }
208    let n: String = it
209        .next()?
210        .chars()
211        .take_while(char::is_ascii_digit)
212        .collect();
213    (!n.is_empty()).then(|| format!("{owner}/{repo}#{n}"))
214}
215
216fn notice_key(pr: &str) -> String {
217    format!("release-pr:{pr}")
218}
219
220/// Fixed per stage: no counts, times or job names.
221fn rerun_message(pr: &str) -> String {
222    format!("Release PR {pr}: failed checks were rerun once")
223}
224
225fn human_message(pr: &str) -> String {
226    format!("Release PR {pr} is stuck and needs a human")
227}
228
229fn question_detail(snap: &RollupView, why: Why, st: &WatchState) -> String {
230    let mut s = String::new();
231    s.push_str(&format!("Pull request: {}\n\n", snap.url));
232    s.push_str(match why {
233        Why::StillRed => "Checks are still red after the failed jobs were rerun once.\n\n",
234        Why::Stalled => "The pull request has made no progress for too long.\n\n",
235    });
236    let failed: Vec<_> = snap
237        .checks
238        .iter()
239        .filter(|c| is_failed(c.verdict))
240        .collect();
241    if failed.is_empty() {
242        s.push_str("No check is failing.\n");
243    } else {
244        s.push_str("Failed jobs:\n");
245        for c in failed {
246            match &c.url {
247                Some(u) => s.push_str(&format!("- {} ({u})\n", c.name)),
248                None => s.push_str(&format!("- {}\n", c.name)),
249            }
250        }
251    }
252    s.push_str(&format!(
253        "\nWorkflow runs rerun so far: {}.\n\n\
254         - `{RERUN_AGAIN}`: rerun the failed jobs once more.\n\
255         - `{HOLD}`: stay quiet until a check or the head changes.\n\
256         - `{LEAVE_IT}`: stop watching this pull request.\n\n\
257         Nothing is merged or pushed by magi; silence is a hold.\n",
258        st.reruns.len()
259    ));
260    s
261}
262
263type Fut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
264
265/// Every forge access, so tests inject a fake instead of `gh`.
266pub(crate) trait ReleaseForge: Send + Sync {
267    /// `(branch, url)` of every open release-bump pull request in `repo`.
268    fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>>;
269    fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>>;
270    fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>>;
271}
272
273/// The real forge: `gh`.
274pub(crate) struct GhForge;
275
276impl ReleaseForge for GhForge {
277    fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
278        Box::pin(crate::bump::list_open_release_prs(repo))
279    }
280
281    fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>> {
282        Box::pin(async move {
283            let args = [
284                "pr".to_owned(),
285                "view".to_owned(),
286                url.to_owned(),
287                "--json".to_owned(),
288                "url,number,state,headRefOid,statusCheckRollup".to_owned(),
289            ];
290            let (ok, out) = land::gh(repo, &args).await?;
291            if !ok {
292                bail!("gh pr view {url}: {out}");
293            }
294            land::parse_rollup(&out)
295        })
296    }
297
298    fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
299        Box::pin(async move {
300            let args = [
301                "run".to_owned(),
302                "rerun".to_owned(),
303                run.to_owned(),
304                "--failed".to_owned(),
305            ];
306            let (ok, out) = land::gh(repo, &args).await?;
307            if !ok {
308                bail!("gh run rerun {run}: {out}");
309            }
310            Ok(())
311        })
312    }
313}
314
315/// The watcher: a forge, a magi home, and nothing else.
316pub(crate) struct Watcher {
317    forge: Box<dyn ReleaseForge>,
318    home: PathBuf,
319}
320
321impl Watcher {
322    pub(crate) fn new(forge: Box<dyn ReleaseForge>, home: PathBuf) -> Self {
323        Self { forge, home }
324    }
325
326    fn dir(&self) -> PathBuf {
327        self.home.join("release-watch")
328    }
329
330    fn state_path(&self, pr: &str) -> PathBuf {
331        self.dir().join(format!("{}.json", notices::id_of(pr)))
332    }
333
334    fn load(&self, pr: &str) -> WatchState {
335        std::fs::read_to_string(self.state_path(pr))
336            .ok()
337            .and_then(|s| serde_json::from_str(&s).ok())
338            .unwrap_or_default()
339    }
340
341    fn save(&self, pr: &str, st: &WatchState) -> bool {
342        let write = || -> Result<()> {
343            std::fs::create_dir_all(self.dir())?;
344            let path = self.state_path(pr);
345            let tmp = path.with_extension("json.tmp");
346            std::fs::write(&tmp, serde_json::to_string_pretty(st)?)?;
347            std::fs::rename(&tmp, &path)?;
348            Ok(())
349        };
350        match write() {
351            Ok(()) => true,
352            Err(e) => {
353                tracing::warn!("could not save the release watch for {pr}: {e:#}");
354                false
355            }
356        }
357    }
358
359    /// Stored records, for pull requests that left the open list.
360    fn stored(&self) -> Vec<WatchState> {
361        let Ok(rd) = std::fs::read_dir(self.dir()) else {
362            return Vec::new();
363        };
364        rd.flatten()
365            .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
366            .filter_map(|e| std::fs::read_to_string(e.path()).ok())
367            .filter_map(|s| serde_json::from_str(&s).ok())
368            .collect()
369    }
370
371    fn questions(&self) -> Questions {
372        Questions::at(self.home.join("questions"))
373    }
374
375    fn raise(&self, pr: &str, message: String) {
376        notices::raise_in(&self.home, Notice::warn(&notice_key(pr), message));
377    }
378
379    /// One pass over every repo. `stall_secs` is `[daemon] release_stall_minutes`
380    /// in seconds. Best-effort throughout.
381    pub(crate) async fn lap(
382        &self,
383        repos: &[PathBuf],
384        stall_secs: i64,
385        now: i64,
386        halt: &(dyn Fn() -> bool + Sync),
387    ) {
388        let stored = self.stored();
389        // A repo with a watch record is covered even when nothing else names it.
390        let mut repos = repos.to_vec();
391        for s in &stored {
392            let p = PathBuf::from(&s.repo);
393            if !s.repo.is_empty() && !repos.contains(&p) {
394                repos.push(p);
395            }
396        }
397        for repo in &repos {
398            if halt() {
399                return;
400            }
401            let mut urls: Vec<String> = match self.forge.list(repo).await {
402                Ok(v) => v.into_iter().map(|(_, u)| u).collect(),
403                Err(e) => {
404                    tracing::warn!(
405                        "could not list release pull requests in {}: {e:#}",
406                        repo.display()
407                    );
408                    Vec::new()
409                }
410            };
411            // A pull request that left the open list is kept until its end is
412            // confirmed, never assumed.
413            let here = repo.to_string_lossy();
414            for s in stored.iter().filter(|s| s.repo == here.as_ref()) {
415                if !urls.contains(&s.url) {
416                    urls.push(s.url.clone());
417                }
418            }
419            for url in urls {
420                if halt() {
421                    return;
422                }
423                self.watch(repo, &url, stall_secs, now, halt).await;
424            }
425        }
426    }
427
428    /// Look at one pull request and act once.
429    pub(crate) async fn watch(
430        &self,
431        repo: &Path,
432        url: &str,
433        stall_secs: i64,
434        now: i64,
435        halt: &(dyn Fn() -> bool + Sync),
436    ) {
437        let Some(pr) = pr_key(url) else {
438            tracing::warn!("not a pull request url: {url}");
439            return;
440        };
441        let mut st = self.load(&pr);
442        st.repo = repo.to_string_lossy().into_owned();
443        st.url = url.to_owned();
444
445        let snap = match self.forge.snapshot(repo, url).await {
446            Ok(s) => Some(s),
447            Err(e) => {
448                tracing::warn!("could not read {url}: {e:#}");
449                None
450            }
451        };
452        if decide(snap.as_ref(), &st, now, stall_secs) == Step::Done {
453            self.finish(&pr, &st);
454            return;
455        }
456        let Some(snap) = snap else {
457            return;
458        };
459
460        // The owner's word comes first.
461        if let Some(id) = st.question.clone() {
462            match self.questions().get(&id) {
463                Ok(q) if q.status.open() => return,
464                Ok(q) => self.apply_answer(&pr, &mut st, &snap, &q),
465                Err(_) => st.question = None,
466            }
467            // Whatever was decided, the record is saved before anything else.
468            observe(&mut st, &snap, now);
469            self.save(&pr, &st);
470        } else {
471            observe(&mut st, &snap, now);
472        }
473
474        match decide(Some(&snap), &st, now, stall_secs) {
475            Step::Done | Step::Unknown | Step::Wait => {}
476            Step::Rerun(runs) => {
477                // Durable before the call: at most once per failing run.
478                for r in &runs {
479                    st.reruns.insert(r.clone(), now);
480                }
481                // Without a durable record the call must not happen: a restart
482                // would rerun the same run again.
483                if !self.save(&pr, &st) || halt() {
484                    return;
485                }
486                let mut ok = false;
487                for r in &runs {
488                    match self.forge.rerun(repo, r).await {
489                        Ok(()) => ok = true,
490                        Err(e) => tracing::warn!("could not rerun run {r} of {url}: {e:#}"),
491                    }
492                }
493                if ok {
494                    self.raise(&pr, rerun_message(&pr));
495                }
496            }
497            Step::Escalate(why) => {
498                self.raise(&pr, human_message(&pr));
499                let mut q = Question::new(
500                    String::new(),
501                    crate::bump::NOTICE_NODE.to_owned(),
502                    "release-watch".to_owned(),
503                    format!("Release PR {pr} is stuck: what now?"),
504                    question_detail(&snap, why, &st),
505                    vec![RERUN_AGAIN.to_owned(), HOLD.to_owned(), LEAVE_IT.to_owned()],
506                );
507                match self.questions().put(&mut q) {
508                    Ok(()) => st.question = Some(q.id.clone()),
509                    Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
510                }
511            }
512        }
513        self.save(&pr, &st);
514    }
515
516    /// Apply a settled question's outcome to the record, once.
517    fn apply_answer(&self, pr: &str, st: &mut WatchState, snap: &RollupView, q: &Question) {
518        st.question = None;
519        if st.applied.contains(&q.id) {
520            return;
521        }
522        st.applied.push(q.id.clone());
523        if st.applied.len() > 20 {
524            st.applied.remove(0);
525        }
526        match &q.answer {
527            Some(Answer::Choice(c)) if c == RERUN_AGAIN => {
528                st.held = None;
529                for c in snap.checks.iter().filter(|c| is_failed(c.verdict)) {
530                    if let Some(r) = &c.run {
531                        st.reruns.remove(r);
532                    }
533                }
534            }
535            Some(Answer::Choice(c)) if c == LEAVE_IT => {
536                st.ignored = true;
537                if let Err(e) = Notices::at(self.home.join("notifications"))
538                    .dismiss(&notices::id_of(&notice_key(pr)))
539                {
540                    tracing::warn!("could not dismiss the notice for {pr}: {e:#}");
541                }
542            }
543            // `hold`, no answer, abandoned: quiet until something changes.
544            _ => st.held = Some(fingerprint(snap)),
545        }
546    }
547
548    /// The pull request is confirmed merged or closed: clear everything.
549    fn finish(&self, pr: &str, st: &WatchState) {
550        let _ =
551            Notices::at(self.home.join("notifications")).dismiss(&notices::id_of(&notice_key(pr)));
552        if let Some(id) = &st.question {
553            let _ = self.questions().update(id, |q| {
554                q.abandon("the release pull request is no longer open");
555                Ok(())
556            });
557        }
558        let _ = std::fs::remove_file(self.state_path(pr));
559    }
560}
561
562/// The task `magi serve` spawns. `settings` returns the checkouts to cover and
563/// `[daemon] release_stall_minutes`; it is re-read every lap so a config
564/// change is picked up, and `0` minutes switches the lap off.
565pub(crate) async fn run(
566    watcher: Watcher,
567    settings: impl Fn() -> (Vec<PathBuf>, u64),
568    stop: crate::daemon::Stop,
569) {
570    let halt = {
571        let stop = stop.clone();
572        move || stop.stopped()
573    };
574    while !stop.stopped() {
575        let (repos, minutes) = settings();
576        if minutes > 0 {
577            let now = jiff::Timestamp::now().as_second();
578            watcher
579                .lap(&repos, (minutes as i64).saturating_mul(60), now, &halt)
580                .await;
581        }
582        let mut slept = Duration::ZERO;
583        while slept < LAP && !stop.stopped() {
584            tokio::time::sleep(Duration::from_secs(1)).await;
585            slept += Duration::from_secs(1);
586        }
587    }
588}
589
590#[cfg(test)]
591mod tests {
592    use super::*;
593    use crate::land::CheckView;
594    use anyhow::Context;
595    use std::sync::Mutex;
596
597    const URL: &str = "https://github.com/o/r/pull/7";
598
599    fn check(name: &str, v: Verdict, run: Option<&str>) -> CheckView {
600        CheckView {
601            name: name.to_owned(),
602            verdict: v,
603            run: run.map(str::to_owned),
604            url: run.map(|r| format!("https://github.com/o/r/actions/runs/{r}/job/1")),
605        }
606    }
607
608    fn snap(state: PrLifecycle, head: &str, checks: Vec<CheckView>) -> RollupView {
609        RollupView {
610            url: URL.to_owned(),
611            number: 7,
612            state,
613            head: head.to_owned(),
614            checks,
615        }
616    }
617
618    fn open(checks: Vec<CheckView>) -> RollupView {
619        snap(PrLifecycle::Open, "h1", checks)
620    }
621
622    fn st_at(progress: i64) -> WatchState {
623        WatchState {
624            progress_at: progress,
625            ..WatchState::default()
626        }
627    }
628
629    #[test]
630    fn unreadable_is_unknown_and_merged_or_closed_is_done() {
631        assert_eq!(decide(None, &st_at(0), 10, 60), Step::Unknown);
632        for s in [PrLifecycle::Merged, PrLifecycle::Closed] {
633            assert_eq!(
634                decide(Some(&snap(s, "h", vec![])), &st_at(0), 10, 60),
635                Step::Done
636            );
637        }
638    }
639
640    #[test]
641    fn a_failed_run_is_rerun_even_while_another_run_is_pending() {
642        let s = open(vec![
643            check("win", Verdict::Fail, Some("11")),
644            check("lint", Verdict::Pending, Some("12")),
645        ]);
646        assert_eq!(
647            decide(Some(&s), &st_at(0), 1, 3600),
648            Step::Rerun(vec!["11".into()])
649        );
650    }
651
652    #[test]
653    fn a_run_with_a_job_still_pending_is_not_complete() {
654        let s = open(vec![
655            check("win", Verdict::Fail, Some("11")),
656            check("mac", Verdict::Pending, Some("11")),
657        ]);
658        assert_eq!(decide(Some(&s), &st_at(0), 1, 3600), Step::Wait);
659    }
660
661    #[test]
662    fn one_rerun_per_run_then_grace_then_escalation() {
663        let s = open(vec![check("win", Verdict::Fail, Some("11"))]);
664        let mut st = st_at(0);
665        st.reruns.insert("11".into(), 100);
666        assert_eq!(decide(Some(&s), &st, 100 + RERUN_GRACE - 1, 0), Step::Wait);
667        assert_eq!(
668            decide(Some(&s), &st, 100 + RERUN_GRACE, 0),
669            Step::Escalate(Why::StillRed)
670        );
671    }
672
673    #[test]
674    fn a_failure_with_no_run_id_goes_straight_to_a_human_once_settled() {
675        let s = open(vec![check("ci/ext", Verdict::Fail, None)]);
676        assert_eq!(
677            decide(Some(&s), &st_at(0), 1, 0),
678            Step::Escalate(Why::StillRed)
679        );
680        let s = open(vec![
681            check("ci/ext", Verdict::Fail, None),
682            check("x", Verdict::Pending, Some("5")),
683        ]);
684        assert_eq!(decide(Some(&s), &st_at(0), 1, 0), Step::Wait);
685    }
686
687    #[test]
688    fn no_progress_for_the_bounded_time_escalates_and_zero_disables_it() {
689        let s = open(vec![check("a", Verdict::Pass, Some("1"))]);
690        assert_eq!(decide(Some(&s), &st_at(0), 3599, 3600), Step::Wait);
691        assert_eq!(
692            decide(Some(&s), &st_at(0), 3600, 3600),
693            Step::Escalate(Why::Stalled)
694        );
695        assert_eq!(decide(Some(&s), &st_at(0), 99_999, 0), Step::Wait);
696    }
697
698    #[test]
699    fn held_and_ignored_stay_quiet_until_the_fingerprint_moves() {
700        let s = open(vec![check("a", Verdict::Fail, None)]);
701        let mut st = st_at(0);
702        st.held = Some(fingerprint(&s));
703        assert_eq!(decide(Some(&s), &st, 99_999, 60), Step::Wait);
704        let moved = open(vec![
705            check("a", Verdict::Pass, None),
706            check("b", Verdict::Fail, None),
707        ]);
708        assert_eq!(
709            decide(Some(&moved), &st, 99_999, 60),
710            Step::Escalate(Why::StillRed)
711        );
712        st.ignored = true;
713        assert_eq!(decide(Some(&moved), &st, 99_999, 60), Step::Wait);
714    }
715
716    #[test]
717    fn observe_restarts_the_clock_only_on_change_and_forgets_reruns_on_a_new_head() {
718        let mut st = st_at(0);
719        let a = open(vec![check("a", Verdict::Pending, Some("1"))]);
720        observe(&mut st, &a, 10);
721        assert_eq!(st.progress_at, 10);
722        observe(&mut st, &a, 50);
723        assert_eq!(st.progress_at, 10);
724        st.reruns.insert("1".into(), 5);
725        observe(
726            &mut st,
727            &snap(PrLifecycle::Open, "h2", a.checks.clone()),
728            60,
729        );
730        assert!(st.reruns.is_empty());
731        assert_eq!(st.progress_at, 60);
732    }
733
734    #[test]
735    fn pr_key_names_owner_repo_and_number() {
736        assert_eq!(pr_key(URL).as_deref(), Some("o/r#7"));
737        assert_eq!(pr_key("https://github.com/o/r/issues/7"), None);
738        assert_eq!(pr_key("nonsense"), None);
739    }
740
741    #[test]
742    fn notice_wording_does_not_vary_between_polls() {
743        assert_eq!(rerun_message("o/r#7"), rerun_message("o/r#7"));
744        assert!(
745            !human_message("o/r#7")
746                .chars()
747                .any(|c| c.is_ascii_digit() && c != '7')
748        );
749    }
750
751    #[test]
752    fn the_release_question_is_not_claimed_by_other_machinery() {
753        let q = Question::new(
754            String::new(),
755            crate::bump::NOTICE_NODE.into(),
756            "release-watch".into(),
757            "s".into(),
758            String::new(),
759            vec![HOLD.into()],
760        );
761        assert_eq!(crate::deputy::kind_of(&q), None);
762        let n = Notice::warn("release-pr:o/r#7", "m").about([String::new()]);
763        assert!(!notices::covers(&q, &n));
764    }
765
766    #[derive(Default)]
767    struct Fake {
768        snap: Mutex<Option<RollupView>>,
769        reruns: Mutex<Vec<String>>,
770    }
771
772    impl ReleaseForge for std::sync::Arc<Fake> {
773        fn list<'a>(&'a self, _: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
774            Box::pin(async { Ok(vec![("chore/release-v1.0.0".to_owned(), URL.to_owned())]) })
775        }
776        fn snapshot<'a>(&'a self, _: &'a Path, _: &'a str) -> Fut<'a, Result<RollupView>> {
777            let s = self.snap.lock().unwrap().clone();
778            Box::pin(async move { s.context("unreadable") })
779        }
780        fn rerun<'a>(&'a self, _: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
781            self.reruns.lock().unwrap().push(run.to_owned());
782            Box::pin(async { Ok(()) })
783        }
784    }
785
786    fn rig() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher) {
787        let dir = tempfile::tempdir().unwrap();
788        let fake = std::sync::Arc::new(Fake::default());
789        let w = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
790        (dir, fake, w)
791    }
792
793    fn red() -> RollupView {
794        open(vec![check("win", Verdict::Fail, Some("11"))])
795    }
796
797    #[tokio::test]
798    async fn reruns_once_survives_a_restart_then_escalates_with_one_question() {
799        let (dir, fake, w) = rig();
800        *fake.snap.lock().unwrap() = Some(red());
801        let repo = PathBuf::from("/nowhere");
802        let no = || false;
803        w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
804        assert_eq!(*fake.reruns.lock().unwrap(), vec!["11".to_owned()]);
805
806        // A restart is a new Watcher over the same home: no second rerun.
807        let w2 = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
808        w2.lap(
809            std::slice::from_ref(&repo),
810            3600,
811            1000 + RERUN_GRACE - 1,
812            &no,
813        )
814        .await;
815        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
816        assert!(w2.questions().list().is_empty());
817
818        w2.lap(
819            std::slice::from_ref(&repo),
820            3600,
821            1000 + RERUN_GRACE + 1,
822            &no,
823        )
824        .await;
825        w2.lap(
826            std::slice::from_ref(&repo),
827            3600,
828            1000 + RERUN_GRACE + 400,
829            &no,
830        )
831        .await;
832        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
833        let qs = w2.questions().list();
834        assert_eq!(qs.len(), 1, "asked once, not every lap");
835        assert_eq!(qs[0].node, crate::bump::NOTICE_NODE);
836        assert!(qs[0].detail.contains("win") && qs[0].detail.contains(URL));
837        assert_eq!(qs[0].choices, vec![RERUN_AGAIN, HOLD, LEAVE_IT]);
838        let ns = Notices::at(dir.path().join("notifications")).list();
839        assert_eq!(ns.len(), 1);
840        assert_eq!(ns[0].key, "release-pr:o/r#7");
841    }
842
843    async fn escalated() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher, String) {
844        let (dir, fake, w) = rig();
845        *fake.snap.lock().unwrap() = Some(red());
846        let repo = PathBuf::from("/nowhere");
847        let no = || false;
848        w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
849        w.lap(&[repo], 3600, 2000, &no).await;
850        let id = w.questions().list()[0].id.clone();
851        (dir, fake, w, id)
852    }
853
854    fn answer(w: &Watcher, id: &str, choice: &str) {
855        w.questions()
856            .update(id, |q| q.answer(Answer::Choice(choice.to_owned())))
857            .unwrap();
858    }
859
860    #[tokio::test]
861    async fn rerun_again_reruns_exactly_once_more() {
862        let (_d, fake, w, id) = escalated().await;
863        answer(&w, &id, RERUN_AGAIN);
864        let repo = PathBuf::from("/nowhere");
865        let no = || false;
866        w.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
867        w.lap(&[repo], 3600, 3001, &no).await;
868        assert_eq!(fake.reruns.lock().unwrap().len(), 2);
869    }
870
871    #[tokio::test]
872    async fn hold_and_silence_do_not_reask_and_leave_it_dismisses() {
873        let (d, fake, w, id) = escalated().await;
874        answer(&w, &id, HOLD);
875        let repo = PathBuf::from("/nowhere");
876        let no = || false;
877        for t in [3000, 90_000, 200_000] {
878            w.lap(std::slice::from_ref(&repo), 3600, t, &no).await;
879        }
880        assert_eq!(w.questions().list().len(), 1, "held: no second question");
881        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
882
883        let (d2, _f2, w2, id2) = escalated().await;
884        answer(&w2, &id2, LEAVE_IT);
885        w2.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
886        let ns = Notices::at(d2.path().join("notifications")).list();
887        assert!(ns.is_empty(), "dismissed notices are hidden");
888        drop(d);
889    }
890
891    #[tokio::test]
892    async fn merged_clears_the_notice_the_question_and_the_record() {
893        let (d, fake, w, _id) = escalated().await;
894        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
895        let repo = PathBuf::from("/nowhere");
896        w.lap(&[repo], 3600, 5000, &(|| false)).await;
897        assert!(
898            Notices::at(d.path().join("notifications"))
899                .list()
900                .is_empty()
901        );
902        assert!(w.questions().list().iter().all(|q| !q.status.open()));
903        assert!(w.stored().is_empty());
904    }
905
906    #[tokio::test]
907    async fn an_unreadable_forge_changes_nothing() {
908        let (d, fake, w) = rig();
909        *fake.snap.lock().unwrap() = None;
910        w.lap(&[PathBuf::from("/nowhere")], 60, 99_999, &(|| false))
911            .await;
912        assert!(fake.reruns.lock().unwrap().is_empty());
913        assert!(w.questions().list().is_empty());
914        assert!(!d.path().join("release-watch").exists());
915    }
916
917    #[tokio::test]
918    async fn no_rerun_when_the_record_cannot_be_saved() {
919        let (d, fake, w) = rig();
920        *fake.snap.lock().unwrap() = Some(red());
921        // A file where the directory should be makes every save fail.
922        std::fs::write(d.path().join("release-watch"), "x").unwrap();
923        w.lap(&[PathBuf::from("/nowhere")], 3600, 1000, &(|| false))
924            .await;
925        assert!(fake.reruns.lock().unwrap().is_empty());
926    }
927
928    #[tokio::test]
929    async fn a_park_stops_the_lap_before_any_call() {
930        let (_d, fake, w) = rig();
931        *fake.snap.lock().unwrap() = Some(red());
932        w.lap(&[PathBuf::from("/nowhere")], 60, 1000, &(|| true))
933            .await;
934        assert!(fake.reruns.lock().unwrap().is_empty());
935    }
936}