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//! - **`[release] mode = "local"` replaces all of the above for that
29//!   repository's release pull requests.** No check is awaited, reran or
30//!   stalled on. The watcher asks the owner (`merge` / `hold`, bound to the
31//!   head it observed), merges with `--match-head-commit`, and after the merge
32//!   drives [`crate::release_local`] - tag, then the configured commands - with
33//!   its progress in the watch record. A failure is a notice and a question
34//!   (`retry` / `leave it`), never a blind retry. See [`decide_local`].
35//! - **Notice wording is fixed per stage** and carries only the pull request,
36//!   so a repeated poll never relights it; everything that varies (job names,
37//!   links) rides in the question.
38
39use std::collections::BTreeMap;
40use std::future::Future;
41use std::path::{Path, PathBuf};
42use std::pin::Pin;
43use std::time::Duration;
44
45use anyhow::{Result, bail};
46use serde::{Deserialize, Serialize};
47
48use crate::ask::{Answer, Question, Questions};
49use crate::land::{self, PrLifecycle, RollupView, Verdict};
50use crate::notices::{self, Notice, Notices};
51use crate::release_local::{self, Job};
52use crate::run::RunState;
53
54/// Pause between laps.
55const LAP: Duration = Duration::from_secs(300);
56
57/// How long after a rerun the forge may still report the old failure before
58/// that is read as "still red". `gh run rerun` returns before the checks flip
59/// to pending.
60const RERUN_GRACE: i64 = 180;
61
62/// Choice: rerun the failed jobs once more. The three choices are matched
63/// verbatim when applied.
64pub const RERUN_AGAIN: &str = "rerun again";
65/// Choice: stay quiet until a check or the head changes.
66pub const HOLD: &str = "hold";
67/// Choice: stop watching this pull request.
68pub const LEAVE_IT: &str = "leave it";
69/// Choice (local mode): run the failed release again from where it stopped.
70pub const RETRY: &str = "retry";
71
72/// What one look at one pull request concludes.
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub(crate) enum Step {
75    /// Merged or closed: stop watching.
76    Done,
77    /// Nothing to do yet.
78    Wait,
79    /// The forge answer could not be read; change nothing.
80    Unknown,
81    /// Rerun the failed jobs of these workflow runs, once each.
82    Rerun(Vec<String>),
83    /// A human is needed.
84    Escalate(Why),
85}
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
88pub(crate) enum Why {
89    /// Still red after the one rerun, or red with nothing to rerun.
90    StillRed,
91    /// No progress for longer than `[daemon] release_stall_minutes`.
92    Stalled,
93}
94
95/// Everything remembered about one watched pull request.
96#[derive(Debug, Clone, Default, Serialize, Deserialize)]
97#[serde(default)]
98pub(crate) struct WatchState {
99    /// Checkout the pull request belongs to.
100    pub repo: String,
101    pub url: String,
102    /// Head commit the rest of this record is about.
103    pub head: String,
104    /// Workflow runs already rerun for this head, with when (unix seconds).
105    pub reruns: BTreeMap<String, i64>,
106    /// Digest of the head and every check's verdict, and when it last changed.
107    pub fingerprint: String,
108    pub progress_at: i64,
109    /// The open question asked about this pull request.
110    pub question: Option<String>,
111    /// Fingerprint the owner chose to hold: no new question until it changes.
112    pub held: Option<String>,
113    /// The owner said to leave it.
114    pub ignored: bool,
115    /// Question ids whose answer was already applied.
116    pub applied: Vec<String>,
117    /// Local mode: the head the open (or last) approval question was about.
118    pub asked_head: Option<String>,
119    /// Local mode: the head the owner held or whose merge was refused; no new
120    /// approval question until the head moves.
121    pub held_head: Option<String>,
122    /// Local mode: the run that opened this release PR, so a failed release can
123    /// hold the task that run finished (the daemon marks it Done at the merge).
124    pub run: Option<String>,
125    /// Local mode: the task this watcher put on hold for a failed release, so a
126    /// successful retry can give it back (and a restart does not hold it twice).
127    pub held_task: Option<String>,
128    /// Local mode: the release after the merge.
129    pub job: Option<Job>,
130}
131
132fn is_failed(v: Verdict) -> bool {
133    v == Verdict::Fail
134}
135
136/// Digest of what "progress" means: the head and every check's verdict.
137pub(crate) fn fingerprint(snap: &RollupView) -> String {
138    let mut parts: Vec<String> = snap
139        .checks
140        .iter()
141        .map(|c| format!("{}={:?}", c.name, c.verdict))
142        .collect();
143    parts.sort();
144    format!("{}|{}", snap.head, parts.join(","))
145}
146
147/// Fold a fresh snapshot into the record: a new head forgets the old head's
148/// reruns, and a changed fingerprint restarts the stall clock.
149pub(crate) fn observe(st: &mut WatchState, snap: &RollupView, now: i64) {
150    if st.head != snap.head {
151        st.reruns.clear();
152        st.head = snap.head.clone();
153    }
154    let fp = fingerprint(snap);
155    if st.fingerprint != fp {
156        st.fingerprint = fp;
157        st.progress_at = now;
158    }
159}
160
161/// Decide what to do about a pull request. `snap` is `None` when the forge
162/// could not be read. `stall_secs == 0` disables the stall rule.
163pub(crate) fn decide(
164    snap: Option<&RollupView>,
165    st: &WatchState,
166    now: i64,
167    stall_secs: i64,
168) -> Step {
169    let Some(snap) = snap else {
170        return Step::Unknown;
171    };
172    if snap.state != PrLifecycle::Open {
173        return Step::Done;
174    }
175    if st.ignored || st.held.as_deref() == Some(fingerprint(snap).as_str()) {
176        return Step::Wait;
177    }
178
179    let any_pending = snap.checks.iter().any(|c| c.verdict == Verdict::Pending);
180    // A run is complete once none of its own checks is still pending; another
181    // run's pending job says nothing about it.
182    let complete = |run: &str| {
183        !snap
184            .checks
185            .iter()
186            .any(|c| c.verdict == Verdict::Pending && c.run.as_deref() == Some(run))
187    };
188    let mut fresh = Vec::new();
189    let mut spent_recently = false;
190    let mut spent = false;
191    let mut external = false;
192    for c in snap.checks.iter().filter(|c| is_failed(c.verdict)) {
193        match c.run.as_deref() {
194            None => external = true,
195            Some(run) if !complete(run) => {}
196            Some(run) => match st.reruns.get(run) {
197                None => {
198                    if !fresh.iter().any(|r: &String| r == run) {
199                        fresh.push(run.to_owned());
200                    }
201                }
202                Some(&at) if now - at < RERUN_GRACE => spent_recently = true,
203                Some(_) => spent = true,
204            },
205        }
206    }
207    if !fresh.is_empty() {
208        return Step::Rerun(fresh);
209    }
210    if spent_recently {
211        return Step::Wait;
212    }
213    if spent || (external && !any_pending) {
214        return Step::Escalate(Why::StillRed);
215    }
216    if stall_secs > 0 && now - st.progress_at >= stall_secs {
217        return Step::Escalate(Why::Stalled);
218    }
219    Step::Wait
220}
221
222/// `owner/repo#N` out of a pull request url, the stable name for notices.
223pub(crate) fn pr_key(url: &str) -> Option<String> {
224    let rest = url.split("://").nth(1)?;
225    let mut it = rest.split('/');
226    let _host = it.next()?;
227    let owner = it.next()?;
228    let repo = it.next()?;
229    if it.next()? != "pull" {
230        return None;
231    }
232    let n: String = it
233        .next()?
234        .chars()
235        .take_while(char::is_ascii_digit)
236        .collect();
237    (!n.is_empty()).then(|| format!("{owner}/{repo}#{n}"))
238}
239
240fn notice_key(pr: &str) -> String {
241    format!("release-pr:{pr}")
242}
243
244/// Fixed per stage: no counts, times or job names.
245fn rerun_message(pr: &str) -> String {
246    format!("Release PR {pr}: failed checks were rerun once")
247}
248
249fn human_message(pr: &str) -> String {
250    format!("Release PR {pr} is stuck and needs a human")
251}
252
253fn question_detail(snap: &RollupView, why: Why, st: &WatchState) -> String {
254    let mut s = String::new();
255    s.push_str(&format!("Pull request: {}\n\n", snap.url));
256    s.push_str(match why {
257        Why::StillRed => "Checks are still red after the failed jobs were rerun once.\n\n",
258        Why::Stalled => "The pull request has made no progress for too long.\n\n",
259    });
260    let failed: Vec<_> = snap
261        .checks
262        .iter()
263        .filter(|c| is_failed(c.verdict))
264        .collect();
265    if failed.is_empty() {
266        s.push_str("No check is failing.\n");
267    } else {
268        s.push_str("Failed jobs:\n");
269        for c in failed {
270            match &c.url {
271                Some(u) => s.push_str(&format!("- {} ({u})\n", c.name)),
272                None => s.push_str(&format!("- {}\n", c.name)),
273            }
274        }
275    }
276    s.push_str(&format!(
277        "\nWorkflow runs rerun so far: {}.\n\n\
278         - `{RERUN_AGAIN}`: rerun the failed jobs once more.\n\
279         - `{HOLD}`: stay quiet until a check or the head changes.\n\
280         - `{LEAVE_IT}`: stop watching this pull request.\n\n\
281         Nothing is merged or pushed by magi; silence is a hold.\n",
282        st.reruns.len()
283    ));
284    s
285}
286
287/// Where the owner's approval stands, as far as [`decide_local`] cares.
288#[derive(Debug, Clone, Copy, PartialEq, Eq)]
289pub(crate) enum Approved {
290    /// No question was filed (or it is gone).
291    NoQuestion,
292    /// Filed, not answered.
293    Open,
294    /// The owner said `merge`.
295    Merge,
296    /// The owner said anything else, or it was abandoned: silence is a hold.
297    Hold,
298}
299
300/// What one look at an open local-mode release pull request concludes.
301#[derive(Debug, Clone, PartialEq, Eq)]
302pub(crate) enum LocalStep {
303    /// Nothing to do now.
304    Wait,
305    /// Ask the owner whether to merge this head.
306    Ask,
307    /// Merge, pinned to this head.
308    Merge(String),
309    /// Remember the owner's no for this head and stay quiet until it moves.
310    Hold(String),
311}
312
313/// Decide an *open* local-mode pull request. No I/O, and no CI: nothing is
314/// awaited, rerun or stalled on, because nothing will ever report. The head is
315/// what the owner approved: an answer about an older head asks again.
316pub(crate) fn decide_local(head: &str, st: &WatchState, approved: Approved) -> LocalStep {
317    if st.ignored || st.held_head.as_deref() == Some(head) {
318        return LocalStep::Wait;
319    }
320    let about_this_head = st.asked_head.as_deref() == Some(head);
321    match approved {
322        Approved::Open => LocalStep::Wait,
323        Approved::NoQuestion => LocalStep::Ask,
324        Approved::Merge if about_this_head => LocalStep::Merge(head.to_owned()),
325        Approved::Hold if about_this_head => LocalStep::Hold(head.to_owned()),
326        // The head moved after the question was filed.
327        Approved::Merge | Approved::Hold => LocalStep::Ask,
328    }
329}
330
331/// Start watching a release pull request before the first lap sees it, so one
332/// merged within a lap of being opened still gets released. Never overwrites.
333pub(crate) fn register(home: &Path, repo: &Path, url: &str, run: &str) {
334    let Some(pr) = pr_key(url) else {
335        return;
336    };
337    let w = Watcher::new(Box::new(GhForge), home.to_path_buf());
338    if w.state_path(&pr).exists() {
339        return;
340    }
341    let st = WatchState {
342        repo: repo.to_string_lossy().into_owned(),
343        url: url.to_owned(),
344        run: Some(run.to_owned()),
345        ..WatchState::default()
346    };
347    w.save(&pr, &st);
348}
349
350/// The watch record whose open question is `question_id`, read from the magi
351/// `home`. `None` when no record names it (never guessed from the text).
352pub(crate) fn state_for_question(home: &Path, question_id: &str) -> Option<WatchState> {
353    Watcher::new(Box::new(GhForge), home.to_path_buf())
354        .stored()
355        .into_iter()
356        .find(|s| s.question.as_deref() == Some(question_id))
357}
358
359/// The task a release-watch question's answer could touch: the one that
360/// finished the run which opened the release pull request (`WatchState::run`).
361/// `None` when the record or the run is not known - then there is nothing to
362/// check, and the settle goes on (an answer here never releases a task).
363pub fn task_of_question(
364    home: &Path,
365    q: &Question,
366    queue: &crate::queue::Queue,
367) -> Option<crate::queue::Task> {
368    let run = state_for_question(home, &q.id)?.run?;
369    queue.list().into_iter().find(|t| t.runs.contains(&run))
370}
371
372/// What a deputy is told about a release-watch question: which pull request,
373/// and what each offered choice really does - matched to [`Watcher::apply_answer`]
374/// and [`Watcher::settle_job_question`]. A missing watch record is said to be
375/// missing, never reconstructed.
376pub(crate) fn deputy_brief(q: &Question, home: &Path) -> String {
377    let mut s = String::from(
378        "This question was filed by magi's release watcher about a release pull \
379         request. Whatever the owner picks, only the watcher acts on it, on its \
380         next lap; you apply nothing and never close, merge, rerun or push \
381         anything. Silence is a hold.",
382    );
383    match state_for_question(home, &q.id) {
384        Some(st) => {
385            s.push_str(&format!(
386                "\n\nPull request: {}\nCheckout: {}",
387                st.url, st.repo
388            ));
389        }
390        None => s.push_str(
391            "\n\nThe watcher's record for this question could not be found, so the \
392             pull request is known to you only through the question's own text. \
393             Say so to the owner rather than guessing.",
394        ),
395    }
396    s.push_str("\n\nWhat each option does:");
397    for c in &q.choices {
398        let what = match c.as_str() {
399            RERUN_AGAIN => "forget the reruns already made and rerun the failed jobs once more",
400            HOLD if q.choices.iter().any(|c| c == land::APPROVE) => {
401                "leave the pull request open and unmerged; ask again only when the head moves"
402            }
403            HOLD => "stay quiet until a check or the head changes",
404            LEAVE_IT => {
405                "stop watching this pull request and dismiss its notice. It does NOT \
406                 close the pull request"
407            }
408            RETRY => "run the failed release again from where it stopped",
409            land::APPROVE => {
410                "squash-merge exactly the head the question names - irreversible - \
411                 then tag and release"
412            }
413            _ => "recorded as the answer",
414        };
415        s.push_str(&format!("\n- `{c}`: {what}"));
416    }
417    s.push_str(
418        "\n\nIf the owner's words ask for the pull request itself to be closed, that \
419         is something only they can do by hand; do not read it as a choice unless \
420         it plainly means stop watching.",
421    );
422    s
423}
424
425/// Starts the hold reason of a task held for a failed release, so only that
426/// hold is undone when the release later succeeds.
427const HOLD_PREFIX: &str = "[release] ";
428
429fn approval_message(pr: &str) -> String {
430    format!("Release PR {pr} is waiting for your approval to merge")
431}
432
433fn failed_message(pr: &str) -> String {
434    format!("Release PR {pr} merged but the release is on hold")
435}
436
437fn released_key(pr: &str) -> String {
438    format!("release-done:{pr}")
439}
440
441type Fut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
442
443/// Every forge access, so tests inject a fake instead of `gh`.
444pub(crate) trait ReleaseForge: Send + Sync {
445    /// `(branch, url)` of every open release-bump pull request in `repo`.
446    fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>>;
447    fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>>;
448    fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>>;
449    /// The repository's config, to read `[release]` and `[merge] remote`.
450    fn config<'a>(&'a self, _repo: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
451        Box::pin(async { Ok(crate::config::Config::default()) })
452    }
453    /// Head branch and merge commit of a pull request (local mode).
454    fn info<'a>(&'a self, _repo: &'a Path, url: &'a str) -> Fut<'a, Result<PrInfo>> {
455        Box::pin(async move { bail!("no pull request info for {url}") })
456    }
457    /// Merge `url` (squash, delete branch) pinned to `head`. `Ok(true)` once
458    /// the forge says it merged, `Ok(false)` when it did not.
459    fn merge<'a>(&'a self, _repo: &'a Path, url: &'a str, _head: &'a str) -> Fut<'a, Result<bool>> {
460        Box::pin(async move { bail!("cannot merge {url}") })
461    }
462}
463
464/// What the forge says about a pull request's branch and merge.
465#[derive(Debug, Clone, Default, PartialEq, Eq)]
466pub(crate) struct PrInfo {
467    /// Head branch name, e.g. `chore/release-v1.2.3`.
468    pub branch: String,
469    /// The commit the pull request merged as, once merged.
470    pub merge_commit: Option<String>,
471    /// Pull request title.
472    pub title: String,
473}
474
475/// Parse `gh pr view --json headRefName,mergeCommit,title`. No I/O.
476pub(crate) fn parse_info(json: &str) -> Result<PrInfo> {
477    let v: serde_json::Value = serde_json::from_str(json)?;
478    let text = |k: &str| {
479        v.get(k)
480            .and_then(|x| x.as_str())
481            .unwrap_or_default()
482            .to_owned()
483    };
484    Ok(PrInfo {
485        branch: text("headRefName"),
486        title: text("title"),
487        merge_commit: v
488            .get("mergeCommit")
489            .and_then(|m| m.get("oid"))
490            .and_then(|o| o.as_str())
491            .filter(|o| !o.is_empty())
492            .map(str::to_owned),
493    })
494}
495
496/// The real forge: `gh`.
497pub(crate) struct GhForge;
498
499impl ReleaseForge for GhForge {
500    fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
501        Box::pin(crate::bump::list_open_release_prs(repo))
502    }
503
504    fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>> {
505        Box::pin(async move {
506            let args = [
507                "pr".to_owned(),
508                "view".to_owned(),
509                url.to_owned(),
510                "--json".to_owned(),
511                "url,number,state,headRefOid,statusCheckRollup".to_owned(),
512            ];
513            let (ok, out) = land::gh(repo, &args).await?;
514            if !ok {
515                bail!("gh pr view {url}: {out}");
516            }
517            land::parse_rollup(&out)
518        })
519    }
520
521    fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
522        Box::pin(async move {
523            let args = [
524                "run".to_owned(),
525                "rerun".to_owned(),
526                run.to_owned(),
527                "--failed".to_owned(),
528            ];
529            let (ok, out) = land::gh(repo, &args).await?;
530            if !ok {
531                bail!("gh run rerun {run}: {out}");
532            }
533            Ok(())
534        })
535    }
536
537    fn config<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
538        Box::pin(async move { Ok(crate::config::Config::discover(repo, None)?.0) })
539    }
540
541    fn info<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<PrInfo>> {
542        Box::pin(async move {
543            let args = [
544                "pr".to_owned(),
545                "view".to_owned(),
546                url.to_owned(),
547                "--json".to_owned(),
548                "headRefName,mergeCommit,title".to_owned(),
549            ];
550            let (ok, out) = land::gh(repo, &args).await?;
551            if !ok {
552                bail!("gh pr view {url}: {out}");
553            }
554            parse_info(&out)
555        })
556    }
557
558    fn merge<'a>(&'a self, repo: &'a Path, url: &'a str, head: &'a str) -> Fut<'a, Result<bool>> {
559        Box::pin(async move {
560            let title = self
561                .info(repo, url)
562                .await
563                .map(|i| i.title)
564                .unwrap_or_default();
565            let title = if title.trim().is_empty() {
566                "chore: release".to_owned()
567            } else {
568                title
569            };
570            // Same shape as the direct merge in `bump`, pinned to the head the
571            // owner approved: a push in between is refused by the forge.
572            let mut argv = crate::bump::bump_merge_argv(url, &title);
573            argv.push("--match-head-commit".to_owned());
574            argv.push(head.to_owned());
575            let (ok, out) = land::gh(repo, &argv).await?;
576            if ok {
577                return Ok(true);
578            }
579            // jj keeps HEAD detached, so `--delete-branch` exits non-zero after
580            // the merge happened: the forge decides.
581            let after = land::lifecycle(repo, url).await.ok();
582            if land::merged_after_all(&argv, &out, after).is_some() {
583                return Ok(true);
584            }
585            bail!("gh pr merge {url}: {out}")
586        })
587    }
588}
589
590/// The watcher: a forge, a magi home, and nothing else.
591pub(crate) struct Watcher {
592    forge: Box<dyn ReleaseForge>,
593    home: PathBuf,
594}
595
596impl Watcher {
597    pub(crate) fn new(forge: Box<dyn ReleaseForge>, home: PathBuf) -> Self {
598        Self { forge, home }
599    }
600
601    fn dir(&self) -> PathBuf {
602        self.home.join("release-watch")
603    }
604
605    fn state_path(&self, pr: &str) -> PathBuf {
606        self.dir().join(format!("{}.json", notices::id_of(pr)))
607    }
608
609    fn load(&self, pr: &str) -> WatchState {
610        std::fs::read_to_string(self.state_path(pr))
611            .ok()
612            .and_then(|s| serde_json::from_str(&s).ok())
613            .unwrap_or_default()
614    }
615
616    fn save(&self, pr: &str, st: &WatchState) -> bool {
617        let write = || -> Result<()> {
618            std::fs::create_dir_all(self.dir())?;
619            let path = self.state_path(pr);
620            let tmp = path.with_extension("json.tmp");
621            std::fs::write(&tmp, serde_json::to_string_pretty(st)?)?;
622            std::fs::rename(&tmp, &path)?;
623            Ok(())
624        };
625        match write() {
626            Ok(()) => true,
627            Err(e) => {
628                tracing::warn!("could not save the release watch for {pr}: {e:#}");
629                false
630            }
631        }
632    }
633
634    /// Stored records, for pull requests that left the open list.
635    fn stored(&self) -> Vec<WatchState> {
636        let Ok(rd) = std::fs::read_dir(self.dir()) else {
637            return Vec::new();
638        };
639        rd.flatten()
640            .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
641            .filter_map(|e| std::fs::read_to_string(e.path()).ok())
642            .filter_map(|s| serde_json::from_str(&s).ok())
643            .collect()
644    }
645
646    fn questions(&self) -> Questions {
647        Questions::at(self.home.join("questions"))
648    }
649
650    fn raise(&self, pr: &str, message: String) {
651        notices::raise_in(&self.home, Notice::warn(&notice_key(pr), message));
652    }
653
654    /// One pass over every repo. `stall_secs` is `[daemon] release_stall_minutes`
655    /// in seconds. Best-effort throughout.
656    pub(crate) async fn lap(
657        &self,
658        repos: &[PathBuf],
659        stall_secs: i64,
660        now: i64,
661        halt: &(dyn Fn() -> bool + Sync),
662    ) {
663        let stored = self.stored();
664        // A repo with a watch record is covered even when nothing else names it.
665        let mut repos = repos.to_vec();
666        for s in &stored {
667            let p = PathBuf::from(&s.repo);
668            if !s.repo.is_empty() && !repos.contains(&p) {
669                repos.push(p);
670            }
671        }
672        for repo in &repos {
673            if halt() {
674                return;
675            }
676            let mut urls: Vec<String> = match self.forge.list(repo).await {
677                Ok(v) => v.into_iter().map(|(_, u)| u).collect(),
678                Err(e) => {
679                    tracing::warn!(
680                        "could not list release pull requests in {}: {e:#}",
681                        repo.display()
682                    );
683                    Vec::new()
684                }
685            };
686            // A pull request that left the open list is kept until its end is
687            // confirmed, never assumed.
688            let here = repo.to_string_lossy();
689            for s in stored.iter().filter(|s| s.repo == here.as_ref()) {
690                if !urls.contains(&s.url) {
691                    urls.push(s.url.clone());
692                }
693            }
694            for url in urls {
695                if halt() {
696                    return;
697                }
698                self.watch(repo, &url, stall_secs, now, halt).await;
699            }
700        }
701    }
702
703    /// Look at one pull request and act once.
704    pub(crate) async fn watch(
705        &self,
706        repo: &Path,
707        url: &str,
708        stall_secs: i64,
709        now: i64,
710        halt: &(dyn Fn() -> bool + Sync),
711    ) {
712        let Some(pr) = pr_key(url) else {
713            tracing::warn!("not a pull request url: {url}");
714            return;
715        };
716        let mut st = self.load(&pr);
717        st.repo = repo.to_string_lossy().into_owned();
718        st.url = url.to_owned();
719
720        let snap = match self.forge.snapshot(repo, url).await {
721            Ok(s) => Some(s),
722            Err(e) => {
723                tracing::warn!("could not read {url}: {e:#}");
724                None
725            }
726        };
727        match self.forge.config(repo).await {
728            Ok(cfg) if cfg.release.is_local() => {
729                self.watch_local(repo, &pr, st, snap, &cfg).await;
730                return;
731            }
732            Ok(_) => {}
733            Err(e) => tracing::warn!(
734                "could not read the config of {}: {e:#}; watching {url} as an Actions release",
735                repo.display()
736            ),
737        }
738        if decide(snap.as_ref(), &st, now, stall_secs) == Step::Done {
739            self.finish(&pr, &st);
740            return;
741        }
742        let Some(snap) = snap else {
743            return;
744        };
745
746        // The owner's word comes first.
747        if let Some(id) = st.question.clone() {
748            match self.questions().get(&id) {
749                Ok(q) if q.status.open() => return,
750                Ok(q) => self.apply_answer(&pr, &mut st, &snap, &q),
751                Err(_) => st.question = None,
752            }
753            // Whatever was decided, the record is saved before anything else.
754            observe(&mut st, &snap, now);
755            self.save(&pr, &st);
756        } else {
757            observe(&mut st, &snap, now);
758        }
759
760        match decide(Some(&snap), &st, now, stall_secs) {
761            Step::Done | Step::Unknown | Step::Wait => {}
762            Step::Rerun(runs) => {
763                // Durable before the call: at most once per failing run.
764                for r in &runs {
765                    st.reruns.insert(r.clone(), now);
766                }
767                // Without a durable record the call must not happen: a restart
768                // would rerun the same run again.
769                if !self.save(&pr, &st) || halt() {
770                    return;
771                }
772                let mut ok = false;
773                for r in &runs {
774                    match self.forge.rerun(repo, r).await {
775                        Ok(()) => ok = true,
776                        Err(e) => tracing::warn!("could not rerun run {r} of {url}: {e:#}"),
777                    }
778                }
779                if ok {
780                    self.raise(&pr, rerun_message(&pr));
781                }
782            }
783            Step::Escalate(why) => {
784                self.raise(&pr, human_message(&pr));
785                let mut q = Question::new(
786                    String::new(),
787                    crate::bump::NOTICE_NODE.to_owned(),
788                    "release-watch".to_owned(),
789                    format!("Release PR {pr} is stuck: what now?"),
790                    question_detail(&snap, why, &st),
791                    vec![RERUN_AGAIN.to_owned(), HOLD.to_owned(), LEAVE_IT.to_owned()],
792                );
793                match self.questions().put(&mut q) {
794                    Ok(()) => st.question = Some(q.id.clone()),
795                    Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
796                }
797            }
798        }
799        self.save(&pr, &st);
800    }
801
802    /// The local-mode path of [`Watcher::watch`]: approval, merge, release.
803    async fn watch_local(
804        &self,
805        repo: &Path,
806        pr: &str,
807        mut st: WatchState,
808        snap: Option<RollupView>,
809        cfg: &crate::config::Config,
810    ) {
811        let Some(snap) = snap else {
812            return;
813        };
814        match snap.state {
815            PrLifecycle::Closed => self.finish(pr, &st),
816            PrLifecycle::Merged => self.release_merged(repo, pr, st, cfg).await,
817            PrLifecycle::Open => {
818                let approved = match st.question.clone() {
819                    None => Approved::NoQuestion,
820                    Some(id) => match self.questions().get(&id) {
821                        Ok(q) if q.status.open() => Approved::Open,
822                        Ok(q) => match &q.answer {
823                            Some(Answer::Choice(c))
824                                if land::approval(Some(c)) == land::Approval::Merge =>
825                            {
826                                Approved::Merge
827                            }
828                            _ => Approved::Hold,
829                        },
830                        Err(_) => Approved::NoQuestion,
831                    },
832                };
833                let step = decide_local(&snap.head, &st, approved);
834                if matches!(
835                    step,
836                    LocalStep::Merge(_) | LocalStep::Hold(_) | LocalStep::Ask
837                ) {
838                    // A settled (or stale) question is consumed exactly once.
839                    st.question = None;
840                }
841                match step {
842                    LocalStep::Wait => {}
843                    LocalStep::Hold(head) => st.held_head = Some(head),
844                    LocalStep::Ask => {
845                        st.held_head = None;
846                        st.asked_head = Some(snap.head.clone());
847                        self.raise(pr, approval_message(pr));
848                        let mut q = Question::new(
849                            String::new(),
850                            crate::bump::NOTICE_NODE.to_owned(),
851                            "release-watch".to_owned(),
852                            format!("Release PR {pr}: merge it and release?"),
853                            format!(
854                                "Pull request: {}\nHead: {}\n\n\
855                                 `[release] mode = \"local\"`: no CI is awaited. `{}` merges \
856                                 exactly this head, then magi tags the merge commit and runs the \
857                                 configured release commands. `{}` leaves the pull request \
858                                 open. Silence is a hold.\n",
859                                snap.url,
860                                snap.head,
861                                land::APPROVE,
862                                land::HOLD
863                            ),
864                            vec![land::APPROVE.to_owned(), land::HOLD.to_owned()],
865                        );
866                        match self.questions().put(&mut q) {
867                            Ok(()) => st.question = Some(q.id.clone()),
868                            Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
869                        }
870                    }
871                    LocalStep::Merge(head) => {
872                        match self.forge.merge(repo, &snap.url, &head).await {
873                            Ok(_) => {
874                                // Release at once rather than a lap later.
875                                self.save(pr, &st);
876                                self.release_merged(repo, pr, st, cfg).await;
877                                return;
878                            }
879                            Err(e) => {
880                                tracing::warn!("could not merge {}: {e:#}", snap.url);
881                                st.held_head = Some(head);
882                                self.raise(
883                                    pr,
884                                    format!(
885                                        "Release PR {pr} could not be merged; merge it by hand"
886                                    ),
887                                );
888                            }
889                        }
890                    }
891                }
892                self.save(pr, &st);
893            }
894        }
895    }
896
897    /// The pull request is merged: create the job once, then drive it.
898    async fn release_merged(
899        &self,
900        repo: &Path,
901        pr: &str,
902        mut st: WatchState,
903        cfg: &crate::config::Config,
904    ) {
905        if st.ignored {
906            return self.finish(pr, &st);
907        }
908        if st.job.is_none() {
909            let info = match self.forge.info(repo, &st.url).await {
910                Ok(i) => i,
911                Err(e) => {
912                    tracing::warn!("could not read the merged {}: {e:#}", st.url);
913                    return;
914                }
915            };
916            let (Some(version), Some(commit)) = (
917                release_local::version_from_branch(&info.branch),
918                info.merge_commit,
919            ) else {
920                // Not a release branch, or the forge has not recorded the merge
921                // commit yet: nothing to release (yet).
922                if release_local::version_from_branch(&info.branch).is_none() {
923                    self.finish(pr, &st);
924                }
925                return;
926            };
927            st.job = Some(Job::new(&version, &st.url, &commit));
928            if !self.save(pr, &st) {
929                return;
930            }
931        }
932        let Some(mut job) = st.job.take() else {
933            return;
934        };
935        if job.finished {
936            // Saved as finished, but a stop may have come before the task was
937            // given back: do that now, and keep the record until it is.
938            st.job = Some(job);
939            if !self.record_on_run(&st) {
940                self.save(pr, &st);
941                return;
942            }
943            return self.complete(pr, &mut st);
944        }
945        if job.interrupted() {
946            job.failed = Some(format!(
947                "magi stopped while `{}` was running; it may have partly run, so it is not repeated on its own",
948                job.running.clone().unwrap_or_default()
949            ));
950        }
951        if let Some(why) = job.failed.clone() {
952            // Reconcile first: a stop between saving the failure and holding
953            // the task (or an interrupted step) must not leave the task Done.
954            self.hold_task(&mut st, &format!("release {} is on hold: {why}", job.tag()));
955            // Held: only the owner's `retry` runs it again.
956            let retry = self.settle_job_question(pr, &mut st, &mut job);
957            if !retry {
958                st.job = Some(job);
959                self.save(pr, &st);
960                return;
961            }
962        }
963        let shell = cfg.shell();
964        let env = release_local::Env {
965            repo,
966            home: &self.home,
967            key: pr,
968            remote: &cfg.merge.remote,
969            shell: &shell,
970            release: &cfg.release,
971        };
972        let result = {
973            let mut save = |j: &Job| {
974                let mut copy = st.clone();
975                copy.job = Some(j.clone());
976                self.save(pr, &copy)
977            };
978            release_local::run_job(&env, &mut job, &mut save).await
979        };
980        match result {
981            Ok(()) => {
982                self.raise_released(pr, &job);
983                st.job = Some(job);
984                if self.record_on_run(&st) {
985                    self.complete(pr, &mut st);
986                } else {
987                    // Finished and saved as such: the next lap retries only
988                    // the run record, never the commands.
989                    self.save(pr, &st);
990                }
991            }
992            Err(e) => {
993                job.failed = Some(format!("{e:#}"));
994                self.hold_task(&mut st, &format!("release {} failed: {e:#}", job.tag()));
995                self.raise(pr, failed_message(pr));
996                self.file_job_question(pr, &mut st, &job);
997                st.job = Some(job);
998                self.save(pr, &st);
999                self.record_on_run(&st);
1000            }
1001        }
1002    }
1003
1004    /// Copy the job onto the run that opened the release PR, so what the
1005    /// release did outlives the watcher's record. Returns whether the record
1006    /// may go: `false` only when a readable run could not be saved. No run, or
1007    /// a run that is gone, is warned about and counts as done.
1008    fn record_on_run(&self, st: &WatchState) -> bool {
1009        let (Some(run), Some(job)) = (st.run.as_deref(), st.job.as_ref()) else {
1010            return true;
1011        };
1012        let mut state = match RunState::load_under(run, &self.home) {
1013            Ok(s) => s,
1014            Err(e) => {
1015                tracing::warn!("could not record the release on run {run}: {e:#}");
1016                return true;
1017            }
1018        };
1019        let bump = state.release_bump.get_or_insert_with(Default::default);
1020        if bump.release.as_ref() == Some(job) {
1021            return true;
1022        }
1023        bump.release = Some(job.clone());
1024        state.event(
1025            "release",
1026            format!(
1027                "release {} {}",
1028                job.tag(),
1029                if job.finished {
1030                    "finished"
1031                } else if let Some(why) = &job.failed {
1032                    why
1033                } else {
1034                    "stopped"
1035                }
1036            ),
1037        );
1038        match state.save_under(&self.home) {
1039            Ok(()) => true,
1040            Err(e) => {
1041                tracing::warn!("could not save the release on run {run}: {e:#}");
1042                false
1043            }
1044        }
1045    }
1046
1047    /// Hold the task whose run opened this release PR: the daemon marked it
1048    /// Done when the change merged, but the release is not done. Best-effort;
1049    /// a task the owner already holds or that cannot be found is left alone.
1050    fn hold_task(&self, st: &mut WatchState, reason: &str) {
1051        if st.held_task.is_some() {
1052            return;
1053        }
1054        let Some(run) = st.run.clone() else {
1055            return;
1056        };
1057        let q = crate::queue::Queue::at(self.home.join("queue"));
1058        for mut t in q.list() {
1059            if t.runs.contains(&run) {
1060                if t.status != crate::queue::TaskStatus::Held {
1061                    t.hold_machine(Some(format!("{HOLD_PREFIX}{reason}")));
1062                    match q.put(&mut t) {
1063                        Ok(()) => st.held_task = Some(t.id.clone()),
1064                        Err(e) => tracing::warn!("could not hold task {} for {run}: {e:#}", t.id),
1065                    }
1066                }
1067                return;
1068            }
1069        }
1070    }
1071
1072    /// Give back the hold [`Watcher::hold_task`] put on a task, once the
1073    /// release succeeded. Only a machine hold carrying our marker is undone: a
1074    /// task the owner has since re-held or changed is left alone.
1075    ///
1076    /// Returns whether nothing is left to give back: `false` only when the
1077    /// write failed, in which case `held_task` is kept so a later lap retries.
1078    #[must_use]
1079    fn unhold_task(&self, st: &mut WatchState) -> bool {
1080        let Some(id) = st.held_task.take() else {
1081            return true;
1082        };
1083        let q = crate::queue::Queue::at(self.home.join("queue"));
1084        let mut t = match q.get(&id) {
1085            Ok(t) => t,
1086            // Only a task file confirmed absent is "nothing to restore";
1087            // any other failure may be transient, so keep the record.
1088            Err(_) if matches!(q.path_of(&id).try_exists(), Ok(false)) => return true,
1089            Err(e) => {
1090                tracing::warn!("could not read task {id} to restore it after the release: {e:#}");
1091                st.held_task = Some(id);
1092                return false;
1093            }
1094        };
1095        let ours = t.status == crate::queue::TaskStatus::Held
1096            && t.hold_source == Some(crate::queue::HoldSource::Machine)
1097            && t.hold_reason
1098                .as_deref()
1099                .is_some_and(|r| r.starts_with(HOLD_PREFIX));
1100        if ours {
1101            t.succeed();
1102            if let Err(e) = q.put(&mut t) {
1103                tracing::warn!("could not restore task {id} after the release: {e:#}");
1104                st.held_task = Some(id);
1105                return false;
1106            }
1107        }
1108        true
1109    }
1110
1111    /// The release is done (`st.job` is finished): give the task back, then
1112    /// drop the record. If the task cannot be restored the record stays, so
1113    /// the next lap tries again instead of leaving the task Held for good.
1114    fn complete(&self, pr: &str, st: &mut WatchState) {
1115        if self.unhold_task(st) {
1116            self.finish(pr, st);
1117        } else {
1118            self.save(pr, st);
1119        }
1120    }
1121
1122    fn raise_released(&self, pr: &str, job: &Job) {
1123        notices::raise_in(
1124            &self.home,
1125            Notice::info(&released_key(pr), format!("Released {} ({pr})", job.tag())),
1126        );
1127    }
1128
1129    /// Ask what to do about a held job. One question per failure.
1130    fn file_job_question(&self, pr: &str, st: &mut WatchState, job: &Job) {
1131        let last = job
1132            .log
1133            .last()
1134            .map(|l| format!("\nLast step: {}\n\n{}\n", l.name, l.tail))
1135            .unwrap_or_default();
1136        let mut q = Question::new(
1137            String::new(),
1138            crate::bump::NOTICE_NODE.to_owned(),
1139            "release-watch".to_owned(),
1140            format!("Release {} of {pr} is on hold: retry?", job.tag()),
1141            format!(
1142                "Pull request: {}\nWhy it stopped: {}\n{last}\n\
1143                 Full output is under the magi home's release-local directory.\n\n\
1144                 - `{RETRY}`: run it again from where it stopped (the tag is not \
1145                 recreated and finished commands are skipped).\n\
1146                 - `{LEAVE_IT}`: stop watching; release by hand.\n\n\
1147                 Nothing is retried on its own; silence is a hold.\n",
1148                job.pr_url,
1149                job.failed.clone().unwrap_or_default()
1150            ),
1151            vec![RETRY.to_owned(), LEAVE_IT.to_owned()],
1152        );
1153        match self.questions().put(&mut q) {
1154            Ok(()) => st.question = Some(q.id.clone()),
1155            Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
1156        }
1157    }
1158
1159    /// Read the owner's word on a held job. `true` means `retry`: the job was
1160    /// resumed. Files the question when there is none yet and the owner has not
1161    /// already held this exact failure.
1162    fn settle_job_question(&self, pr: &str, st: &mut WatchState, job: &mut Job) -> bool {
1163        let why = job.failed.clone().unwrap_or_default();
1164        if let Some(id) = st.question.clone() {
1165            match self.questions().get(&id) {
1166                Ok(q) if q.status.open() => return false,
1167                Ok(q) => {
1168                    st.question = None;
1169                    match &q.answer {
1170                        Some(Answer::Choice(c)) if c == RETRY => {
1171                            job.resume();
1172                            st.held = None;
1173                            return true;
1174                        }
1175                        Some(Answer::Choice(c)) if c == LEAVE_IT => {
1176                            st.ignored = true;
1177                            let _ = Notices::at(self.home.join("notifications"))
1178                                .dismiss(&notices::id_of(&notice_key(pr)));
1179                        }
1180                        // Abandoned or anything else: this failure stays held.
1181                        _ => st.held = Some(why),
1182                    }
1183                    return false;
1184                }
1185                Err(_) => st.question = None,
1186            }
1187        }
1188        if !st.ignored && st.held.as_deref() != Some(why.as_str()) {
1189            self.raise(pr, failed_message(pr));
1190            self.file_job_question(pr, st, job);
1191        }
1192        false
1193    }
1194
1195    /// Apply a settled question's outcome to the record, once.
1196    fn apply_answer(&self, pr: &str, st: &mut WatchState, snap: &RollupView, q: &Question) {
1197        st.question = None;
1198        if st.applied.contains(&q.id) {
1199            return;
1200        }
1201        st.applied.push(q.id.clone());
1202        if st.applied.len() > 20 {
1203            st.applied.remove(0);
1204        }
1205        match &q.answer {
1206            Some(Answer::Choice(c)) if c == RERUN_AGAIN => {
1207                st.held = None;
1208                for c in snap.checks.iter().filter(|c| is_failed(c.verdict)) {
1209                    if let Some(r) = &c.run {
1210                        st.reruns.remove(r);
1211                    }
1212                }
1213            }
1214            Some(Answer::Choice(c)) if c == LEAVE_IT => {
1215                st.ignored = true;
1216                if let Err(e) = Notices::at(self.home.join("notifications"))
1217                    .dismiss(&notices::id_of(&notice_key(pr)))
1218                {
1219                    tracing::warn!("could not dismiss the notice for {pr}: {e:#}");
1220                }
1221            }
1222            // `hold`, no answer, abandoned: quiet until something changes.
1223            _ => st.held = Some(fingerprint(snap)),
1224        }
1225    }
1226
1227    /// The pull request is confirmed merged or closed: clear everything.
1228    fn finish(&self, pr: &str, st: &WatchState) {
1229        let _ =
1230            Notices::at(self.home.join("notifications")).dismiss(&notices::id_of(&notice_key(pr)));
1231        if let Some(id) = &st.question {
1232            let _ = self.questions().update(id, |q| {
1233                q.abandon("the release pull request is no longer open");
1234                Ok(())
1235            });
1236        }
1237        let _ = std::fs::remove_file(self.state_path(pr));
1238    }
1239}
1240
1241/// The task `magi serve` spawns. `settings` returns the checkouts to cover and
1242/// `[daemon] release_stall_minutes`; it is re-read every lap so a config
1243/// change is picked up, and `0` minutes switches the lap off.
1244pub(crate) async fn run(
1245    watcher: Watcher,
1246    settings: impl Fn() -> (Vec<PathBuf>, u64),
1247    stop: crate::daemon::Stop,
1248) {
1249    let halt = {
1250        let stop = stop.clone();
1251        move || stop.stopped()
1252    };
1253    while !stop.stopped() {
1254        let (repos, minutes) = settings();
1255        if minutes > 0 {
1256            let now = jiff::Timestamp::now().as_second();
1257            watcher
1258                .lap(&repos, (minutes as i64).saturating_mul(60), now, &halt)
1259                .await;
1260        }
1261        let mut slept = Duration::ZERO;
1262        while slept < LAP && !stop.stopped() {
1263            tokio::time::sleep(Duration::from_secs(1)).await;
1264            slept += Duration::from_secs(1);
1265        }
1266    }
1267}
1268
1269#[cfg(test)]
1270mod tests {
1271    use super::*;
1272    use crate::land::CheckView;
1273    use anyhow::Context;
1274    use std::sync::Mutex;
1275
1276    const URL: &str = "https://github.com/o/r/pull/7";
1277
1278    fn check(name: &str, v: Verdict, run: Option<&str>) -> CheckView {
1279        CheckView {
1280            name: name.to_owned(),
1281            verdict: v,
1282            run: run.map(str::to_owned),
1283            url: run.map(|r| format!("https://github.com/o/r/actions/runs/{r}/job/1")),
1284        }
1285    }
1286
1287    fn snap(state: PrLifecycle, head: &str, checks: Vec<CheckView>) -> RollupView {
1288        RollupView {
1289            url: URL.to_owned(),
1290            number: 7,
1291            state,
1292            head: head.to_owned(),
1293            checks,
1294        }
1295    }
1296
1297    fn open(checks: Vec<CheckView>) -> RollupView {
1298        snap(PrLifecycle::Open, "h1", checks)
1299    }
1300
1301    fn st_at(progress: i64) -> WatchState {
1302        WatchState {
1303            progress_at: progress,
1304            ..WatchState::default()
1305        }
1306    }
1307
1308    #[test]
1309    fn unreadable_is_unknown_and_merged_or_closed_is_done() {
1310        assert_eq!(decide(None, &st_at(0), 10, 60), Step::Unknown);
1311        for s in [PrLifecycle::Merged, PrLifecycle::Closed] {
1312            assert_eq!(
1313                decide(Some(&snap(s, "h", vec![])), &st_at(0), 10, 60),
1314                Step::Done
1315            );
1316        }
1317    }
1318
1319    #[test]
1320    fn a_failed_run_is_rerun_even_while_another_run_is_pending() {
1321        let s = open(vec![
1322            check("win", Verdict::Fail, Some("11")),
1323            check("lint", Verdict::Pending, Some("12")),
1324        ]);
1325        assert_eq!(
1326            decide(Some(&s), &st_at(0), 1, 3600),
1327            Step::Rerun(vec!["11".into()])
1328        );
1329    }
1330
1331    #[test]
1332    fn a_run_with_a_job_still_pending_is_not_complete() {
1333        let s = open(vec![
1334            check("win", Verdict::Fail, Some("11")),
1335            check("mac", Verdict::Pending, Some("11")),
1336        ]);
1337        assert_eq!(decide(Some(&s), &st_at(0), 1, 3600), Step::Wait);
1338    }
1339
1340    #[test]
1341    fn one_rerun_per_run_then_grace_then_escalation() {
1342        let s = open(vec![check("win", Verdict::Fail, Some("11"))]);
1343        let mut st = st_at(0);
1344        st.reruns.insert("11".into(), 100);
1345        assert_eq!(decide(Some(&s), &st, 100 + RERUN_GRACE - 1, 0), Step::Wait);
1346        assert_eq!(
1347            decide(Some(&s), &st, 100 + RERUN_GRACE, 0),
1348            Step::Escalate(Why::StillRed)
1349        );
1350    }
1351
1352    #[test]
1353    fn a_failure_with_no_run_id_goes_straight_to_a_human_once_settled() {
1354        let s = open(vec![check("ci/ext", Verdict::Fail, None)]);
1355        assert_eq!(
1356            decide(Some(&s), &st_at(0), 1, 0),
1357            Step::Escalate(Why::StillRed)
1358        );
1359        let s = open(vec![
1360            check("ci/ext", Verdict::Fail, None),
1361            check("x", Verdict::Pending, Some("5")),
1362        ]);
1363        assert_eq!(decide(Some(&s), &st_at(0), 1, 0), Step::Wait);
1364    }
1365
1366    #[test]
1367    fn no_progress_for_the_bounded_time_escalates_and_zero_disables_it() {
1368        let s = open(vec![check("a", Verdict::Pass, Some("1"))]);
1369        assert_eq!(decide(Some(&s), &st_at(0), 3599, 3600), Step::Wait);
1370        assert_eq!(
1371            decide(Some(&s), &st_at(0), 3600, 3600),
1372            Step::Escalate(Why::Stalled)
1373        );
1374        assert_eq!(decide(Some(&s), &st_at(0), 99_999, 0), Step::Wait);
1375    }
1376
1377    #[test]
1378    fn held_and_ignored_stay_quiet_until_the_fingerprint_moves() {
1379        let s = open(vec![check("a", Verdict::Fail, None)]);
1380        let mut st = st_at(0);
1381        st.held = Some(fingerprint(&s));
1382        assert_eq!(decide(Some(&s), &st, 99_999, 60), Step::Wait);
1383        let moved = open(vec![
1384            check("a", Verdict::Pass, None),
1385            check("b", Verdict::Fail, None),
1386        ]);
1387        assert_eq!(
1388            decide(Some(&moved), &st, 99_999, 60),
1389            Step::Escalate(Why::StillRed)
1390        );
1391        st.ignored = true;
1392        assert_eq!(decide(Some(&moved), &st, 99_999, 60), Step::Wait);
1393    }
1394
1395    #[test]
1396    fn observe_restarts_the_clock_only_on_change_and_forgets_reruns_on_a_new_head() {
1397        let mut st = st_at(0);
1398        let a = open(vec![check("a", Verdict::Pending, Some("1"))]);
1399        observe(&mut st, &a, 10);
1400        assert_eq!(st.progress_at, 10);
1401        observe(&mut st, &a, 50);
1402        assert_eq!(st.progress_at, 10);
1403        st.reruns.insert("1".into(), 5);
1404        observe(
1405            &mut st,
1406            &snap(PrLifecycle::Open, "h2", a.checks.clone()),
1407            60,
1408        );
1409        assert!(st.reruns.is_empty());
1410        assert_eq!(st.progress_at, 60);
1411    }
1412
1413    #[test]
1414    fn pr_key_names_owner_repo_and_number() {
1415        assert_eq!(pr_key(URL).as_deref(), Some("o/r#7"));
1416        assert_eq!(pr_key("https://github.com/o/r/issues/7"), None);
1417        assert_eq!(pr_key("nonsense"), None);
1418    }
1419
1420    #[test]
1421    fn notice_wording_does_not_vary_between_polls() {
1422        assert_eq!(rerun_message("o/r#7"), rerun_message("o/r#7"));
1423        assert!(
1424            !human_message("o/r#7")
1425                .chars()
1426                .any(|c| c.is_ascii_digit() && c != '7')
1427        );
1428    }
1429
1430    #[test]
1431    fn the_release_question_is_not_claimed_by_other_machinery() {
1432        let q = Question::new(
1433            String::new(),
1434            crate::bump::NOTICE_NODE.into(),
1435            "release-watch".into(),
1436            "s".into(),
1437            String::new(),
1438            vec![HOLD.into()],
1439        );
1440        // The watcher's question is served by a deputy, but still covers nothing.
1441        assert_eq!(
1442            crate::deputy::kind_of(&q),
1443            Some(crate::deputy::Kind::Release)
1444        );
1445        let n = Notice::warn("release-pr:o/r#7", "m").about([String::new()]);
1446        assert!(!notices::covers(&q, &n));
1447        // A choice-less notice on the same node (seat `bump`) gets no deputy.
1448        let mut plain = q.clone();
1449        plain.seat = "bump".into();
1450        assert_eq!(crate::deputy::kind_of(&plain), None);
1451    }
1452
1453    #[derive(Default)]
1454    struct Fake {
1455        snap: Mutex<Option<RollupView>>,
1456        reruns: Mutex<Vec<String>>,
1457        local: Mutex<bool>,
1458        info: Mutex<Option<PrInfo>>,
1459        merges: Mutex<Vec<String>>,
1460    }
1461
1462    impl ReleaseForge for std::sync::Arc<Fake> {
1463        fn list<'a>(&'a self, _: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
1464            Box::pin(async { Ok(vec![("chore/release-v1.0.0".to_owned(), URL.to_owned())]) })
1465        }
1466        fn snapshot<'a>(&'a self, _: &'a Path, _: &'a str) -> Fut<'a, Result<RollupView>> {
1467            let s = self.snap.lock().unwrap().clone();
1468            Box::pin(async move { s.context("unreadable") })
1469        }
1470        fn rerun<'a>(&'a self, _: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
1471            self.reruns.lock().unwrap().push(run.to_owned());
1472            Box::pin(async { Ok(()) })
1473        }
1474        fn config<'a>(&'a self, _: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
1475            let mut c = crate::config::Config::default();
1476            if *self.local.lock().unwrap() {
1477                c.release.mode = crate::config::ReleaseMode::Local;
1478                c.release.commands = vec!["true".to_owned()];
1479            }
1480            Box::pin(async move { Ok(c) })
1481        }
1482        fn info<'a>(&'a self, _: &'a Path, _: &'a str) -> Fut<'a, Result<PrInfo>> {
1483            let i = self.info.lock().unwrap().clone();
1484            Box::pin(async move { i.context("no info") })
1485        }
1486        fn merge<'a>(&'a self, _: &'a Path, _: &'a str, head: &'a str) -> Fut<'a, Result<bool>> {
1487            self.merges.lock().unwrap().push(head.to_owned());
1488            Box::pin(async { Ok(true) })
1489        }
1490    }
1491
1492    fn rig() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher) {
1493        let dir = tempfile::tempdir().unwrap();
1494        let fake = std::sync::Arc::new(Fake::default());
1495        let w = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
1496        (dir, fake, w)
1497    }
1498
1499    fn red() -> RollupView {
1500        open(vec![check("win", Verdict::Fail, Some("11"))])
1501    }
1502
1503    #[tokio::test]
1504    async fn reruns_once_survives_a_restart_then_escalates_with_one_question() {
1505        let (dir, fake, w) = rig();
1506        *fake.snap.lock().unwrap() = Some(red());
1507        let repo = PathBuf::from("/nowhere");
1508        let no = || false;
1509        w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
1510        assert_eq!(*fake.reruns.lock().unwrap(), vec!["11".to_owned()]);
1511
1512        // A restart is a new Watcher over the same home: no second rerun.
1513        let w2 = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
1514        w2.lap(
1515            std::slice::from_ref(&repo),
1516            3600,
1517            1000 + RERUN_GRACE - 1,
1518            &no,
1519        )
1520        .await;
1521        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1522        assert!(w2.questions().list().is_empty());
1523
1524        w2.lap(
1525            std::slice::from_ref(&repo),
1526            3600,
1527            1000 + RERUN_GRACE + 1,
1528            &no,
1529        )
1530        .await;
1531        w2.lap(
1532            std::slice::from_ref(&repo),
1533            3600,
1534            1000 + RERUN_GRACE + 400,
1535            &no,
1536        )
1537        .await;
1538        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1539        let qs = w2.questions().list();
1540        assert_eq!(qs.len(), 1, "asked once, not every lap");
1541        assert_eq!(qs[0].node, crate::bump::NOTICE_NODE);
1542        assert!(qs[0].detail.contains("win") && qs[0].detail.contains(URL));
1543        assert_eq!(qs[0].choices, vec![RERUN_AGAIN, HOLD, LEAVE_IT]);
1544        let ns = Notices::at(dir.path().join("notifications")).list();
1545        assert_eq!(ns.len(), 1);
1546        assert_eq!(ns[0].key, "release-pr:o/r#7");
1547    }
1548
1549    async fn escalated() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher, String) {
1550        let (dir, fake, w) = rig();
1551        *fake.snap.lock().unwrap() = Some(red());
1552        let repo = PathBuf::from("/nowhere");
1553        let no = || false;
1554        w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
1555        w.lap(&[repo], 3600, 2000, &no).await;
1556        let id = w.questions().list()[0].id.clone();
1557        (dir, fake, w, id)
1558    }
1559
1560    fn answer(w: &Watcher, id: &str, choice: &str) {
1561        w.questions()
1562            .update(id, |q| q.answer(Answer::Choice(choice.to_owned())))
1563            .unwrap();
1564    }
1565
1566    #[tokio::test]
1567    async fn rerun_again_reruns_exactly_once_more() {
1568        let (_d, fake, w, id) = escalated().await;
1569        answer(&w, &id, RERUN_AGAIN);
1570        let repo = PathBuf::from("/nowhere");
1571        let no = || false;
1572        w.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
1573        w.lap(&[repo], 3600, 3001, &no).await;
1574        assert_eq!(fake.reruns.lock().unwrap().len(), 2);
1575    }
1576
1577    #[tokio::test]
1578    async fn hold_and_silence_do_not_reask_and_leave_it_dismisses() {
1579        let (d, fake, w, id) = escalated().await;
1580        answer(&w, &id, HOLD);
1581        let repo = PathBuf::from("/nowhere");
1582        let no = || false;
1583        for t in [3000, 90_000, 200_000] {
1584            w.lap(std::slice::from_ref(&repo), 3600, t, &no).await;
1585        }
1586        assert_eq!(w.questions().list().len(), 1, "held: no second question");
1587        assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1588
1589        let (d2, _f2, w2, id2) = escalated().await;
1590        answer(&w2, &id2, LEAVE_IT);
1591        w2.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
1592        let ns = Notices::at(d2.path().join("notifications")).list();
1593        assert!(ns.is_empty(), "dismissed notices are hidden");
1594        drop(d);
1595    }
1596
1597    /// Give the question a deputy seat the way `Deputies::attach` + a turn do.
1598    fn with_deputy(w: &Watcher, id: &str, home: &Path) -> String {
1599        let seat = crate::agent::SeatState::new(&crate::ask::deputy_seat_key(id), "a", 1);
1600        let key = seat.key.clone();
1601        let brief = deputy_brief(&w.questions().get(id).unwrap(), home);
1602        w.questions()
1603            .update(id, |q| {
1604                let mut d = crate::ask::Deputy::new(brief);
1605                d.seat = Some(seat);
1606                q.deputy = Some(d);
1607                Ok(())
1608            })
1609            .unwrap();
1610        key
1611    }
1612
1613    #[tokio::test]
1614    async fn the_brief_names_the_pull_request_and_what_leave_it_really_does() {
1615        let (d, _f, w, id) = escalated().await;
1616        let q = w.questions().get(&id).unwrap();
1617        let b = deputy_brief(&q, d.path());
1618        assert!(b.contains(URL), "{b}");
1619        assert!(
1620            b.contains("`leave it`") && b.contains("does NOT close"),
1621            "{b}"
1622        );
1623        assert!(b.contains("`rerun again`") && b.contains("`hold`"), "{b}");
1624        assert!(!b.contains("`merge`"), "no merge on an escalation: {b}");
1625
1626        // No record naming the question: said, not invented.
1627        let empty = tempfile::tempdir().unwrap();
1628        let b = deputy_brief(&q, empty.path());
1629        assert!(b.contains("could not be found") && !b.contains(URL), "{b}");
1630    }
1631
1632    #[tokio::test]
1633    async fn a_settled_leave_it_is_applied_once_by_the_watcher_and_closes_nothing() {
1634        let (d, fake, w, id) = escalated().await;
1635        let seat = with_deputy(&w, &id, d.path());
1636        let say = "クローズしていいよ。private repo だから、何回やっても失敗しちゃうから";
1637        w.questions().update(&id, |q| q.say(say)).unwrap();
1638        w.questions()
1639            .update(&id, |q| {
1640                q.settle_by_deputy(&seat, LEAVE_IT, "クローズしていいよ")
1641            })
1642            .unwrap();
1643        let q = w.questions().get(&id).unwrap();
1644        assert_eq!(q.answer, Some(Answer::Choice(LEAVE_IT.to_owned())));
1645
1646        let repo = PathBuf::from("/nowhere");
1647        let no = || false;
1648        for t in [3000, 3100] {
1649            w.lap(std::slice::from_ref(&repo), 3600, t, &no).await;
1650        }
1651        let st = w.stored().pop().unwrap();
1652        assert!(st.ignored, "watching stopped");
1653        assert_eq!(st.applied, vec![id.clone()], "applied exactly once");
1654        assert!(
1655            Notices::at(d.path().join("notifications"))
1656                .list()
1657                .is_empty()
1658        );
1659        assert_eq!(fake.reruns.lock().unwrap().len(), 1, "no extra rerun");
1660        assert_eq!(w.questions().list().len(), 1, "no second question");
1661    }
1662
1663    #[tokio::test]
1664    async fn the_task_of_a_release_question_comes_from_the_watch_record() {
1665        let (d, _f, w, id) = escalated().await;
1666        let q = w.questions().get(&id).unwrap();
1667        let queue = crate::queue::Queue::at(d.path().join("queue"));
1668        assert!(
1669            task_of_question(d.path(), &q, &queue).is_none(),
1670            "no run known"
1671        );
1672        assert!(task_of_question(d.path(), &q, &queue).is_none());
1673        let mut st = w.stored().pop().unwrap();
1674        st.run = Some("run-1".into());
1675        w.save("o/r#7", &st);
1676        // A known run with no task is still nothing to refuse on.
1677        assert!(task_of_question(d.path(), &q, &queue).is_none());
1678        let mut t = crate::queue::Task::new(
1679            "t".into(),
1680            "i".into(),
1681            PathBuf::from("/nowhere"),
1682            crate::queue::Source::Human,
1683        );
1684        t.runs.push("run-1".into());
1685        queue.put(&mut t).unwrap();
1686        assert_eq!(task_of_question(d.path(), &q, &queue).unwrap().id, t.id);
1687    }
1688
1689    #[tokio::test]
1690    async fn merged_clears_the_notice_the_question_and_the_record() {
1691        let (d, fake, w, _id) = escalated().await;
1692        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1693        let repo = PathBuf::from("/nowhere");
1694        w.lap(&[repo], 3600, 5000, &(|| false)).await;
1695        assert!(
1696            Notices::at(d.path().join("notifications"))
1697                .list()
1698                .is_empty()
1699        );
1700        assert!(w.questions().list().iter().all(|q| !q.status.open()));
1701        assert!(w.stored().is_empty());
1702    }
1703
1704    #[tokio::test]
1705    async fn an_unreadable_forge_changes_nothing() {
1706        let (d, fake, w) = rig();
1707        *fake.snap.lock().unwrap() = None;
1708        w.lap(&[PathBuf::from("/nowhere")], 60, 99_999, &(|| false))
1709            .await;
1710        assert!(fake.reruns.lock().unwrap().is_empty());
1711        assert!(w.questions().list().is_empty());
1712        assert!(!d.path().join("release-watch").exists());
1713    }
1714
1715    #[tokio::test]
1716    async fn no_rerun_when_the_record_cannot_be_saved() {
1717        let (d, fake, w) = rig();
1718        *fake.snap.lock().unwrap() = Some(red());
1719        // A file where the directory should be makes every save fail.
1720        std::fs::write(d.path().join("release-watch"), "x").unwrap();
1721        w.lap(&[PathBuf::from("/nowhere")], 3600, 1000, &(|| false))
1722            .await;
1723        assert!(fake.reruns.lock().unwrap().is_empty());
1724    }
1725
1726    #[tokio::test]
1727    async fn a_park_stops_the_lap_before_any_call() {
1728        let (_d, fake, w) = rig();
1729        *fake.snap.lock().unwrap() = Some(red());
1730        w.lap(&[PathBuf::from("/nowhere")], 60, 1000, &(|| true))
1731            .await;
1732        assert!(fake.reruns.lock().unwrap().is_empty());
1733    }
1734
1735    // ---- [release] mode = "local" ----
1736
1737    fn local_st(asked: Option<&str>, held: Option<&str>) -> WatchState {
1738        WatchState {
1739            asked_head: asked.map(str::to_owned),
1740            held_head: held.map(str::to_owned),
1741            ..WatchState::default()
1742        }
1743    }
1744
1745    #[test]
1746    fn local_decisions_never_wait_on_ci_and_bind_approval_to_the_head() {
1747        use Approved::*;
1748        let none = local_st(None, None);
1749        assert_eq!(decide_local("h1", &none, NoQuestion), LocalStep::Ask);
1750        assert_eq!(decide_local("h1", &none, Open), LocalStep::Wait);
1751        let asked = local_st(Some("h1"), None);
1752        assert_eq!(
1753            decide_local("h1", &asked, Merge),
1754            LocalStep::Merge("h1".to_owned())
1755        );
1756        assert_eq!(
1757            decide_local("h1", &asked, Hold),
1758            LocalStep::Hold("h1".to_owned())
1759        );
1760        // The head moved after the owner answered: ask again, never merge.
1761        assert_eq!(decide_local("h2", &asked, Merge), LocalStep::Ask);
1762        assert_eq!(decide_local("h2", &asked, Hold), LocalStep::Ask);
1763        // A held head stays quiet until it moves; ignored stays quiet.
1764        let held = local_st(Some("h1"), Some("h1"));
1765        assert_eq!(decide_local("h1", &held, NoQuestion), LocalStep::Wait);
1766        assert_eq!(decide_local("h2", &held, NoQuestion), LocalStep::Ask);
1767        let mut ign = local_st(None, None);
1768        ign.ignored = true;
1769        assert_eq!(decide_local("h1", &ign, NoQuestion), LocalStep::Wait);
1770    }
1771
1772    #[test]
1773    fn pr_info_is_read_from_gh_json() {
1774        let i = parse_info(
1775            r#"{"headRefName":"chore/release-v1.2.3","title":"chore: release v1.2.3","mergeCommit":{"oid":"abc"}}"#,
1776        )
1777        .unwrap();
1778        assert_eq!(i.branch, "chore/release-v1.2.3");
1779        assert_eq!(i.merge_commit.as_deref(), Some("abc"));
1780        let open = parse_info(r#"{"headRefName":"b","title":"t","mergeCommit":null}"#).unwrap();
1781        assert_eq!(open.merge_commit, None);
1782    }
1783
1784    #[tokio::test]
1785    async fn a_local_release_asks_once_without_ci_reruns_or_stall_escalation() {
1786        let (_d, fake, w) = rig();
1787        *fake.local.lock().unwrap() = true;
1788        // Red checks that would rerun in Actions mode, and a stall clock far
1789        // past the limit: neither matters here.
1790        *fake.snap.lock().unwrap() = Some(red());
1791        let repo = PathBuf::from("/nowhere");
1792        let no = || false;
1793        w.lap(std::slice::from_ref(&repo), 60, 1_000_000, &no).await;
1794        w.lap(std::slice::from_ref(&repo), 60, 2_000_000, &no).await;
1795        assert!(fake.reruns.lock().unwrap().is_empty());
1796        let qs = w.questions().list();
1797        assert_eq!(qs.len(), 1, "one approval question, not one per lap");
1798        assert_eq!(qs[0].choices, vec![land::APPROVE, land::HOLD]);
1799        assert!(fake.merges.lock().unwrap().is_empty(), "silence is a hold");
1800    }
1801
1802    #[tokio::test]
1803    async fn merge_is_pinned_to_the_approved_head_and_a_moved_head_asks_again() {
1804        let (_d, fake, w) = rig();
1805        *fake.local.lock().unwrap() = true;
1806        *fake.snap.lock().unwrap() = Some(open(vec![]));
1807        let repo = PathBuf::from("/nowhere");
1808        let no = || false;
1809        w.lap(std::slice::from_ref(&repo), 60, 1, &no).await;
1810        let id = w.questions().list()[0].id.clone();
1811        // The head moves before the owner answers.
1812        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Open, "h2", vec![]));
1813        answer(&w, &id, land::APPROVE);
1814        w.lap(std::slice::from_ref(&repo), 60, 2, &no).await;
1815        assert!(fake.merges.lock().unwrap().is_empty());
1816        assert_eq!(w.questions().list().len(), 2, "asked again about h2");
1817        // Answering the new question merges exactly h2.
1818        let id2 = w
1819            .questions()
1820            .list()
1821            .into_iter()
1822            .find(|q| q.status.open())
1823            .unwrap()
1824            .id;
1825        answer(&w, &id2, land::APPROVE);
1826        w.lap(std::slice::from_ref(&repo), 60, 3, &no).await;
1827        assert_eq!(*fake.merges.lock().unwrap(), vec!["h2".to_owned()]);
1828    }
1829
1830    #[tokio::test]
1831    async fn a_failed_release_holds_with_one_notice_and_one_question_and_never_retries() {
1832        let (d, fake, w) = rig();
1833        *fake.local.lock().unwrap() = true;
1834        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1835        *fake.info.lock().unwrap() = Some(PrInfo {
1836            branch: "chore/release-v1.0.0".to_owned(),
1837            merge_commit: Some("deadbeef".to_owned()),
1838            title: "chore: release v1.0.0".to_owned(),
1839        });
1840        // The checkout is not a repository, so the release cannot start.
1841        let repo = d.path().join("nowhere");
1842        let no = || false;
1843        w.lap(std::slice::from_ref(&repo), 60, 1, &no).await;
1844        let st = w.load("o/r#7");
1845        let job = st.job.expect("the job is recorded");
1846        assert!(job.failed.is_some() && !job.tag_done && !job.finished);
1847        let qs = w.questions().list();
1848        assert_eq!(qs.len(), 1);
1849        assert_eq!(qs[0].choices, vec![RETRY, LEAVE_IT]);
1850        // Later laps change nothing: no retry, no second question.
1851        w.lap(std::slice::from_ref(&repo), 60, 2, &no).await;
1852        assert_eq!(w.questions().list().len(), 1);
1853        assert_eq!(Notices::at(d.path().join("notifications")).list().len(), 1);
1854    }
1855
1856    #[tokio::test]
1857    async fn a_failed_release_holds_the_task_whose_run_opened_the_pr() {
1858        let (d, fake, w) = rig();
1859        *fake.local.lock().unwrap() = true;
1860        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1861        *fake.info.lock().unwrap() = Some(PrInfo {
1862            branch: "chore/release-v1.0.0".to_owned(),
1863            merge_commit: Some("deadbeef".to_owned()),
1864            title: "t".to_owned(),
1865        });
1866        let repo = d.path().join("nowhere");
1867        let q = crate::queue::Queue::at(d.path().join("queue"));
1868        let mut t = crate::queue::Task::new(
1869            "t".to_owned(),
1870            "i".to_owned(),
1871            repo.clone(),
1872            crate::queue::Source::Human,
1873        );
1874        t.runs.push("run1".to_owned());
1875        t.status = crate::queue::TaskStatus::Done;
1876        q.put(&mut t).unwrap();
1877        register(d.path(), &repo, URL, "run1");
1878        w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1879        let held = q.get(&t.id).unwrap();
1880        assert_eq!(held.status, crate::queue::TaskStatus::Held);
1881        assert!(
1882            held.hold_reason
1883                .unwrap_or_default()
1884                .contains("release v1.0.0")
1885        );
1886    }
1887
1888    #[tokio::test]
1889    async fn unhold_gives_the_task_back_only_when_the_hold_is_ours() {
1890        let (d, _fake, w) = rig();
1891        let q = crate::queue::Queue::at(d.path().join("queue"));
1892        let mut t = crate::queue::Task::new(
1893            "t".to_owned(),
1894            "i".to_owned(),
1895            PathBuf::from("/r"),
1896            crate::queue::Source::Human,
1897        );
1898        t.runs.push("run1".to_owned());
1899        t.status = crate::queue::TaskStatus::Done;
1900        q.put(&mut t).unwrap();
1901        let mut st = WatchState {
1902            run: Some("run1".to_owned()),
1903            ..WatchState::default()
1904        };
1905        w.hold_task(&mut st, "failed");
1906        assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1907        // A restart reconciling again does not hold twice.
1908        w.hold_task(&mut st, "failed again");
1909        assert!(w.unhold_task(&mut st));
1910        assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Done);
1911        assert!(st.held_task.is_none());
1912
1913        // The owner re-held it by hand with their own reason: left alone.
1914        let mut h = q.get(&t.id).unwrap();
1915        st.held_task = Some(h.id.clone());
1916        h.hold_manual(Some("mine".to_owned()));
1917        q.put(&mut h).unwrap();
1918        assert!(w.unhold_task(&mut st));
1919        assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1920    }
1921
1922    #[tokio::test]
1923    async fn a_restart_after_the_failure_was_saved_still_holds_the_task() {
1924        let (d, fake, w) = rig();
1925        *fake.local.lock().unwrap() = true;
1926        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1927        let repo = d.path().join("nowhere");
1928        let q = crate::queue::Queue::at(d.path().join("queue"));
1929        let mut t = crate::queue::Task::new(
1930            "t".to_owned(),
1931            "i".to_owned(),
1932            repo.clone(),
1933            crate::queue::Source::Human,
1934        );
1935        t.runs.push("run1".to_owned());
1936        t.status = crate::queue::TaskStatus::Done;
1937        q.put(&mut t).unwrap();
1938        // The state a crash leaves: failure persisted, task never held.
1939        let mut job = Job::new("1.0.0", URL, "deadbeef");
1940        job.failed = Some("command 1 exited 3".to_owned());
1941        let st = WatchState {
1942            repo: repo.to_string_lossy().into_owned(),
1943            url: URL.to_owned(),
1944            run: Some("run1".to_owned()),
1945            job: Some(job),
1946            ..WatchState::default()
1947        };
1948        w.save("o/r#7", &st);
1949        w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1950        assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1951    }
1952
1953    #[tokio::test]
1954    async fn a_restart_after_the_job_finished_still_gives_the_task_back() {
1955        let (d, fake, w) = rig();
1956        *fake.local.lock().unwrap() = true;
1957        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1958        let repo = d.path().join("nowhere");
1959        let q = crate::queue::Queue::at(d.path().join("queue"));
1960        let mut t = crate::queue::Task::new(
1961            "t".to_owned(),
1962            "i".to_owned(),
1963            repo.clone(),
1964            crate::queue::Source::Human,
1965        );
1966        t.runs.push("run1".to_owned());
1967        t.hold_machine(Some(format!("{HOLD_PREFIX}release v1.0.0 failed")));
1968        q.put(&mut t).unwrap();
1969        // The state a crash leaves: finished saved, task never given back.
1970        let mut job = Job::new("1.0.0", URL, "deadbeef");
1971        job.finished = true;
1972        let st = WatchState {
1973            repo: repo.to_string_lossy().into_owned(),
1974            url: URL.to_owned(),
1975            run: Some("run1".to_owned()),
1976            job: Some(job),
1977            held_task: Some(t.id.clone()),
1978            ..WatchState::default()
1979        };
1980        w.save("o/r#7", &st);
1981        w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1982        assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Done);
1983        assert!(!w.state_path("o/r#7").exists());
1984    }
1985
1986    #[tokio::test]
1987    async fn a_pull_request_registered_at_open_is_picked_up_even_if_unlisted() {
1988        let (d, _fake, w) = rig();
1989        register(d.path(), Path::new("/r"), URL, "run1");
1990        let st = w.load("o/r#7");
1991        assert_eq!(st.url, URL);
1992        assert_eq!(st.repo, "/r");
1993        // Never overwrites a live record.
1994        let mut live = st.clone();
1995        live.head = "keep".to_owned();
1996        w.save("o/r#7", &live);
1997        register(d.path(), Path::new("/r"), URL, "run1");
1998        assert_eq!(w.load("o/r#7").head, "keep");
1999    }
2000
2001    #[tokio::test]
2002    async fn a_finished_release_is_kept_on_the_run_and_a_missing_run_still_completes() {
2003        let (d, fake, w) = rig();
2004        *fake.local.lock().unwrap() = true;
2005        *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
2006        let repo = d.path().join("nowhere");
2007        let mut run = RunState::new(
2008            repo.clone(),
2009            "main".to_owned(),
2010            "0123456789abcdef".to_owned(),
2011            "t".to_owned(),
2012            crate::config::Config::default(),
2013        );
2014        run.id = "run1".to_owned();
2015        run.save_under(d.path()).unwrap();
2016        let mut job = Job::new("1.0.0", URL, "deadbeef");
2017        job.finished = true;
2018        job.log.push(release_local::StepLog {
2019            name: "cmd".to_owned(),
2020            code: Some(0),
2021            tail: "ok".to_owned(),
2022            output: None,
2023        });
2024        let mk = |run: &str| WatchState {
2025            repo: repo.to_string_lossy().into_owned(),
2026            url: URL.to_owned(),
2027            run: Some(run.to_owned()),
2028            job: Some(job.clone()),
2029            ..WatchState::default()
2030        };
2031        w.save("o/r#7", &mk("run1"));
2032        w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
2033        assert!(!w.state_path("o/r#7").exists());
2034        let kept = RunState::load_under("run1", d.path()).unwrap();
2035        assert_eq!(kept.release_bump.unwrap().release, Some(job.clone()));
2036
2037        // A run that is gone must not keep the watcher record forever.
2038        w.save("o/r#7", &mk("gone"));
2039        w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
2040        assert!(!w.state_path("o/r#7").exists());
2041    }
2042}