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
350pub(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
359pub 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
372pub(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
425const 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
443pub(crate) trait ReleaseForge: Send + Sync {
445 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 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 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 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#[derive(Debug, Clone, Default, PartialEq, Eq)]
466pub(crate) struct PrInfo {
467 pub branch: String,
469 pub merge_commit: Option<String>,
471 pub title: String,
473}
474
475pub(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
496pub(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 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 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
590pub(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 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(¬ice_key(pr), message));
652 }
653
654 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 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 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 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 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 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 for r in &runs {
765 st.reruns.insert(r.clone(), now);
766 }
767 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 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 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 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 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 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 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 self.hold_task(&mut st, &format!("release {} is on hold: {why}", job.tag()));
955 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, ©)
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 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 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 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 #[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 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 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 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 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(¬ices::id_of(¬ice_key(pr)));
1179 }
1180 _ => 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 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(¬ices::id_of(¬ice_key(pr)))
1218 {
1219 tracing::warn!("could not dismiss the notice for {pr}: {e:#}");
1220 }
1221 }
1222 _ => st.held = Some(fingerprint(snap)),
1224 }
1225 }
1226
1227 fn finish(&self, pr: &str, st: &WatchState) {
1229 let _ =
1230 Notices::at(self.home.join("notifications")).dismiss(¬ices::id_of(¬ice_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
1241pub(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 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 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 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 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 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 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 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 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 assert_eq!(decide_local("h2", &asked, Merge), LocalStep::Ask);
1762 assert_eq!(decide_local("h2", &asked, Hold), LocalStep::Ask);
1763 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 *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 *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 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 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 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 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 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 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 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 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 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}