1use std::path::{Path, PathBuf};
8use std::time::Duration;
9
10use anyhow::{Context, Result};
11use jiff::Timestamp;
12use serde::{Deserialize, Serialize};
13
14use crate::config::{Update, UpdateMode};
15
16pub const NO_AUTOUPDATE_ENV: &str = "MAGI_NO_AUTOUPDATE";
20
21pub fn default_interval() -> Duration {
23 kaishin::default_interval()
24}
25
26const MIN_INTERVAL: Duration = Duration::from_secs(60);
37
38pub fn effective_interval(cfg: &Update) -> Duration {
45 let interval = cfg
46 .interval
47 .as_deref()
48 .and_then(|s| kaishin::parse_interval(s).ok())
49 .unwrap_or_else(default_interval);
50 interval.max(MIN_INTERVAL)
51}
52
53pub fn disabled_by_env() -> bool {
55 match std::env::var(NO_AUTOUPDATE_ENV) {
56 Ok(v) => {
57 let v = v.trim();
58 !(v.is_empty() || v == "0" || v.eq_ignore_ascii_case("false"))
59 }
60 Err(_) => false,
61 }
62}
63
64const OWNER: &str = "yukimemi";
66const REPO: &str = "magi";
68const BIN: &str = "magi";
70const CRATE: &str = "magi-cli";
72
73pub fn repo_name() -> &'static str {
77 REPO
78}
79
80fn options() -> kaishin::KaishinOptions {
88 kaishin::KaishinOptions::new(OWNER, REPO, BIN, env!("CARGO_PKG_VERSION")).crate_name(CRATE)
89}
90
91fn state_path() -> Option<PathBuf> {
94 dirs::cache_dir().map(|d| d.join("magi").join("last_update_check.json"))
95}
96
97pub async fn run_self_update(yes: bool, check_only: bool, non_interactive: bool) -> Result<()> {
99 let opts = kaishin::UpdateOptions::new()
100 .yes(yes)
101 .check_only(check_only)
102 .non_interactive(non_interactive);
103 kaishin::run_self_update(&options(), opts).await
104}
105
106pub enum Pending {
108 Cached {
110 checker: Checker,
112 latest: kaishin::LatestRelease,
114 },
115 Notify {
117 checker: Checker,
119 handle: tokio::task::JoinHandle<Result<Option<kaishin::LatestRelease>>>,
121 },
122 Install {
124 handle: tokio::task::JoinHandle<Result<Option<kaishin::LatestRelease>>>,
126 },
127}
128
129#[derive(Clone)]
131pub struct Checker {
132 inner: kaishin::Checker,
133}
134
135impl Checker {
136 pub fn new(cfg: &Update) -> Option<Self> {
150 if cfg.mode == UpdateMode::Off {
151 return None;
152 }
153 let mut inner = kaishin::Checker::new(BIN, options());
154 if let Some(path) = state_path() {
155 inner = inner.state_path(path);
156 }
157 Some(Self {
158 inner: inner.interval(effective_interval(cfg)),
159 })
160 }
161
162 pub fn should_check(&self) -> bool {
164 self.inner.should_check()
165 }
166
167 pub async fn newer_release(&self) -> Result<Option<kaishin::LatestRelease>> {
179 self.inner.check_and_save().await
180 }
181
182 pub fn cached_update(&self) -> Option<kaishin::LatestRelease> {
184 self.inner.cached_update()
185 }
186
187 pub fn format_banner(&self, latest: &kaishin::LatestRelease) -> String {
189 self.inner.format_banner(latest)
190 }
191
192 #[cfg(test)]
197 pub(crate) fn for_test(interval: Duration, state_path: PathBuf) -> Self {
198 let opts = kaishin::KaishinOptions::new(OWNER, REPO, BIN, env!("CARGO_PKG_VERSION"));
199 Self {
200 inner: kaishin::Checker::new(BIN, opts)
201 .state_path(state_path)
202 .interval(interval),
203 }
204 }
205}
206
207#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
213#[serde(rename_all = "snake_case")]
214pub enum Stage {
215 Downloading,
223 Replaced,
226 Parking,
230 Restarting,
234 Done,
237 Failed,
240}
241
242impl Stage {
243 #[must_use]
245 pub fn terminal(self) -> bool {
246 matches!(self, Self::Done | Self::Failed)
247 }
248}
249
250#[derive(Debug, Clone, Serialize, Deserialize)]
260pub struct Progress {
261 pub stage: Stage,
263 pub from: String,
265 pub to: Option<String>,
267 #[serde(default)]
269 pub parked_run: Option<String>,
270 pub started_at: Timestamp,
272 pub updated_at: Timestamp,
274 #[serde(default)]
276 pub detail: Option<String>,
277}
278
279impl Progress {
280 #[must_use]
282 pub fn new(from: String, to: String) -> Self {
283 let now = Timestamp::now();
284 Self {
285 stage: Stage::Downloading,
286 from,
287 to: Some(to),
288 parked_run: None,
289 started_at: now,
290 updated_at: now,
291 detail: None,
292 }
293 }
294
295 pub fn advance(&mut self, stage: Stage) {
297 self.stage = stage;
298 self.updated_at = Timestamp::now();
299 self.detail = None;
301 }
302
303 pub fn fail(&mut self, detail: impl Into<String>) {
305 self.stage = Stage::Failed;
306 self.updated_at = Timestamp::now();
307 self.detail = Some(detail.into());
308 }
309}
310
311#[must_use]
314pub fn progress_path(home: &Path) -> PathBuf {
315 home.join("upgrade.json")
316}
317
318pub fn write_progress(home: &Path, progress: &Progress) -> Result<()> {
324 let path = progress_path(home);
325 if let Some(parent) = path.parent() {
326 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
327 }
328 let body = serde_json::to_string_pretty(progress).context("serialize upgrade progress")?;
329 let tmp = path.with_extension("json.tmp");
330 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
331 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
332 Ok(())
333}
334
335#[must_use]
337pub fn read_progress(home: &Path) -> Option<Progress> {
338 let body = std::fs::read_to_string(progress_path(home)).ok()?;
339 serde_json::from_str(&body).ok()
340}
341
342#[must_use]
345pub fn log_path(home: &Path) -> PathBuf {
346 home.join("upgrade.log")
347}
348
349pub const LOG_MAX_BYTES: u64 = 256 * 1024;
351
352pub const STALL_AFTER_SECS: i64 = 120;
354
355pub const PARKING_STALL_AFTER_SECS: i64 = 70 * 60;
358
359pub const HEARTBEAT_SECS: i64 = 60;
361
362pub const WATCHDOG_POLL: Duration = Duration::from_secs(30);
364
365pub fn append_bounded(path: &Path, line: &str, max: u64) -> std::io::Result<()> {
369 use std::io::Write as _;
370 if std::fs::metadata(path).is_ok_and(|m| m.len() >= max) {
371 let mut old = path.as_os_str().to_owned();
372 old.push(".1");
373 std::fs::rename(path, PathBuf::from(old))?;
374 }
375 if let Some(parent) = path.parent() {
376 std::fs::create_dir_all(parent)?;
377 }
378 let mut file = std::fs::OpenOptions::new()
379 .create(true)
380 .append(true)
381 .open(path)?;
382 writeln!(file, "{line}")
383}
384
385pub fn log_step(home: &Path, msg: &str) {
389 tracing::info!("handover: {msg}");
390 log_line(home, "INFO", msg);
391}
392
393pub fn log_warn(home: &Path, msg: &str) {
395 tracing::warn!("handover: {msg}");
396 log_line(home, "WARN", msg);
397}
398
399fn log_line(home: &Path, level: &str, msg: &str) {
400 let line = format!(
401 "{} pid={} {level} {msg}",
402 Timestamp::now(),
403 std::process::id()
404 );
405 if let Err(e) = append_bounded(&log_path(home), &line, LOG_MAX_BYTES) {
406 tracing::warn!("could not append to {}: {e}", log_path(home).display());
407 }
408}
409
410pub fn write_progress_logged(home: &Path, progress: &Progress) {
414 if let Err(e) = write_progress(home, progress) {
415 log_warn(home, &format!("could not write upgrade.json: {e:#}"));
416 }
417}
418
419#[derive(Debug, Clone, PartialEq, Eq)]
421pub struct Stall {
422 pub stage: Stage,
424 pub age_secs: i64,
426 pub waiting_on: String,
428}
429
430#[must_use]
433pub fn stage_age_secs(progress: &Progress, now: Timestamp) -> i64 {
434 (now.as_second() - progress.updated_at.as_second()).max(0)
435}
436
437fn waiting_on(progress: &Progress) -> String {
439 match progress.stage {
440 Stage::Replaced => "serve() observing the handover signal and calling hand_over \
441 (hand_over has not recorded `parking`)"
442 .to_owned(),
443 Stage::Parking => match &progress.parked_run {
444 Some(run) => format!("the loop to finish run {run} at its next node boundary"),
445 None => "the loop to stop (no run was recorded as in flight)".to_owned(),
446 },
447 Stage::Restarting => "spawn_successor returning and this process exiting".to_owned(),
448 Stage::Downloading => "the release download and binary replacement".to_owned(),
449 Stage::Done | Stage::Failed => String::new(),
450 }
451}
452
453#[must_use]
455pub fn stall(progress: &Progress, now: Timestamp) -> Option<Stall> {
456 let limit = match progress.stage {
457 Stage::Replaced | Stage::Restarting => STALL_AFTER_SECS,
458 Stage::Parking => PARKING_STALL_AFTER_SECS,
459 Stage::Downloading | Stage::Done | Stage::Failed => return None,
460 };
461 let age_secs = stage_age_secs(progress, now);
462 (age_secs > limit).then(|| Stall {
463 stage: progress.stage,
464 age_secs,
465 waiting_on: waiting_on(progress),
466 })
467}
468
469#[derive(Debug, Clone, PartialEq, Eq)]
471pub struct Beat {
472 pub stage: Stage,
474 pub warn: bool,
476 pub message: String,
478}
479
480#[derive(Debug, Default)]
482pub struct Watchdog {
483 last: Option<(Stage, Timestamp)>,
484}
485
486impl Watchdog {
487 pub fn tick(&mut self, progress: &Progress, now: Timestamp) -> Option<Beat> {
490 if progress.stage.terminal() {
491 self.last = None;
492 return None;
493 }
494 if self.last.is_some_and(|(stage, _)| stage != progress.stage) {
495 self.last = None;
496 }
497 let stalled = stall(progress, now);
498 if stalled.is_none() && progress.stage != Stage::Parking {
501 return None;
502 }
503 if let Some((_, at)) = self.last
504 && now.as_second() - at.as_second() < HEARTBEAT_SECS
505 {
506 return None;
507 }
508 self.last = Some((progress.stage, now));
509 let age = stage_age_secs(progress, now);
510 let (warn, message) = match &stalled {
511 Some(s) => (
512 true,
513 format!(
514 "stuck in {:?} for {} min {} s, waiting on {}",
515 s.stage,
516 s.age_secs / 60,
517 s.age_secs % 60,
518 s.waiting_on
519 ),
520 ),
521 None => (
522 false,
523 format!(
524 "parking for {} min {} s, waiting on {}",
525 age / 60,
526 age % 60,
527 waiting_on(progress)
528 ),
529 ),
530 };
531 Some(Beat {
532 stage: progress.stage,
533 warn,
534 message,
535 })
536 }
537}
538
539#[derive(Debug, Serialize, Deserialize)]
542struct Note {
543 stage: Stage,
544 stage_since: Timestamp,
545 message: String,
546}
547
548fn note_path(home: &Path) -> PathBuf {
549 home.join("upgrade.note.json")
550}
551
552fn write_note(home: &Path, progress: &Progress, message: &str) {
553 let note = Note {
554 stage: progress.stage,
555 stage_since: progress.updated_at,
556 message: message.to_owned(),
557 };
558 let path = note_path(home);
559 let tmp = path.with_extension("json.tmp");
560 let written = serde_json::to_string(¬e)
561 .map_err(std::io::Error::other)
562 .and_then(|body| std::fs::write(&tmp, body))
563 .and_then(|()| std::fs::rename(&tmp, &path));
564 if let Err(e) = written {
565 tracing::warn!("could not write {}: {e}", path.display());
566 }
567}
568
569#[must_use]
572pub fn read_note(home: &Path, progress: &Progress) -> Option<String> {
573 let body = std::fs::read_to_string(note_path(home)).ok()?;
574 let note: Note = serde_json::from_str(&body).ok()?;
575 (note.stage == progress.stage && note.stage_since == progress.updated_at)
576 .then_some(note.message)
577}
578
579pub fn spawn_watchdog(home: PathBuf) {
583 let spawned = std::thread::Builder::new()
584 .name("upgrade-watchdog".to_owned())
585 .spawn(move || {
586 let mut dog = Watchdog::default();
587 loop {
588 std::thread::sleep(WATCHDOG_POLL);
589 let Some(progress) = read_progress(&home) else {
590 continue;
591 };
592 let Some(beat) = dog.tick(&progress, Timestamp::now()) else {
593 continue;
594 };
595 if beat.warn {
596 log_warn(&home, &beat.message);
597 } else {
598 log_step(&home, &beat.message);
599 }
600 write_note(&home, &progress, &beat.message);
605 }
606 });
607 if let Err(e) = spawned {
608 tracing::warn!("could not start the upgrade watchdog: {e}");
609 }
610}
611
612pub fn reconcile_after_restart(home: &Path) {
624 let Some(mut progress) = read_progress(home) else {
625 return;
626 };
627 if progress.stage.terminal() {
628 return;
629 }
630 let running = env!("CARGO_PKG_VERSION");
631 if progress
637 .to
638 .as_deref()
639 .is_some_and(|to| to.trim_start_matches('v') == running)
640 {
641 progress.advance(Stage::Done);
642 } else {
643 let to = progress
644 .to
645 .clone()
646 .unwrap_or_else(|| "the expected release".to_owned());
647 progress.fail(format!(
648 "this process came up on {running}, not {to} - the upgrade may \
649 not have replaced the binary"
650 ));
651 }
652 write_progress_logged(home, &progress);
653}
654
655pub fn spawn(cfg: &Update, rt: &tokio::runtime::Handle) -> Option<Pending> {
657 if disabled_by_env() || cfg.mode == UpdateMode::Off {
658 return None;
659 }
660 let checker = Checker::new(cfg)?;
661 match cfg.mode {
662 UpdateMode::Off => None,
663 UpdateMode::Notify => {
664 if !checker.should_check() {
665 let latest = checker.cached_update()?;
666 return Some(Pending::Cached { checker, latest });
667 }
668 let inner = checker.inner.clone();
669 let handle = rt.spawn(async move { inner.check_and_save().await });
670 Some(Pending::Notify { checker, handle })
671 }
672 UpdateMode::Install => {
673 let inner = checker.inner.clone();
674 let handle = rt.spawn(async move { inner.auto_update().await });
675 Some(Pending::Install { handle })
676 }
677 }
678}
679
680pub async fn finalize(pending: Option<Pending>, budget: Duration) {
685 let Some(pending) = pending else {
686 return;
687 };
688 match pending {
689 Pending::Cached { checker, latest } => {
690 eprintln!("{}", checker.format_banner(&latest));
691 }
692 Pending::Notify { checker, handle } => {
693 if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
694 eprintln!("{}", checker.format_banner(&latest));
695 }
696 }
697 Pending::Install { handle } => {
698 if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
699 eprintln!("magi updated itself to {}", latest.tag_name);
700 }
701 }
702 }
703}
704
705#[cfg(test)]
706mod tests {
707 use super::*;
708
709 #[test]
714 fn effective_interval_floors_a_configured_interval_below_githubs_rate_limit() {
715 let cfg = Update {
716 mode: UpdateMode::Notify,
717 interval: Some("1s".to_owned()),
718 };
719 assert_eq!(
720 effective_interval(&cfg),
721 MIN_INTERVAL,
722 "an interval that would exceed GitHub's rate limit under continuous \
723 polling must be floored rather than honoured verbatim"
724 );
725
726 let sane = Update {
727 mode: UpdateMode::Notify,
728 interval: Some("2h".to_owned()),
729 };
730 assert_eq!(
731 effective_interval(&sane),
732 Duration::from_secs(2 * 60 * 60),
733 "an interval already above the floor must pass through unchanged"
734 );
735 }
736
737 #[test]
738 fn env_kill_switch_semantics() {
739 unsafe {
741 std::env::remove_var(NO_AUTOUPDATE_ENV);
742 }
743 assert!(!disabled_by_env());
744 for (value, disabled) in [
745 ("1", true),
746 ("true", true),
747 ("yes", true),
748 ("0", false),
749 ("false", false),
750 ("FALSE", false),
751 ("", false),
752 (" ", false),
753 ] {
754 unsafe {
755 std::env::set_var(NO_AUTOUPDATE_ENV, value);
756 }
757 assert_eq!(
758 disabled_by_env(),
759 disabled,
760 "MAGI_NO_AUTOUPDATE={value:?} should {} disable",
761 if disabled { "" } else { "not" }
762 );
763 }
764 unsafe {
765 std::env::remove_var(NO_AUTOUPDATE_ENV);
766 }
767 }
768
769 #[test]
770 fn off_mode_never_spawns() {
771 let rt = tokio::runtime::Builder::new_current_thread()
772 .enable_all()
773 .build()
774 .unwrap();
775 let cfg = Update {
776 mode: UpdateMode::Off,
777 interval: None,
778 };
779 assert!(spawn(&cfg, rt.handle()).is_none());
780 }
781
782 #[test]
783 fn state_path_lives_under_the_cache_dir() {
784 let path = state_path().expect("a cache dir on every supported platform");
785 assert!(path.ends_with("magi/last_update_check.json"));
786 let data = dirs::data_local_dir().unwrap_or_default();
787 assert!(
788 !path.starts_with(&data) || dirs::cache_dir() == dirs::data_local_dir(),
789 "throttle state must not sit in the run history directory"
790 );
791 }
792
793 #[tokio::test]
794 async fn finalize_of_nothing_is_a_no_op() {
795 finalize(None, Duration::from_millis(1)).await;
796 }
797
798 #[test]
808 fn checking_is_off_for_every_caller_when_the_config_says_off() {
809 assert!(
810 Checker::new(&Update {
811 mode: UpdateMode::Off,
812 interval: None,
813 })
814 .is_none(),
815 "an operator who writes mode = \"off\" means it"
816 );
817 for mode in [UpdateMode::Notify, UpdateMode::Install] {
818 assert!(
819 Checker::new(&Update {
820 mode,
821 interval: None,
822 })
823 .is_some(),
824 "{mode:?} still asks the forge"
825 );
826 }
827 }
828
829 #[test]
840 fn cached_update_answers_from_disk_with_no_network_call() {
841 let dir = tempfile::tempdir().expect("temp dir");
842 let path = dir.path().join("state.json");
843 let opts = kaishin::KaishinOptions::new("yukimemi", "magi", "magi", "0.1.0");
844 let checker = Checker {
845 inner: kaishin::Checker::new("magi", opts).state_path(path.clone()),
846 };
847
848 assert!(
849 checker.cached_update().is_none(),
850 "no state file yet must read as \"unknown\", not an error"
851 );
852
853 let state = kaishin::UpdateCheckState {
854 last_checked_unix: 0,
855 last_known_latest: Some("v9.9.9".to_owned()),
856 last_known_url: Some("https://example.invalid/9.9.9".to_owned()),
857 };
858 kaishin::save_check_state(&path, &state).expect("seed the state file");
859
860 let latest = checker.cached_update().expect("a newer release was cached");
861 assert_eq!(latest.tag_name, "v9.9.9");
862 }
863
864 #[test]
865 fn reconcile_after_restart_confirms_a_matching_version() {
866 let home = tempfile::tempdir().expect("temp home");
871 let mut progress = Progress::new(
872 "0.1.0".to_owned(),
873 format!("v{}", env!("CARGO_PKG_VERSION")),
874 );
875 progress.advance(Stage::Restarting);
876 write_progress(home.path(), &progress).expect("seed progress");
877
878 reconcile_after_restart(home.path());
879
880 let after = read_progress(home.path()).expect("progress on disk");
881 assert_eq!(
882 after.stage,
883 Stage::Done,
884 "the successor is running exactly the release that was asked for, \
885 `v` prefix and all"
886 );
887 }
888
889 #[test]
890 fn reconcile_after_restart_flags_a_mismatched_version() {
891 let home = tempfile::tempdir().expect("temp home");
892 let mut progress = Progress::new("0.1.0".to_owned(), "v9.9.9".to_owned());
893 progress.advance(Stage::Restarting);
894 write_progress(home.path(), &progress).expect("seed progress");
895
896 reconcile_after_restart(home.path());
897
898 let after = read_progress(home.path()).expect("progress on disk");
899 assert_eq!(after.stage, Stage::Failed);
900 assert!(
901 after.detail.is_some_and(|d| d.contains("9.9.9")),
902 "the operator needs to know which release it did not come back on"
903 );
904 }
905
906 #[test]
907 fn reconcile_after_restart_leaves_a_settled_record_alone() {
908 let home = tempfile::tempdir().expect("temp home");
909 let mut progress = Progress::new("0.1.0".to_owned(), "9.9.9".to_owned());
910 progress.advance(Stage::Done);
911 write_progress(home.path(), &progress).expect("seed progress");
912
913 reconcile_after_restart(home.path());
914
915 let after = read_progress(home.path()).expect("progress on disk");
916 assert_eq!(
917 after.stage,
918 Stage::Done,
919 "an already-settled record must not be rewritten by a later, unrelated start"
920 );
921 }
922
923 #[test]
924 fn reconcile_after_restart_with_nothing_on_disk_is_a_quiet_no_op() {
925 let home = tempfile::tempdir().expect("temp home");
926 reconcile_after_restart(home.path());
927 assert!(read_progress(home.path()).is_none());
928 }
929
930 fn at(secs: i64) -> Timestamp {
931 Timestamp::from_second(secs).expect("timestamp")
932 }
933
934 fn staged(stage: Stage, since: i64) -> Progress {
935 let mut p = Progress::new("0.1.0".to_owned(), "v0.2.0".to_owned());
936 p.stage = stage;
937 p.updated_at = at(since);
938 p
939 }
940
941 #[test]
942 fn stall_has_a_threshold_per_stage_and_is_silent_when_terminal() {
943 let p = staged(Stage::Replaced, 1000);
944 assert!(stall(&p, at(1000 + STALL_AFTER_SECS)).is_none());
945 let s = stall(&p, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
946 assert_eq!(s.stage, Stage::Replaced);
947 assert_eq!(s.age_secs, STALL_AFTER_SECS + 1);
948 assert!(s.waiting_on.contains("hand_over"), "{}", s.waiting_on);
949
950 let p = staged(Stage::Restarting, 1000);
951 assert!(stall(&p, at(1000 + STALL_AFTER_SECS + 1)).is_some());
952
953 let p = staged(Stage::Parking, 1000);
954 assert!(
955 stall(&p, at(1000 + 3600)).is_none(),
956 "an hour of parking is normal"
957 );
958 assert!(stall(&p, at(1000 + PARKING_STALL_AFTER_SECS + 1)).is_some());
959
960 for stage in [Stage::Done, Stage::Failed, Stage::Downloading] {
961 assert!(stall(&staged(stage, 0), at(1_000_000)).is_none());
962 }
963 }
964
965 #[test]
966 fn a_clock_that_went_backwards_is_age_zero() {
967 let p = staged(Stage::Replaced, 5000);
968 assert_eq!(stage_age_secs(&p, at(100)), 0);
969 assert!(stall(&p, at(100)).is_none());
970 }
971
972 #[test]
973 fn a_note_is_kept_beside_the_record_and_matches_only_its_own_stage() {
974 let home = tempfile::tempdir().expect("temp home");
975 let p = staged(Stage::Replaced, 1000);
976 write_progress(home.path(), &p).expect("write");
977 write_note(home.path(), &p, "stuck");
978 assert_eq!(read_note(home.path(), &p).as_deref(), Some("stuck"));
979 assert_eq!(read_progress(home.path()).unwrap().updated_at, at(1000));
980 assert!(read_note(home.path(), &staged(Stage::Parking, 1000)).is_none());
981 assert!(read_note(home.path(), &staged(Stage::Replaced, 2000)).is_none());
982 }
983
984 #[test]
985 fn the_watchdog_speaks_once_a_minute_and_resets_on_a_new_stage() {
986 let mut dog = Watchdog::default();
987 let p = staged(Stage::Replaced, 1000);
988 assert!(dog.tick(&p, at(1060)).is_none(), "not stalled yet");
989 let beat = dog.tick(&p, at(1200)).expect("stalled");
990 assert!(beat.warn);
991 assert!(dog.tick(&p, at(1230)).is_none(), "spoke 30 s ago");
992 assert!(dog.tick(&p, at(1260)).is_some(), "a minute later");
993
994 let parking = staged(Stage::Parking, 1260);
995 let beat = dog
996 .tick(&parking, at(1270))
997 .expect("parking heartbeat at once");
998 assert!(!beat.warn, "a short park is not a warning");
999 assert!(dog.tick(&parking, at(1300)).is_none());
1000 let late = dog
1001 .tick(&parking, at(1260 + PARKING_STALL_AFTER_SECS + 1))
1002 .expect("past the parking ceiling");
1003 assert!(late.warn);
1004
1005 let done = staged(Stage::Done, 0);
1006 assert!(dog.tick(&done, at(9_999_999)).is_none());
1007 }
1008
1009 #[test]
1010 fn the_upgrade_log_appends_and_rotates_to_one_generation() {
1011 let dir = tempfile::tempdir().expect("temp dir");
1012 let path = dir.path().join("upgrade.log");
1013 append_bounded(&path, "one", 16).expect("append");
1014 append_bounded(&path, "two", 16).expect("append");
1015 assert_eq!(std::fs::read_to_string(&path).unwrap(), "one\ntwo\n");
1016 append_bounded(&path, "three-and-more", 16).expect("append");
1017 append_bounded(&path, "four", 16).expect("append");
1018 assert_eq!(std::fs::read_to_string(&path).unwrap(), "four\n");
1019 let old = dir.path().join("upgrade.log.1");
1020 assert!(
1021 std::fs::read_to_string(old)
1022 .unwrap()
1023 .contains("three-and-more")
1024 );
1025 }
1026}