1use 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
54const LAP: Duration = Duration::from_secs(300);
56
57const RERUN_GRACE: i64 = 180;
61
62pub const RERUN_AGAIN: &str = "rerun again";
65pub const HOLD: &str = "hold";
67pub const LEAVE_IT: &str = "leave it";
69pub const RETRY: &str = "retry";
71
72#[derive(Debug, Clone, PartialEq, Eq)]
74pub(crate) enum Step {
75 Done,
77 Wait,
79 Unknown,
81 Rerun(Vec<String>),
83 Escalate(Why),
85}
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
88pub(crate) enum Why {
89 StillRed,
91 Stalled,
93}
94
95#[derive(Debug, Clone, Default, Serialize, Deserialize)]
97#[serde(default)]
98pub(crate) struct WatchState {
99 pub repo: String,
101 pub url: String,
102 pub head: String,
104 pub reruns: BTreeMap<String, i64>,
106 pub fingerprint: String,
108 pub progress_at: i64,
109 pub question: Option<String>,
111 pub held: Option<String>,
113 pub ignored: bool,
115 pub applied: Vec<String>,
117 pub asked_head: Option<String>,
119 pub held_head: Option<String>,
122 pub run: Option<String>,
125 pub held_task: Option<String>,
128 pub job: Option<Job>,
130}
131
132fn is_failed(v: Verdict) -> bool {
133 v == Verdict::Fail
134}
135
136pub(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
147pub(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
161pub(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 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
222pub(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
244fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
289pub(crate) enum Approved {
290 NoQuestion,
292 Open,
294 Merge,
296 Hold,
298}
299
300#[derive(Debug, Clone, PartialEq, Eq)]
302pub(crate) enum LocalStep {
303 Wait,
305 Ask,
307 Merge(String),
309 Hold(String),
311}
312
313pub(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 Approved::Merge | Approved::Hold => LocalStep::Ask,
328 }
329}
330
331pub(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
350const HOLD_PREFIX: &str = "[release] ";
353
354fn approval_message(pr: &str) -> String {
355 format!("Release PR {pr} is waiting for your approval to merge")
356}
357
358fn failed_message(pr: &str) -> String {
359 format!("Release PR {pr} merged but the release is on hold")
360}
361
362fn released_key(pr: &str) -> String {
363 format!("release-done:{pr}")
364}
365
366type Fut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
367
368pub(crate) trait ReleaseForge: Send + Sync {
370 fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>>;
372 fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>>;
373 fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>>;
374 fn config<'a>(&'a self, _repo: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
376 Box::pin(async { Ok(crate::config::Config::default()) })
377 }
378 fn info<'a>(&'a self, _repo: &'a Path, url: &'a str) -> Fut<'a, Result<PrInfo>> {
380 Box::pin(async move { bail!("no pull request info for {url}") })
381 }
382 fn merge<'a>(&'a self, _repo: &'a Path, url: &'a str, _head: &'a str) -> Fut<'a, Result<bool>> {
385 Box::pin(async move { bail!("cannot merge {url}") })
386 }
387}
388
389#[derive(Debug, Clone, Default, PartialEq, Eq)]
391pub(crate) struct PrInfo {
392 pub branch: String,
394 pub merge_commit: Option<String>,
396 pub title: String,
398}
399
400pub(crate) fn parse_info(json: &str) -> Result<PrInfo> {
402 let v: serde_json::Value = serde_json::from_str(json)?;
403 let text = |k: &str| {
404 v.get(k)
405 .and_then(|x| x.as_str())
406 .unwrap_or_default()
407 .to_owned()
408 };
409 Ok(PrInfo {
410 branch: text("headRefName"),
411 title: text("title"),
412 merge_commit: v
413 .get("mergeCommit")
414 .and_then(|m| m.get("oid"))
415 .and_then(|o| o.as_str())
416 .filter(|o| !o.is_empty())
417 .map(str::to_owned),
418 })
419}
420
421pub(crate) struct GhForge;
423
424impl ReleaseForge for GhForge {
425 fn list<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
426 Box::pin(crate::bump::list_open_release_prs(repo))
427 }
428
429 fn snapshot<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<RollupView>> {
430 Box::pin(async move {
431 let args = [
432 "pr".to_owned(),
433 "view".to_owned(),
434 url.to_owned(),
435 "--json".to_owned(),
436 "url,number,state,headRefOid,statusCheckRollup".to_owned(),
437 ];
438 let (ok, out) = land::gh(repo, &args).await?;
439 if !ok {
440 bail!("gh pr view {url}: {out}");
441 }
442 land::parse_rollup(&out)
443 })
444 }
445
446 fn rerun<'a>(&'a self, repo: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
447 Box::pin(async move {
448 let args = [
449 "run".to_owned(),
450 "rerun".to_owned(),
451 run.to_owned(),
452 "--failed".to_owned(),
453 ];
454 let (ok, out) = land::gh(repo, &args).await?;
455 if !ok {
456 bail!("gh run rerun {run}: {out}");
457 }
458 Ok(())
459 })
460 }
461
462 fn config<'a>(&'a self, repo: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
463 Box::pin(async move { Ok(crate::config::Config::discover(repo, None)?.0) })
464 }
465
466 fn info<'a>(&'a self, repo: &'a Path, url: &'a str) -> Fut<'a, Result<PrInfo>> {
467 Box::pin(async move {
468 let args = [
469 "pr".to_owned(),
470 "view".to_owned(),
471 url.to_owned(),
472 "--json".to_owned(),
473 "headRefName,mergeCommit,title".to_owned(),
474 ];
475 let (ok, out) = land::gh(repo, &args).await?;
476 if !ok {
477 bail!("gh pr view {url}: {out}");
478 }
479 parse_info(&out)
480 })
481 }
482
483 fn merge<'a>(&'a self, repo: &'a Path, url: &'a str, head: &'a str) -> Fut<'a, Result<bool>> {
484 Box::pin(async move {
485 let title = self
486 .info(repo, url)
487 .await
488 .map(|i| i.title)
489 .unwrap_or_default();
490 let title = if title.trim().is_empty() {
491 "chore: release".to_owned()
492 } else {
493 title
494 };
495 let mut argv = crate::bump::bump_merge_argv(url, &title);
498 argv.push("--match-head-commit".to_owned());
499 argv.push(head.to_owned());
500 let (ok, out) = land::gh(repo, &argv).await?;
501 if ok {
502 return Ok(true);
503 }
504 let after = land::lifecycle(repo, url).await.ok();
507 if land::merged_after_all(&argv, &out, after).is_some() {
508 return Ok(true);
509 }
510 bail!("gh pr merge {url}: {out}")
511 })
512 }
513}
514
515pub(crate) struct Watcher {
517 forge: Box<dyn ReleaseForge>,
518 home: PathBuf,
519}
520
521impl Watcher {
522 pub(crate) fn new(forge: Box<dyn ReleaseForge>, home: PathBuf) -> Self {
523 Self { forge, home }
524 }
525
526 fn dir(&self) -> PathBuf {
527 self.home.join("release-watch")
528 }
529
530 fn state_path(&self, pr: &str) -> PathBuf {
531 self.dir().join(format!("{}.json", notices::id_of(pr)))
532 }
533
534 fn load(&self, pr: &str) -> WatchState {
535 std::fs::read_to_string(self.state_path(pr))
536 .ok()
537 .and_then(|s| serde_json::from_str(&s).ok())
538 .unwrap_or_default()
539 }
540
541 fn save(&self, pr: &str, st: &WatchState) -> bool {
542 let write = || -> Result<()> {
543 std::fs::create_dir_all(self.dir())?;
544 let path = self.state_path(pr);
545 let tmp = path.with_extension("json.tmp");
546 std::fs::write(&tmp, serde_json::to_string_pretty(st)?)?;
547 std::fs::rename(&tmp, &path)?;
548 Ok(())
549 };
550 match write() {
551 Ok(()) => true,
552 Err(e) => {
553 tracing::warn!("could not save the release watch for {pr}: {e:#}");
554 false
555 }
556 }
557 }
558
559 fn stored(&self) -> Vec<WatchState> {
561 let Ok(rd) = std::fs::read_dir(self.dir()) else {
562 return Vec::new();
563 };
564 rd.flatten()
565 .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
566 .filter_map(|e| std::fs::read_to_string(e.path()).ok())
567 .filter_map(|s| serde_json::from_str(&s).ok())
568 .collect()
569 }
570
571 fn questions(&self) -> Questions {
572 Questions::at(self.home.join("questions"))
573 }
574
575 fn raise(&self, pr: &str, message: String) {
576 notices::raise_in(&self.home, Notice::warn(¬ice_key(pr), message));
577 }
578
579 pub(crate) async fn lap(
582 &self,
583 repos: &[PathBuf],
584 stall_secs: i64,
585 now: i64,
586 halt: &(dyn Fn() -> bool + Sync),
587 ) {
588 let stored = self.stored();
589 let mut repos = repos.to_vec();
591 for s in &stored {
592 let p = PathBuf::from(&s.repo);
593 if !s.repo.is_empty() && !repos.contains(&p) {
594 repos.push(p);
595 }
596 }
597 for repo in &repos {
598 if halt() {
599 return;
600 }
601 let mut urls: Vec<String> = match self.forge.list(repo).await {
602 Ok(v) => v.into_iter().map(|(_, u)| u).collect(),
603 Err(e) => {
604 tracing::warn!(
605 "could not list release pull requests in {}: {e:#}",
606 repo.display()
607 );
608 Vec::new()
609 }
610 };
611 let here = repo.to_string_lossy();
614 for s in stored.iter().filter(|s| s.repo == here.as_ref()) {
615 if !urls.contains(&s.url) {
616 urls.push(s.url.clone());
617 }
618 }
619 for url in urls {
620 if halt() {
621 return;
622 }
623 self.watch(repo, &url, stall_secs, now, halt).await;
624 }
625 }
626 }
627
628 pub(crate) async fn watch(
630 &self,
631 repo: &Path,
632 url: &str,
633 stall_secs: i64,
634 now: i64,
635 halt: &(dyn Fn() -> bool + Sync),
636 ) {
637 let Some(pr) = pr_key(url) else {
638 tracing::warn!("not a pull request url: {url}");
639 return;
640 };
641 let mut st = self.load(&pr);
642 st.repo = repo.to_string_lossy().into_owned();
643 st.url = url.to_owned();
644
645 let snap = match self.forge.snapshot(repo, url).await {
646 Ok(s) => Some(s),
647 Err(e) => {
648 tracing::warn!("could not read {url}: {e:#}");
649 None
650 }
651 };
652 match self.forge.config(repo).await {
653 Ok(cfg) if cfg.release.is_local() => {
654 self.watch_local(repo, &pr, st, snap, &cfg).await;
655 return;
656 }
657 Ok(_) => {}
658 Err(e) => tracing::warn!(
659 "could not read the config of {}: {e:#}; watching {url} as an Actions release",
660 repo.display()
661 ),
662 }
663 if decide(snap.as_ref(), &st, now, stall_secs) == Step::Done {
664 self.finish(&pr, &st);
665 return;
666 }
667 let Some(snap) = snap else {
668 return;
669 };
670
671 if let Some(id) = st.question.clone() {
673 match self.questions().get(&id) {
674 Ok(q) if q.status.open() => return,
675 Ok(q) => self.apply_answer(&pr, &mut st, &snap, &q),
676 Err(_) => st.question = None,
677 }
678 observe(&mut st, &snap, now);
680 self.save(&pr, &st);
681 } else {
682 observe(&mut st, &snap, now);
683 }
684
685 match decide(Some(&snap), &st, now, stall_secs) {
686 Step::Done | Step::Unknown | Step::Wait => {}
687 Step::Rerun(runs) => {
688 for r in &runs {
690 st.reruns.insert(r.clone(), now);
691 }
692 if !self.save(&pr, &st) || halt() {
695 return;
696 }
697 let mut ok = false;
698 for r in &runs {
699 match self.forge.rerun(repo, r).await {
700 Ok(()) => ok = true,
701 Err(e) => tracing::warn!("could not rerun run {r} of {url}: {e:#}"),
702 }
703 }
704 if ok {
705 self.raise(&pr, rerun_message(&pr));
706 }
707 }
708 Step::Escalate(why) => {
709 self.raise(&pr, human_message(&pr));
710 let mut q = Question::new(
711 String::new(),
712 crate::bump::NOTICE_NODE.to_owned(),
713 "release-watch".to_owned(),
714 format!("Release PR {pr} is stuck: what now?"),
715 question_detail(&snap, why, &st),
716 vec![RERUN_AGAIN.to_owned(), HOLD.to_owned(), LEAVE_IT.to_owned()],
717 );
718 match self.questions().put(&mut q) {
719 Ok(()) => st.question = Some(q.id.clone()),
720 Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
721 }
722 }
723 }
724 self.save(&pr, &st);
725 }
726
727 async fn watch_local(
729 &self,
730 repo: &Path,
731 pr: &str,
732 mut st: WatchState,
733 snap: Option<RollupView>,
734 cfg: &crate::config::Config,
735 ) {
736 let Some(snap) = snap else {
737 return;
738 };
739 match snap.state {
740 PrLifecycle::Closed => self.finish(pr, &st),
741 PrLifecycle::Merged => self.release_merged(repo, pr, st, cfg).await,
742 PrLifecycle::Open => {
743 let approved = match st.question.clone() {
744 None => Approved::NoQuestion,
745 Some(id) => match self.questions().get(&id) {
746 Ok(q) if q.status.open() => Approved::Open,
747 Ok(q) => match &q.answer {
748 Some(Answer::Choice(c))
749 if land::approval(Some(c)) == land::Approval::Merge =>
750 {
751 Approved::Merge
752 }
753 _ => Approved::Hold,
754 },
755 Err(_) => Approved::NoQuestion,
756 },
757 };
758 let step = decide_local(&snap.head, &st, approved);
759 if matches!(
760 step,
761 LocalStep::Merge(_) | LocalStep::Hold(_) | LocalStep::Ask
762 ) {
763 st.question = None;
765 }
766 match step {
767 LocalStep::Wait => {}
768 LocalStep::Hold(head) => st.held_head = Some(head),
769 LocalStep::Ask => {
770 st.held_head = None;
771 st.asked_head = Some(snap.head.clone());
772 self.raise(pr, approval_message(pr));
773 let mut q = Question::new(
774 String::new(),
775 crate::bump::NOTICE_NODE.to_owned(),
776 "release-watch".to_owned(),
777 format!("Release PR {pr}: merge it and release?"),
778 format!(
779 "Pull request: {}\nHead: {}\n\n\
780 `[release] mode = \"local\"`: no CI is awaited. `{}` merges \
781 exactly this head, then magi tags the merge commit and runs the \
782 configured release commands. `{}` leaves the pull request \
783 open. Silence is a hold.\n",
784 snap.url,
785 snap.head,
786 land::APPROVE,
787 land::HOLD
788 ),
789 vec![land::APPROVE.to_owned(), land::HOLD.to_owned()],
790 );
791 match self.questions().put(&mut q) {
792 Ok(()) => st.question = Some(q.id.clone()),
793 Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
794 }
795 }
796 LocalStep::Merge(head) => {
797 match self.forge.merge(repo, &snap.url, &head).await {
798 Ok(_) => {
799 self.save(pr, &st);
801 self.release_merged(repo, pr, st, cfg).await;
802 return;
803 }
804 Err(e) => {
805 tracing::warn!("could not merge {}: {e:#}", snap.url);
806 st.held_head = Some(head);
807 self.raise(
808 pr,
809 format!(
810 "Release PR {pr} could not be merged; merge it by hand"
811 ),
812 );
813 }
814 }
815 }
816 }
817 self.save(pr, &st);
818 }
819 }
820 }
821
822 async fn release_merged(
824 &self,
825 repo: &Path,
826 pr: &str,
827 mut st: WatchState,
828 cfg: &crate::config::Config,
829 ) {
830 if st.ignored {
831 return self.finish(pr, &st);
832 }
833 if st.job.is_none() {
834 let info = match self.forge.info(repo, &st.url).await {
835 Ok(i) => i,
836 Err(e) => {
837 tracing::warn!("could not read the merged {}: {e:#}", st.url);
838 return;
839 }
840 };
841 let (Some(version), Some(commit)) = (
842 release_local::version_from_branch(&info.branch),
843 info.merge_commit,
844 ) else {
845 if release_local::version_from_branch(&info.branch).is_none() {
848 self.finish(pr, &st);
849 }
850 return;
851 };
852 st.job = Some(Job::new(&version, &st.url, &commit));
853 if !self.save(pr, &st) {
854 return;
855 }
856 }
857 let Some(mut job) = st.job.take() else {
858 return;
859 };
860 if job.finished {
861 st.job = Some(job);
864 if !self.record_on_run(&st) {
865 self.save(pr, &st);
866 return;
867 }
868 return self.complete(pr, &mut st);
869 }
870 if job.interrupted() {
871 job.failed = Some(format!(
872 "magi stopped while `{}` was running; it may have partly run, so it is not repeated on its own",
873 job.running.clone().unwrap_or_default()
874 ));
875 }
876 if let Some(why) = job.failed.clone() {
877 self.hold_task(&mut st, &format!("release {} is on hold: {why}", job.tag()));
880 let retry = self.settle_job_question(pr, &mut st, &mut job);
882 if !retry {
883 st.job = Some(job);
884 self.save(pr, &st);
885 return;
886 }
887 }
888 let shell = cfg.shell();
889 let env = release_local::Env {
890 repo,
891 home: &self.home,
892 key: pr,
893 remote: &cfg.merge.remote,
894 shell: &shell,
895 release: &cfg.release,
896 };
897 let result = {
898 let mut save = |j: &Job| {
899 let mut copy = st.clone();
900 copy.job = Some(j.clone());
901 self.save(pr, ©)
902 };
903 release_local::run_job(&env, &mut job, &mut save).await
904 };
905 match result {
906 Ok(()) => {
907 self.raise_released(pr, &job);
908 st.job = Some(job);
909 if self.record_on_run(&st) {
910 self.complete(pr, &mut st);
911 } else {
912 self.save(pr, &st);
915 }
916 }
917 Err(e) => {
918 job.failed = Some(format!("{e:#}"));
919 self.hold_task(&mut st, &format!("release {} failed: {e:#}", job.tag()));
920 self.raise(pr, failed_message(pr));
921 self.file_job_question(pr, &mut st, &job);
922 st.job = Some(job);
923 self.save(pr, &st);
924 self.record_on_run(&st);
925 }
926 }
927 }
928
929 fn record_on_run(&self, st: &WatchState) -> bool {
934 let (Some(run), Some(job)) = (st.run.as_deref(), st.job.as_ref()) else {
935 return true;
936 };
937 let mut state = match RunState::load_under(run, &self.home) {
938 Ok(s) => s,
939 Err(e) => {
940 tracing::warn!("could not record the release on run {run}: {e:#}");
941 return true;
942 }
943 };
944 let bump = state.release_bump.get_or_insert_with(Default::default);
945 if bump.release.as_ref() == Some(job) {
946 return true;
947 }
948 bump.release = Some(job.clone());
949 state.event(
950 "release",
951 format!(
952 "release {} {}",
953 job.tag(),
954 if job.finished {
955 "finished"
956 } else if let Some(why) = &job.failed {
957 why
958 } else {
959 "stopped"
960 }
961 ),
962 );
963 match state.save_under(&self.home) {
964 Ok(()) => true,
965 Err(e) => {
966 tracing::warn!("could not save the release on run {run}: {e:#}");
967 false
968 }
969 }
970 }
971
972 fn hold_task(&self, st: &mut WatchState, reason: &str) {
976 if st.held_task.is_some() {
977 return;
978 }
979 let Some(run) = st.run.clone() else {
980 return;
981 };
982 let q = crate::queue::Queue::at(self.home.join("queue"));
983 for mut t in q.list() {
984 if t.runs.contains(&run) {
985 if t.status != crate::queue::TaskStatus::Held {
986 t.hold_machine(Some(format!("{HOLD_PREFIX}{reason}")));
987 match q.put(&mut t) {
988 Ok(()) => st.held_task = Some(t.id.clone()),
989 Err(e) => tracing::warn!("could not hold task {} for {run}: {e:#}", t.id),
990 }
991 }
992 return;
993 }
994 }
995 }
996
997 #[must_use]
1004 fn unhold_task(&self, st: &mut WatchState) -> bool {
1005 let Some(id) = st.held_task.take() else {
1006 return true;
1007 };
1008 let q = crate::queue::Queue::at(self.home.join("queue"));
1009 let mut t = match q.get(&id) {
1010 Ok(t) => t,
1011 Err(_) if matches!(q.path_of(&id).try_exists(), Ok(false)) => return true,
1014 Err(e) => {
1015 tracing::warn!("could not read task {id} to restore it after the release: {e:#}");
1016 st.held_task = Some(id);
1017 return false;
1018 }
1019 };
1020 let ours = t.status == crate::queue::TaskStatus::Held
1021 && t.hold_source == Some(crate::queue::HoldSource::Machine)
1022 && t.hold_reason
1023 .as_deref()
1024 .is_some_and(|r| r.starts_with(HOLD_PREFIX));
1025 if ours {
1026 t.succeed();
1027 if let Err(e) = q.put(&mut t) {
1028 tracing::warn!("could not restore task {id} after the release: {e:#}");
1029 st.held_task = Some(id);
1030 return false;
1031 }
1032 }
1033 true
1034 }
1035
1036 fn complete(&self, pr: &str, st: &mut WatchState) {
1040 if self.unhold_task(st) {
1041 self.finish(pr, st);
1042 } else {
1043 self.save(pr, st);
1044 }
1045 }
1046
1047 fn raise_released(&self, pr: &str, job: &Job) {
1048 notices::raise_in(
1049 &self.home,
1050 Notice::info(&released_key(pr), format!("Released {} ({pr})", job.tag())),
1051 );
1052 }
1053
1054 fn file_job_question(&self, pr: &str, st: &mut WatchState, job: &Job) {
1056 let last = job
1057 .log
1058 .last()
1059 .map(|l| format!("\nLast step: {}\n\n{}\n", l.name, l.tail))
1060 .unwrap_or_default();
1061 let mut q = Question::new(
1062 String::new(),
1063 crate::bump::NOTICE_NODE.to_owned(),
1064 "release-watch".to_owned(),
1065 format!("Release {} of {pr} is on hold: retry?", job.tag()),
1066 format!(
1067 "Pull request: {}\nWhy it stopped: {}\n{last}\n\
1068 Full output is under the magi home's release-local directory.\n\n\
1069 - `{RETRY}`: run it again from where it stopped (the tag is not \
1070 recreated and finished commands are skipped).\n\
1071 - `{LEAVE_IT}`: stop watching; release by hand.\n\n\
1072 Nothing is retried on its own; silence is a hold.\n",
1073 job.pr_url,
1074 job.failed.clone().unwrap_or_default()
1075 ),
1076 vec![RETRY.to_owned(), LEAVE_IT.to_owned()],
1077 );
1078 match self.questions().put(&mut q) {
1079 Ok(()) => st.question = Some(q.id.clone()),
1080 Err(e) => tracing::warn!("could not file the question for {pr}: {e:#}"),
1081 }
1082 }
1083
1084 fn settle_job_question(&self, pr: &str, st: &mut WatchState, job: &mut Job) -> bool {
1088 let why = job.failed.clone().unwrap_or_default();
1089 if let Some(id) = st.question.clone() {
1090 match self.questions().get(&id) {
1091 Ok(q) if q.status.open() => return false,
1092 Ok(q) => {
1093 st.question = None;
1094 match &q.answer {
1095 Some(Answer::Choice(c)) if c == RETRY => {
1096 job.resume();
1097 st.held = None;
1098 return true;
1099 }
1100 Some(Answer::Choice(c)) if c == LEAVE_IT => {
1101 st.ignored = true;
1102 let _ = Notices::at(self.home.join("notifications"))
1103 .dismiss(¬ices::id_of(¬ice_key(pr)));
1104 }
1105 _ => st.held = Some(why),
1107 }
1108 return false;
1109 }
1110 Err(_) => st.question = None,
1111 }
1112 }
1113 if !st.ignored && st.held.as_deref() != Some(why.as_str()) {
1114 self.raise(pr, failed_message(pr));
1115 self.file_job_question(pr, st, job);
1116 }
1117 false
1118 }
1119
1120 fn apply_answer(&self, pr: &str, st: &mut WatchState, snap: &RollupView, q: &Question) {
1122 st.question = None;
1123 if st.applied.contains(&q.id) {
1124 return;
1125 }
1126 st.applied.push(q.id.clone());
1127 if st.applied.len() > 20 {
1128 st.applied.remove(0);
1129 }
1130 match &q.answer {
1131 Some(Answer::Choice(c)) if c == RERUN_AGAIN => {
1132 st.held = None;
1133 for c in snap.checks.iter().filter(|c| is_failed(c.verdict)) {
1134 if let Some(r) = &c.run {
1135 st.reruns.remove(r);
1136 }
1137 }
1138 }
1139 Some(Answer::Choice(c)) if c == LEAVE_IT => {
1140 st.ignored = true;
1141 if let Err(e) = Notices::at(self.home.join("notifications"))
1142 .dismiss(¬ices::id_of(¬ice_key(pr)))
1143 {
1144 tracing::warn!("could not dismiss the notice for {pr}: {e:#}");
1145 }
1146 }
1147 _ => st.held = Some(fingerprint(snap)),
1149 }
1150 }
1151
1152 fn finish(&self, pr: &str, st: &WatchState) {
1154 let _ =
1155 Notices::at(self.home.join("notifications")).dismiss(¬ices::id_of(¬ice_key(pr)));
1156 if let Some(id) = &st.question {
1157 let _ = self.questions().update(id, |q| {
1158 q.abandon("the release pull request is no longer open");
1159 Ok(())
1160 });
1161 }
1162 let _ = std::fs::remove_file(self.state_path(pr));
1163 }
1164}
1165
1166pub(crate) async fn run(
1170 watcher: Watcher,
1171 settings: impl Fn() -> (Vec<PathBuf>, u64),
1172 stop: crate::daemon::Stop,
1173) {
1174 let halt = {
1175 let stop = stop.clone();
1176 move || stop.stopped()
1177 };
1178 while !stop.stopped() {
1179 let (repos, minutes) = settings();
1180 if minutes > 0 {
1181 let now = jiff::Timestamp::now().as_second();
1182 watcher
1183 .lap(&repos, (minutes as i64).saturating_mul(60), now, &halt)
1184 .await;
1185 }
1186 let mut slept = Duration::ZERO;
1187 while slept < LAP && !stop.stopped() {
1188 tokio::time::sleep(Duration::from_secs(1)).await;
1189 slept += Duration::from_secs(1);
1190 }
1191 }
1192}
1193
1194#[cfg(test)]
1195mod tests {
1196 use super::*;
1197 use crate::land::CheckView;
1198 use anyhow::Context;
1199 use std::sync::Mutex;
1200
1201 const URL: &str = "https://github.com/o/r/pull/7";
1202
1203 fn check(name: &str, v: Verdict, run: Option<&str>) -> CheckView {
1204 CheckView {
1205 name: name.to_owned(),
1206 verdict: v,
1207 run: run.map(str::to_owned),
1208 url: run.map(|r| format!("https://github.com/o/r/actions/runs/{r}/job/1")),
1209 }
1210 }
1211
1212 fn snap(state: PrLifecycle, head: &str, checks: Vec<CheckView>) -> RollupView {
1213 RollupView {
1214 url: URL.to_owned(),
1215 number: 7,
1216 state,
1217 head: head.to_owned(),
1218 checks,
1219 }
1220 }
1221
1222 fn open(checks: Vec<CheckView>) -> RollupView {
1223 snap(PrLifecycle::Open, "h1", checks)
1224 }
1225
1226 fn st_at(progress: i64) -> WatchState {
1227 WatchState {
1228 progress_at: progress,
1229 ..WatchState::default()
1230 }
1231 }
1232
1233 #[test]
1234 fn unreadable_is_unknown_and_merged_or_closed_is_done() {
1235 assert_eq!(decide(None, &st_at(0), 10, 60), Step::Unknown);
1236 for s in [PrLifecycle::Merged, PrLifecycle::Closed] {
1237 assert_eq!(
1238 decide(Some(&snap(s, "h", vec![])), &st_at(0), 10, 60),
1239 Step::Done
1240 );
1241 }
1242 }
1243
1244 #[test]
1245 fn a_failed_run_is_rerun_even_while_another_run_is_pending() {
1246 let s = open(vec![
1247 check("win", Verdict::Fail, Some("11")),
1248 check("lint", Verdict::Pending, Some("12")),
1249 ]);
1250 assert_eq!(
1251 decide(Some(&s), &st_at(0), 1, 3600),
1252 Step::Rerun(vec!["11".into()])
1253 );
1254 }
1255
1256 #[test]
1257 fn a_run_with_a_job_still_pending_is_not_complete() {
1258 let s = open(vec![
1259 check("win", Verdict::Fail, Some("11")),
1260 check("mac", Verdict::Pending, Some("11")),
1261 ]);
1262 assert_eq!(decide(Some(&s), &st_at(0), 1, 3600), Step::Wait);
1263 }
1264
1265 #[test]
1266 fn one_rerun_per_run_then_grace_then_escalation() {
1267 let s = open(vec![check("win", Verdict::Fail, Some("11"))]);
1268 let mut st = st_at(0);
1269 st.reruns.insert("11".into(), 100);
1270 assert_eq!(decide(Some(&s), &st, 100 + RERUN_GRACE - 1, 0), Step::Wait);
1271 assert_eq!(
1272 decide(Some(&s), &st, 100 + RERUN_GRACE, 0),
1273 Step::Escalate(Why::StillRed)
1274 );
1275 }
1276
1277 #[test]
1278 fn a_failure_with_no_run_id_goes_straight_to_a_human_once_settled() {
1279 let s = open(vec![check("ci/ext", Verdict::Fail, None)]);
1280 assert_eq!(
1281 decide(Some(&s), &st_at(0), 1, 0),
1282 Step::Escalate(Why::StillRed)
1283 );
1284 let s = open(vec![
1285 check("ci/ext", Verdict::Fail, None),
1286 check("x", Verdict::Pending, Some("5")),
1287 ]);
1288 assert_eq!(decide(Some(&s), &st_at(0), 1, 0), Step::Wait);
1289 }
1290
1291 #[test]
1292 fn no_progress_for_the_bounded_time_escalates_and_zero_disables_it() {
1293 let s = open(vec![check("a", Verdict::Pass, Some("1"))]);
1294 assert_eq!(decide(Some(&s), &st_at(0), 3599, 3600), Step::Wait);
1295 assert_eq!(
1296 decide(Some(&s), &st_at(0), 3600, 3600),
1297 Step::Escalate(Why::Stalled)
1298 );
1299 assert_eq!(decide(Some(&s), &st_at(0), 99_999, 0), Step::Wait);
1300 }
1301
1302 #[test]
1303 fn held_and_ignored_stay_quiet_until_the_fingerprint_moves() {
1304 let s = open(vec![check("a", Verdict::Fail, None)]);
1305 let mut st = st_at(0);
1306 st.held = Some(fingerprint(&s));
1307 assert_eq!(decide(Some(&s), &st, 99_999, 60), Step::Wait);
1308 let moved = open(vec![
1309 check("a", Verdict::Pass, None),
1310 check("b", Verdict::Fail, None),
1311 ]);
1312 assert_eq!(
1313 decide(Some(&moved), &st, 99_999, 60),
1314 Step::Escalate(Why::StillRed)
1315 );
1316 st.ignored = true;
1317 assert_eq!(decide(Some(&moved), &st, 99_999, 60), Step::Wait);
1318 }
1319
1320 #[test]
1321 fn observe_restarts_the_clock_only_on_change_and_forgets_reruns_on_a_new_head() {
1322 let mut st = st_at(0);
1323 let a = open(vec![check("a", Verdict::Pending, Some("1"))]);
1324 observe(&mut st, &a, 10);
1325 assert_eq!(st.progress_at, 10);
1326 observe(&mut st, &a, 50);
1327 assert_eq!(st.progress_at, 10);
1328 st.reruns.insert("1".into(), 5);
1329 observe(
1330 &mut st,
1331 &snap(PrLifecycle::Open, "h2", a.checks.clone()),
1332 60,
1333 );
1334 assert!(st.reruns.is_empty());
1335 assert_eq!(st.progress_at, 60);
1336 }
1337
1338 #[test]
1339 fn pr_key_names_owner_repo_and_number() {
1340 assert_eq!(pr_key(URL).as_deref(), Some("o/r#7"));
1341 assert_eq!(pr_key("https://github.com/o/r/issues/7"), None);
1342 assert_eq!(pr_key("nonsense"), None);
1343 }
1344
1345 #[test]
1346 fn notice_wording_does_not_vary_between_polls() {
1347 assert_eq!(rerun_message("o/r#7"), rerun_message("o/r#7"));
1348 assert!(
1349 !human_message("o/r#7")
1350 .chars()
1351 .any(|c| c.is_ascii_digit() && c != '7')
1352 );
1353 }
1354
1355 #[test]
1356 fn the_release_question_is_not_claimed_by_other_machinery() {
1357 let q = Question::new(
1358 String::new(),
1359 crate::bump::NOTICE_NODE.into(),
1360 "release-watch".into(),
1361 "s".into(),
1362 String::new(),
1363 vec![HOLD.into()],
1364 );
1365 assert_eq!(crate::deputy::kind_of(&q), None);
1366 let n = Notice::warn("release-pr:o/r#7", "m").about([String::new()]);
1367 assert!(!notices::covers(&q, &n));
1368 }
1369
1370 #[derive(Default)]
1371 struct Fake {
1372 snap: Mutex<Option<RollupView>>,
1373 reruns: Mutex<Vec<String>>,
1374 local: Mutex<bool>,
1375 info: Mutex<Option<PrInfo>>,
1376 merges: Mutex<Vec<String>>,
1377 }
1378
1379 impl ReleaseForge for std::sync::Arc<Fake> {
1380 fn list<'a>(&'a self, _: &'a Path) -> Fut<'a, Result<Vec<(String, String)>>> {
1381 Box::pin(async { Ok(vec![("chore/release-v1.0.0".to_owned(), URL.to_owned())]) })
1382 }
1383 fn snapshot<'a>(&'a self, _: &'a Path, _: &'a str) -> Fut<'a, Result<RollupView>> {
1384 let s = self.snap.lock().unwrap().clone();
1385 Box::pin(async move { s.context("unreadable") })
1386 }
1387 fn rerun<'a>(&'a self, _: &'a Path, run: &'a str) -> Fut<'a, Result<()>> {
1388 self.reruns.lock().unwrap().push(run.to_owned());
1389 Box::pin(async { Ok(()) })
1390 }
1391 fn config<'a>(&'a self, _: &'a Path) -> Fut<'a, Result<crate::config::Config>> {
1392 let mut c = crate::config::Config::default();
1393 if *self.local.lock().unwrap() {
1394 c.release.mode = crate::config::ReleaseMode::Local;
1395 c.release.commands = vec!["true".to_owned()];
1396 }
1397 Box::pin(async move { Ok(c) })
1398 }
1399 fn info<'a>(&'a self, _: &'a Path, _: &'a str) -> Fut<'a, Result<PrInfo>> {
1400 let i = self.info.lock().unwrap().clone();
1401 Box::pin(async move { i.context("no info") })
1402 }
1403 fn merge<'a>(&'a self, _: &'a Path, _: &'a str, head: &'a str) -> Fut<'a, Result<bool>> {
1404 self.merges.lock().unwrap().push(head.to_owned());
1405 Box::pin(async { Ok(true) })
1406 }
1407 }
1408
1409 fn rig() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher) {
1410 let dir = tempfile::tempdir().unwrap();
1411 let fake = std::sync::Arc::new(Fake::default());
1412 let w = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
1413 (dir, fake, w)
1414 }
1415
1416 fn red() -> RollupView {
1417 open(vec![check("win", Verdict::Fail, Some("11"))])
1418 }
1419
1420 #[tokio::test]
1421 async fn reruns_once_survives_a_restart_then_escalates_with_one_question() {
1422 let (dir, fake, w) = rig();
1423 *fake.snap.lock().unwrap() = Some(red());
1424 let repo = PathBuf::from("/nowhere");
1425 let no = || false;
1426 w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
1427 assert_eq!(*fake.reruns.lock().unwrap(), vec!["11".to_owned()]);
1428
1429 let w2 = Watcher::new(Box::new(fake.clone()), dir.path().to_path_buf());
1431 w2.lap(
1432 std::slice::from_ref(&repo),
1433 3600,
1434 1000 + RERUN_GRACE - 1,
1435 &no,
1436 )
1437 .await;
1438 assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1439 assert!(w2.questions().list().is_empty());
1440
1441 w2.lap(
1442 std::slice::from_ref(&repo),
1443 3600,
1444 1000 + RERUN_GRACE + 1,
1445 &no,
1446 )
1447 .await;
1448 w2.lap(
1449 std::slice::from_ref(&repo),
1450 3600,
1451 1000 + RERUN_GRACE + 400,
1452 &no,
1453 )
1454 .await;
1455 assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1456 let qs = w2.questions().list();
1457 assert_eq!(qs.len(), 1, "asked once, not every lap");
1458 assert_eq!(qs[0].node, crate::bump::NOTICE_NODE);
1459 assert!(qs[0].detail.contains("win") && qs[0].detail.contains(URL));
1460 assert_eq!(qs[0].choices, vec![RERUN_AGAIN, HOLD, LEAVE_IT]);
1461 let ns = Notices::at(dir.path().join("notifications")).list();
1462 assert_eq!(ns.len(), 1);
1463 assert_eq!(ns[0].key, "release-pr:o/r#7");
1464 }
1465
1466 async fn escalated() -> (tempfile::TempDir, std::sync::Arc<Fake>, Watcher, String) {
1467 let (dir, fake, w) = rig();
1468 *fake.snap.lock().unwrap() = Some(red());
1469 let repo = PathBuf::from("/nowhere");
1470 let no = || false;
1471 w.lap(std::slice::from_ref(&repo), 3600, 1000, &no).await;
1472 w.lap(&[repo], 3600, 2000, &no).await;
1473 let id = w.questions().list()[0].id.clone();
1474 (dir, fake, w, id)
1475 }
1476
1477 fn answer(w: &Watcher, id: &str, choice: &str) {
1478 w.questions()
1479 .update(id, |q| q.answer(Answer::Choice(choice.to_owned())))
1480 .unwrap();
1481 }
1482
1483 #[tokio::test]
1484 async fn rerun_again_reruns_exactly_once_more() {
1485 let (_d, fake, w, id) = escalated().await;
1486 answer(&w, &id, RERUN_AGAIN);
1487 let repo = PathBuf::from("/nowhere");
1488 let no = || false;
1489 w.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
1490 w.lap(&[repo], 3600, 3001, &no).await;
1491 assert_eq!(fake.reruns.lock().unwrap().len(), 2);
1492 }
1493
1494 #[tokio::test]
1495 async fn hold_and_silence_do_not_reask_and_leave_it_dismisses() {
1496 let (d, fake, w, id) = escalated().await;
1497 answer(&w, &id, HOLD);
1498 let repo = PathBuf::from("/nowhere");
1499 let no = || false;
1500 for t in [3000, 90_000, 200_000] {
1501 w.lap(std::slice::from_ref(&repo), 3600, t, &no).await;
1502 }
1503 assert_eq!(w.questions().list().len(), 1, "held: no second question");
1504 assert_eq!(fake.reruns.lock().unwrap().len(), 1);
1505
1506 let (d2, _f2, w2, id2) = escalated().await;
1507 answer(&w2, &id2, LEAVE_IT);
1508 w2.lap(std::slice::from_ref(&repo), 3600, 3000, &no).await;
1509 let ns = Notices::at(d2.path().join("notifications")).list();
1510 assert!(ns.is_empty(), "dismissed notices are hidden");
1511 drop(d);
1512 }
1513
1514 #[tokio::test]
1515 async fn merged_clears_the_notice_the_question_and_the_record() {
1516 let (d, fake, w, _id) = escalated().await;
1517 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1518 let repo = PathBuf::from("/nowhere");
1519 w.lap(&[repo], 3600, 5000, &(|| false)).await;
1520 assert!(
1521 Notices::at(d.path().join("notifications"))
1522 .list()
1523 .is_empty()
1524 );
1525 assert!(w.questions().list().iter().all(|q| !q.status.open()));
1526 assert!(w.stored().is_empty());
1527 }
1528
1529 #[tokio::test]
1530 async fn an_unreadable_forge_changes_nothing() {
1531 let (d, fake, w) = rig();
1532 *fake.snap.lock().unwrap() = None;
1533 w.lap(&[PathBuf::from("/nowhere")], 60, 99_999, &(|| false))
1534 .await;
1535 assert!(fake.reruns.lock().unwrap().is_empty());
1536 assert!(w.questions().list().is_empty());
1537 assert!(!d.path().join("release-watch").exists());
1538 }
1539
1540 #[tokio::test]
1541 async fn no_rerun_when_the_record_cannot_be_saved() {
1542 let (d, fake, w) = rig();
1543 *fake.snap.lock().unwrap() = Some(red());
1544 std::fs::write(d.path().join("release-watch"), "x").unwrap();
1546 w.lap(&[PathBuf::from("/nowhere")], 3600, 1000, &(|| false))
1547 .await;
1548 assert!(fake.reruns.lock().unwrap().is_empty());
1549 }
1550
1551 #[tokio::test]
1552 async fn a_park_stops_the_lap_before_any_call() {
1553 let (_d, fake, w) = rig();
1554 *fake.snap.lock().unwrap() = Some(red());
1555 w.lap(&[PathBuf::from("/nowhere")], 60, 1000, &(|| true))
1556 .await;
1557 assert!(fake.reruns.lock().unwrap().is_empty());
1558 }
1559
1560 fn local_st(asked: Option<&str>, held: Option<&str>) -> WatchState {
1563 WatchState {
1564 asked_head: asked.map(str::to_owned),
1565 held_head: held.map(str::to_owned),
1566 ..WatchState::default()
1567 }
1568 }
1569
1570 #[test]
1571 fn local_decisions_never_wait_on_ci_and_bind_approval_to_the_head() {
1572 use Approved::*;
1573 let none = local_st(None, None);
1574 assert_eq!(decide_local("h1", &none, NoQuestion), LocalStep::Ask);
1575 assert_eq!(decide_local("h1", &none, Open), LocalStep::Wait);
1576 let asked = local_st(Some("h1"), None);
1577 assert_eq!(
1578 decide_local("h1", &asked, Merge),
1579 LocalStep::Merge("h1".to_owned())
1580 );
1581 assert_eq!(
1582 decide_local("h1", &asked, Hold),
1583 LocalStep::Hold("h1".to_owned())
1584 );
1585 assert_eq!(decide_local("h2", &asked, Merge), LocalStep::Ask);
1587 assert_eq!(decide_local("h2", &asked, Hold), LocalStep::Ask);
1588 let held = local_st(Some("h1"), Some("h1"));
1590 assert_eq!(decide_local("h1", &held, NoQuestion), LocalStep::Wait);
1591 assert_eq!(decide_local("h2", &held, NoQuestion), LocalStep::Ask);
1592 let mut ign = local_st(None, None);
1593 ign.ignored = true;
1594 assert_eq!(decide_local("h1", &ign, NoQuestion), LocalStep::Wait);
1595 }
1596
1597 #[test]
1598 fn pr_info_is_read_from_gh_json() {
1599 let i = parse_info(
1600 r#"{"headRefName":"chore/release-v1.2.3","title":"chore: release v1.2.3","mergeCommit":{"oid":"abc"}}"#,
1601 )
1602 .unwrap();
1603 assert_eq!(i.branch, "chore/release-v1.2.3");
1604 assert_eq!(i.merge_commit.as_deref(), Some("abc"));
1605 let open = parse_info(r#"{"headRefName":"b","title":"t","mergeCommit":null}"#).unwrap();
1606 assert_eq!(open.merge_commit, None);
1607 }
1608
1609 #[tokio::test]
1610 async fn a_local_release_asks_once_without_ci_reruns_or_stall_escalation() {
1611 let (_d, fake, w) = rig();
1612 *fake.local.lock().unwrap() = true;
1613 *fake.snap.lock().unwrap() = Some(red());
1616 let repo = PathBuf::from("/nowhere");
1617 let no = || false;
1618 w.lap(std::slice::from_ref(&repo), 60, 1_000_000, &no).await;
1619 w.lap(std::slice::from_ref(&repo), 60, 2_000_000, &no).await;
1620 assert!(fake.reruns.lock().unwrap().is_empty());
1621 let qs = w.questions().list();
1622 assert_eq!(qs.len(), 1, "one approval question, not one per lap");
1623 assert_eq!(qs[0].choices, vec![land::APPROVE, land::HOLD]);
1624 assert!(fake.merges.lock().unwrap().is_empty(), "silence is a hold");
1625 }
1626
1627 #[tokio::test]
1628 async fn merge_is_pinned_to_the_approved_head_and_a_moved_head_asks_again() {
1629 let (_d, fake, w) = rig();
1630 *fake.local.lock().unwrap() = true;
1631 *fake.snap.lock().unwrap() = Some(open(vec![]));
1632 let repo = PathBuf::from("/nowhere");
1633 let no = || false;
1634 w.lap(std::slice::from_ref(&repo), 60, 1, &no).await;
1635 let id = w.questions().list()[0].id.clone();
1636 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Open, "h2", vec![]));
1638 answer(&w, &id, land::APPROVE);
1639 w.lap(std::slice::from_ref(&repo), 60, 2, &no).await;
1640 assert!(fake.merges.lock().unwrap().is_empty());
1641 assert_eq!(w.questions().list().len(), 2, "asked again about h2");
1642 let id2 = w
1644 .questions()
1645 .list()
1646 .into_iter()
1647 .find(|q| q.status.open())
1648 .unwrap()
1649 .id;
1650 answer(&w, &id2, land::APPROVE);
1651 w.lap(std::slice::from_ref(&repo), 60, 3, &no).await;
1652 assert_eq!(*fake.merges.lock().unwrap(), vec!["h2".to_owned()]);
1653 }
1654
1655 #[tokio::test]
1656 async fn a_failed_release_holds_with_one_notice_and_one_question_and_never_retries() {
1657 let (d, fake, w) = rig();
1658 *fake.local.lock().unwrap() = true;
1659 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1660 *fake.info.lock().unwrap() = Some(PrInfo {
1661 branch: "chore/release-v1.0.0".to_owned(),
1662 merge_commit: Some("deadbeef".to_owned()),
1663 title: "chore: release v1.0.0".to_owned(),
1664 });
1665 let repo = d.path().join("nowhere");
1667 let no = || false;
1668 w.lap(std::slice::from_ref(&repo), 60, 1, &no).await;
1669 let st = w.load("o/r#7");
1670 let job = st.job.expect("the job is recorded");
1671 assert!(job.failed.is_some() && !job.tag_done && !job.finished);
1672 let qs = w.questions().list();
1673 assert_eq!(qs.len(), 1);
1674 assert_eq!(qs[0].choices, vec![RETRY, LEAVE_IT]);
1675 w.lap(std::slice::from_ref(&repo), 60, 2, &no).await;
1677 assert_eq!(w.questions().list().len(), 1);
1678 assert_eq!(Notices::at(d.path().join("notifications")).list().len(), 1);
1679 }
1680
1681 #[tokio::test]
1682 async fn a_failed_release_holds_the_task_whose_run_opened_the_pr() {
1683 let (d, fake, w) = rig();
1684 *fake.local.lock().unwrap() = true;
1685 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1686 *fake.info.lock().unwrap() = Some(PrInfo {
1687 branch: "chore/release-v1.0.0".to_owned(),
1688 merge_commit: Some("deadbeef".to_owned()),
1689 title: "t".to_owned(),
1690 });
1691 let repo = d.path().join("nowhere");
1692 let q = crate::queue::Queue::at(d.path().join("queue"));
1693 let mut t = crate::queue::Task::new(
1694 "t".to_owned(),
1695 "i".to_owned(),
1696 repo.clone(),
1697 crate::queue::Source::Human,
1698 );
1699 t.runs.push("run1".to_owned());
1700 t.status = crate::queue::TaskStatus::Done;
1701 q.put(&mut t).unwrap();
1702 register(d.path(), &repo, URL, "run1");
1703 w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1704 let held = q.get(&t.id).unwrap();
1705 assert_eq!(held.status, crate::queue::TaskStatus::Held);
1706 assert!(
1707 held.hold_reason
1708 .unwrap_or_default()
1709 .contains("release v1.0.0")
1710 );
1711 }
1712
1713 #[tokio::test]
1714 async fn unhold_gives_the_task_back_only_when_the_hold_is_ours() {
1715 let (d, _fake, w) = rig();
1716 let q = crate::queue::Queue::at(d.path().join("queue"));
1717 let mut t = crate::queue::Task::new(
1718 "t".to_owned(),
1719 "i".to_owned(),
1720 PathBuf::from("/r"),
1721 crate::queue::Source::Human,
1722 );
1723 t.runs.push("run1".to_owned());
1724 t.status = crate::queue::TaskStatus::Done;
1725 q.put(&mut t).unwrap();
1726 let mut st = WatchState {
1727 run: Some("run1".to_owned()),
1728 ..WatchState::default()
1729 };
1730 w.hold_task(&mut st, "failed");
1731 assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1732 w.hold_task(&mut st, "failed again");
1734 assert!(w.unhold_task(&mut st));
1735 assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Done);
1736 assert!(st.held_task.is_none());
1737
1738 let mut h = q.get(&t.id).unwrap();
1740 st.held_task = Some(h.id.clone());
1741 h.hold_manual(Some("mine".to_owned()));
1742 q.put(&mut h).unwrap();
1743 assert!(w.unhold_task(&mut st));
1744 assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1745 }
1746
1747 #[tokio::test]
1748 async fn a_restart_after_the_failure_was_saved_still_holds_the_task() {
1749 let (d, fake, w) = rig();
1750 *fake.local.lock().unwrap() = true;
1751 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1752 let repo = d.path().join("nowhere");
1753 let q = crate::queue::Queue::at(d.path().join("queue"));
1754 let mut t = crate::queue::Task::new(
1755 "t".to_owned(),
1756 "i".to_owned(),
1757 repo.clone(),
1758 crate::queue::Source::Human,
1759 );
1760 t.runs.push("run1".to_owned());
1761 t.status = crate::queue::TaskStatus::Done;
1762 q.put(&mut t).unwrap();
1763 let mut job = Job::new("1.0.0", URL, "deadbeef");
1765 job.failed = Some("command 1 exited 3".to_owned());
1766 let st = WatchState {
1767 repo: repo.to_string_lossy().into_owned(),
1768 url: URL.to_owned(),
1769 run: Some("run1".to_owned()),
1770 job: Some(job),
1771 ..WatchState::default()
1772 };
1773 w.save("o/r#7", &st);
1774 w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1775 assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Held);
1776 }
1777
1778 #[tokio::test]
1779 async fn a_restart_after_the_job_finished_still_gives_the_task_back() {
1780 let (d, fake, w) = rig();
1781 *fake.local.lock().unwrap() = true;
1782 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1783 let repo = d.path().join("nowhere");
1784 let q = crate::queue::Queue::at(d.path().join("queue"));
1785 let mut t = crate::queue::Task::new(
1786 "t".to_owned(),
1787 "i".to_owned(),
1788 repo.clone(),
1789 crate::queue::Source::Human,
1790 );
1791 t.runs.push("run1".to_owned());
1792 t.hold_machine(Some(format!("{HOLD_PREFIX}release v1.0.0 failed")));
1793 q.put(&mut t).unwrap();
1794 let mut job = Job::new("1.0.0", URL, "deadbeef");
1796 job.finished = true;
1797 let st = WatchState {
1798 repo: repo.to_string_lossy().into_owned(),
1799 url: URL.to_owned(),
1800 run: Some("run1".to_owned()),
1801 job: Some(job),
1802 held_task: Some(t.id.clone()),
1803 ..WatchState::default()
1804 };
1805 w.save("o/r#7", &st);
1806 w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1807 assert_eq!(q.get(&t.id).unwrap().status, crate::queue::TaskStatus::Done);
1808 assert!(!w.state_path("o/r#7").exists());
1809 }
1810
1811 #[tokio::test]
1812 async fn a_pull_request_registered_at_open_is_picked_up_even_if_unlisted() {
1813 let (d, _fake, w) = rig();
1814 register(d.path(), Path::new("/r"), URL, "run1");
1815 let st = w.load("o/r#7");
1816 assert_eq!(st.url, URL);
1817 assert_eq!(st.repo, "/r");
1818 let mut live = st.clone();
1820 live.head = "keep".to_owned();
1821 w.save("o/r#7", &live);
1822 register(d.path(), Path::new("/r"), URL, "run1");
1823 assert_eq!(w.load("o/r#7").head, "keep");
1824 }
1825
1826 #[tokio::test]
1827 async fn a_finished_release_is_kept_on_the_run_and_a_missing_run_still_completes() {
1828 let (d, fake, w) = rig();
1829 *fake.local.lock().unwrap() = true;
1830 *fake.snap.lock().unwrap() = Some(snap(PrLifecycle::Merged, "h1", vec![]));
1831 let repo = d.path().join("nowhere");
1832 let mut run = RunState::new(
1833 repo.clone(),
1834 "main".to_owned(),
1835 "0123456789abcdef".to_owned(),
1836 "t".to_owned(),
1837 crate::config::Config::default(),
1838 );
1839 run.id = "run1".to_owned();
1840 run.save_under(d.path()).unwrap();
1841 let mut job = Job::new("1.0.0", URL, "deadbeef");
1842 job.finished = true;
1843 job.log.push(release_local::StepLog {
1844 name: "cmd".to_owned(),
1845 code: Some(0),
1846 tail: "ok".to_owned(),
1847 output: None,
1848 });
1849 let mk = |run: &str| WatchState {
1850 repo: repo.to_string_lossy().into_owned(),
1851 url: URL.to_owned(),
1852 run: Some(run.to_owned()),
1853 job: Some(job.clone()),
1854 ..WatchState::default()
1855 };
1856 w.save("o/r#7", &mk("run1"));
1857 w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1858 assert!(!w.state_path("o/r#7").exists());
1859 let kept = RunState::load_under("run1", d.path()).unwrap();
1860 assert_eq!(kept.release_bump.unwrap().release, Some(job.clone()));
1861
1862 w.save("o/r#7", &mk("gone"));
1864 w.lap(std::slice::from_ref(&repo), 60, 1, &(|| false)).await;
1865 assert!(!w.state_path("o/r#7").exists());
1866 }
1867}