1use 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
45const LAP: Duration = Duration::from_secs(300);
47
48const RERUN_GRACE: i64 = 180;
52
53pub const RERUN_AGAIN: &str = "rerun again";
56pub const HOLD: &str = "hold";
58pub const LEAVE_IT: &str = "leave it";
60
61#[derive(Debug, Clone, PartialEq, Eq)]
63pub(crate) enum Step {
64 Done,
66 Wait,
68 Unknown,
70 Rerun(Vec<String>),
72 Escalate(Why),
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77pub(crate) enum Why {
78 StillRed,
80 Stalled,
82}
83
84#[derive(Debug, Clone, Default, Serialize, Deserialize)]
86#[serde(default)]
87pub(crate) struct WatchState {
88 pub repo: String,
90 pub url: String,
91 pub head: String,
93 pub reruns: BTreeMap<String, i64>,
95 pub fingerprint: String,
97 pub progress_at: i64,
98 pub question: Option<String>,
100 pub held: Option<String>,
102 pub ignored: bool,
104 pub applied: Vec<String>,
106}
107
108fn is_failed(v: Verdict) -> bool {
109 v == Verdict::Fail
110}
111
112pub(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
123pub(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
137pub(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 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
198pub(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
220fn 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
265pub(crate) trait ReleaseForge: Send + Sync {
267 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
273pub(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
315pub(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 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(¬ice_key(pr), message));
377 }
378
379 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 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 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 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 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 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 for r in &runs {
479 st.reruns.insert(r.clone(), now);
480 }
481 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 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(¬ices::id_of(¬ice_key(pr)))
539 {
540 tracing::warn!("could not dismiss the notice for {pr}: {e:#}");
541 }
542 }
543 _ => st.held = Some(fingerprint(snap)),
545 }
546 }
547
548 fn finish(&self, pr: &str, st: &WatchState) {
550 let _ =
551 Notices::at(self.home.join("notifications")).dismiss(¬ices::id_of(¬ice_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
562pub(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 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 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}