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]
246 pub fn rank(self) -> u8 {
247 match self {
248 Self::Downloading => 0,
249 Self::Replaced => 1,
250 Self::Parking => 2,
251 Self::Restarting => 3,
252 Self::Done | Self::Failed => 4,
253 }
254 }
255
256 #[must_use]
258 pub fn terminal(self) -> bool {
259 matches!(self, Self::Done | Self::Failed)
260 }
261
262 #[must_use]
264 pub fn as_str(self) -> &'static str {
265 match self {
266 Self::Downloading => "downloading",
267 Self::Replaced => "replaced",
268 Self::Parking => "parking",
269 Self::Restarting => "restarting",
270 Self::Done => "done",
271 Self::Failed => "failed",
272 }
273 }
274}
275
276#[derive(Debug, Clone, Serialize, Deserialize)]
286pub struct Progress {
287 pub stage: Stage,
289 pub from: String,
291 pub to: Option<String>,
293 #[serde(default)]
295 pub parked_run: Option<String>,
296 pub started_at: Timestamp,
298 pub updated_at: Timestamp,
300 #[serde(default)]
302 pub detail: Option<String>,
303}
304
305impl Progress {
306 #[must_use]
308 pub fn new(from: String, to: String) -> Self {
309 let now = Timestamp::now();
310 Self {
311 stage: Stage::Downloading,
312 from,
313 to: Some(to),
314 parked_run: None,
315 started_at: now,
316 updated_at: now,
317 detail: None,
318 }
319 }
320
321 pub fn advance(&mut self, stage: Stage) {
323 self.stage = stage;
324 self.updated_at = Timestamp::now();
325 self.detail = None;
327 }
328
329 pub fn fail(&mut self, detail: impl Into<String>) {
331 self.stage = Stage::Failed;
332 self.updated_at = Timestamp::now();
333 self.detail = Some(detail.into());
334 }
335}
336
337#[must_use]
340pub fn progress_path(home: &Path) -> PathBuf {
341 home.join("upgrade.json")
342}
343
344static WRITE_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
347
348#[must_use]
357pub fn monotonic(current: Option<&Progress>, candidate: &Progress) -> Progress {
358 match current {
359 Some(cur)
360 if !cur.stage.terminal()
361 && !candidate.stage.terminal()
362 && candidate.stage.rank() <= cur.stage.rank() =>
363 {
364 let mut kept = cur.clone();
365 if candidate.to.is_some() {
366 kept.to.clone_from(&candidate.to);
367 }
368 kept
369 }
370 _ => candidate.clone(),
371 }
372}
373
374pub fn write_progress(home: &Path, progress: &Progress) -> Result<()> {
381 let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
382 write_progress_locked(home, progress)
383}
384
385fn write_progress_locked(home: &Path, progress: &Progress) -> Result<()> {
386 let path = progress_path(home);
387 if let Some(parent) = path.parent() {
388 std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
389 }
390 let to_write = monotonic(read_progress(home).as_ref(), progress);
391 let body = serde_json::to_string_pretty(&to_write).context("serialize upgrade progress")?;
392 let tmp = path.with_extension("json.tmp");
393 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
394 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
395 Ok(())
396}
397
398pub fn fail_progress(home: &Path, detail: &str) -> Result<bool> {
404 let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
405 let Some(mut progress) = read_progress(home) else {
406 return Ok(false);
407 };
408 if matches!(progress.stage, Stage::Parking | Stage::Restarting) {
409 return Ok(false);
410 }
411 progress.fail(detail);
412 write_progress_locked(home, &progress)?;
413 Ok(true)
414}
415
416#[must_use]
418pub fn read_progress(home: &Path) -> Option<Progress> {
419 let body = std::fs::read_to_string(progress_path(home)).ok()?;
420 serde_json::from_str(&body).ok()
421}
422
423#[must_use]
426pub fn log_path(home: &Path) -> PathBuf {
427 home.join("upgrade.log")
428}
429
430pub const LOG_MAX_BYTES: u64 = 256 * 1024;
432
433pub const STALL_AFTER_SECS: i64 = 120;
435
436pub const LEASE_TTL_SECS: i64 = 90;
439
440pub const HEARTBEAT_SECS: i64 = 60;
442
443pub const WATCHDOG_POLL: Duration = Duration::from_secs(30);
445
446pub fn append_bounded(path: &Path, line: &str, max: u64) -> std::io::Result<()> {
450 use std::io::Write as _;
451 if std::fs::metadata(path).is_ok_and(|m| m.len() >= max) {
452 let mut old = path.as_os_str().to_owned();
453 old.push(".1");
454 std::fs::rename(path, PathBuf::from(old))?;
455 }
456 if let Some(parent) = path.parent() {
457 std::fs::create_dir_all(parent)?;
458 }
459 let mut file = std::fs::OpenOptions::new()
460 .create(true)
461 .append(true)
462 .open(path)?;
463 writeln!(file, "{line}")
464}
465
466pub fn log_step(home: &Path, msg: &str) {
470 tracing::info!("handover: {msg}");
471 log_line(home, "INFO", msg);
472}
473
474pub fn log_warn(home: &Path, msg: &str) {
476 tracing::warn!("handover: {msg}");
477 log_line(home, "WARN", msg);
478}
479
480fn log_line(home: &Path, level: &str, msg: &str) {
481 let line = format!(
482 "{} pid={} {level} {msg}",
483 Timestamp::now(),
484 std::process::id()
485 );
486 if let Err(e) = append_bounded(&log_path(home), &line, LOG_MAX_BYTES) {
487 tracing::warn!("could not append to {}: {e}", log_path(home).display());
488 }
489}
490
491pub fn write_progress_logged(home: &Path, progress: &Progress) {
495 if let Err(e) = write_progress(home, progress) {
496 log_warn(home, &format!("could not write upgrade.json: {e:#}"));
497 }
498}
499
500#[derive(Debug, Clone, Serialize, Deserialize)]
503pub struct HandoverLease {
504 pub entered_at: Timestamp,
506 pub beat_at: Timestamp,
508 #[serde(default)]
510 pub parked_run: Option<String>,
511}
512
513impl HandoverLease {
514 #[must_use]
516 pub fn fresh(&self, now: Timestamp) -> bool {
517 now.as_second() - self.beat_at.as_second() <= LEASE_TTL_SECS
518 }
519}
520
521#[must_use]
523pub fn lease_path(home: &Path) -> PathBuf {
524 home.join("upgrade.handover.json")
525}
526
527#[must_use]
529pub fn read_lease(home: &Path) -> Option<HandoverLease> {
530 let body = std::fs::read_to_string(lease_path(home)).ok()?;
531 serde_json::from_str(&body).ok()
532}
533
534fn write_lease(home: &Path, lease: &HandoverLease) {
535 let path = lease_path(home);
536 let tmp = path.with_extension("json.tmp");
537 let written = serde_json::to_string(lease)
538 .map_err(std::io::Error::other)
539 .and_then(|body| std::fs::write(&tmp, body))
540 .and_then(|()| std::fs::rename(&tmp, &path));
541 if let Err(e) = written {
542 log_warn(home, &format!("could not write the handover lease: {e}"));
543 }
544}
545
546#[derive(Debug)]
548pub struct LeaseGuard {
549 home: PathBuf,
550 lease: HandoverLease,
551}
552
553impl LeaseGuard {
554 #[must_use]
556 pub fn enter(home: &Path, parked_run: Option<String>) -> Self {
557 let now = Timestamp::now();
558 let lease = HandoverLease {
559 entered_at: now,
560 beat_at: now,
561 parked_run,
562 };
563 write_lease(home, &lease);
564 Self {
565 home: home.to_owned(),
566 lease,
567 }
568 }
569
570 #[must_use]
576 pub fn enter_parking(home: &Path) -> (Self, bool) {
577 let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
578 let progress = read_progress(home);
579 let this = Self::enter(home, progress.as_ref().and_then(|p| p.parked_run.clone()));
580 let recorded = match progress {
581 Some(mut p) => {
582 p.advance(Stage::Parking);
583 if let Err(e) = write_progress_locked(home, &p) {
584 log_warn(home, &format!("could not write upgrade.json: {e:#}"));
585 }
586 true
587 }
588 None => false,
589 };
590 (this, recorded)
591 }
592
593 pub fn beat(&mut self) {
595 self.lease.beat_at = Timestamp::now();
596 write_lease(&self.home, &self.lease);
597 }
598}
599
600impl Drop for LeaseGuard {
601 fn drop(&mut self) {
602 let _ = std::fs::remove_file(lease_path(&self.home));
603 }
604}
605
606#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
608#[serde(rename_all = "snake_case")]
609pub enum StallKind {
610 NeverEntered,
612 StoppedBeating,
615}
616
617#[derive(Debug, Clone, PartialEq, Eq)]
619pub struct Stall {
620 pub stage: Stage,
622 pub kind: StallKind,
624 pub age_secs: i64,
626 pub waiting_on: String,
628}
629
630#[must_use]
633pub fn stage_age_secs(progress: &Progress, now: Timestamp) -> i64 {
634 (now.as_second() - progress.updated_at.as_second()).max(0)
635}
636
637#[must_use]
639pub fn live_lease<'a>(
640 progress: &Progress,
641 lease: Option<&'a HandoverLease>,
642 now: Timestamp,
643) -> Option<&'a HandoverLease> {
644 lease.filter(|l| {
645 !progress.stage.terminal() && l.fresh(now) && l.entered_at >= progress.started_at
646 })
647}
648
649fn waiting_on(progress: &Progress, lease: Option<&HandoverLease>) -> String {
651 match progress.stage {
652 Stage::Replaced => "hand_over to start (HANDOVER was signalled; hand_over has left no \
653 record that it was entered)"
654 .to_owned(),
655 Stage::Parking => {
656 let run = lease
657 .and_then(|l| l.parked_run.as_ref())
658 .or(progress.parked_run.as_ref());
659 match run {
660 Some(run) => format!("the loop to finish run {run} at its next node boundary"),
661 None => "the loop to stop (no run was recorded as in flight)".to_owned(),
662 }
663 }
664 Stage::Restarting => "spawn_successor returning and this process exiting".to_owned(),
665 Stage::Downloading => "the release download and binary replacement".to_owned(),
666 Stage::Done | Stage::Failed => String::new(),
667 }
668}
669
670#[must_use]
672pub fn waited_secs(lease: &HandoverLease, now: Timestamp) -> i64 {
673 (now.as_second() - lease.entered_at.as_second()).max(0)
674}
675
676#[must_use]
681pub fn stall(progress: &Progress, lease: Option<&HandoverLease>, now: Timestamp) -> Option<Stall> {
682 if progress.stage.terminal() || progress.stage == Stage::Downloading {
683 return None;
684 }
685 if live_lease(progress, lease, now).is_some() {
686 return None;
687 }
688 let stale = lease.filter(|l| l.entered_at >= progress.started_at);
690 let (kind, since) = match (progress.stage, stale) {
691 (_, Some(l)) => (StallKind::StoppedBeating, l.beat_at),
692 (Stage::Replaced, None) => (StallKind::NeverEntered, progress.updated_at),
693 _ => (StallKind::StoppedBeating, progress.updated_at),
694 };
695 let age_secs = (now.as_second() - since.as_second()).max(0);
696 let limit = if stale.is_some() {
697 LEASE_TTL_SECS
698 } else {
699 STALL_AFTER_SECS
700 };
701 (age_secs > limit).then(|| Stall {
702 stage: progress.stage,
703 kind,
704 age_secs,
705 waiting_on: waiting_on(progress, lease),
706 })
707}
708
709#[derive(Debug, Clone, PartialEq, Eq)]
711pub struct Beat {
712 pub stage: Stage,
714 pub warn: bool,
716 pub message: String,
718}
719
720#[derive(Debug, Default)]
722pub struct Watchdog {
723 last: Option<(Stage, Timestamp)>,
724}
725
726impl Watchdog {
727 pub fn tick(
730 &mut self,
731 progress: &Progress,
732 lease: Option<&HandoverLease>,
733 now: Timestamp,
734 ) -> Option<Beat> {
735 if progress.stage.terminal() {
736 self.last = None;
737 return None;
738 }
739 if self.last.is_some_and(|(stage, _)| stage != progress.stage) {
740 self.last = None;
741 }
742 let stalled = stall(progress, lease, now);
743 let alive = live_lease(progress, lease, now);
744 if stalled.is_none() && alive.is_none() && progress.stage != Stage::Parking {
747 return None;
748 }
749 if let Some((_, at)) = self.last
750 && now.as_second() - at.as_second() < HEARTBEAT_SECS
751 {
752 return None;
753 }
754 self.last = Some((progress.stage, now));
755 let age = alive.map_or_else(|| stage_age_secs(progress, now), |l| waited_secs(l, now));
756 let (warn, message) = match &stalled {
757 Some(s) => (
758 true,
759 format!(
760 "stuck in {:?} for {} min {} s without progress, waiting on {}",
761 s.stage,
762 s.age_secs / 60,
763 s.age_secs % 60,
764 s.waiting_on
765 ),
766 ),
767 None => (
768 false,
769 format!(
770 "parking for {} min {} s, waiting on {}",
771 age / 60,
772 age % 60,
773 waiting_on(progress, lease)
774 ),
775 ),
776 };
777 Some(Beat {
778 stage: progress.stage,
779 warn,
780 message,
781 })
782 }
783}
784
785#[derive(Debug, Serialize, Deserialize)]
788struct Note {
789 stage: Stage,
790 stage_since: Timestamp,
791 message: String,
792}
793
794fn note_path(home: &Path) -> PathBuf {
795 home.join("upgrade.note.json")
796}
797
798fn write_note(home: &Path, progress: &Progress, message: &str) {
799 let note = Note {
800 stage: progress.stage,
801 stage_since: progress.updated_at,
802 message: message.to_owned(),
803 };
804 let path = note_path(home);
805 let tmp = path.with_extension("json.tmp");
806 let written = serde_json::to_string(¬e)
807 .map_err(std::io::Error::other)
808 .and_then(|body| std::fs::write(&tmp, body))
809 .and_then(|()| std::fs::rename(&tmp, &path));
810 if let Err(e) = written {
811 tracing::warn!("could not write {}: {e}", path.display());
812 }
813}
814
815#[must_use]
818pub fn read_note(home: &Path, progress: &Progress) -> Option<String> {
819 let body = std::fs::read_to_string(note_path(home)).ok()?;
820 let note: Note = serde_json::from_str(&body).ok()?;
821 (note.stage == progress.stage && note.stage_since == progress.updated_at)
822 .then_some(note.message)
823}
824
825pub fn spawn_watchdog(home: PathBuf) {
829 let spawned = std::thread::Builder::new()
830 .name("upgrade-watchdog".to_owned())
831 .spawn(move || {
832 let mut dog = Watchdog::default();
833 loop {
834 std::thread::sleep(WATCHDOG_POLL);
835 let Some(progress) = read_progress(&home) else {
836 continue;
837 };
838 let lease = read_lease(&home);
839 let Some(beat) = dog.tick(&progress, lease.as_ref(), Timestamp::now()) else {
840 continue;
841 };
842 if beat.warn {
843 log_warn(&home, &beat.message);
844 } else {
845 log_step(&home, &beat.message);
846 }
847 write_note(&home, &progress, &beat.message);
852 }
853 });
854 if let Err(e) = spawned {
855 tracing::warn!("could not start the upgrade watchdog: {e}");
856 }
857}
858
859pub fn reconcile_after_restart(home: &Path) {
871 let Some(mut progress) = read_progress(home) else {
872 return;
873 };
874 let _ = std::fs::remove_file(lease_path(home));
876 if progress.stage.terminal() {
877 return;
878 }
879 let running = env!("CARGO_PKG_VERSION");
880 if progress
886 .to
887 .as_deref()
888 .is_some_and(|to| to.trim_start_matches('v') == running)
889 {
890 progress.advance(Stage::Done);
891 } else {
892 let to = progress
893 .to
894 .clone()
895 .unwrap_or_else(|| "the expected release".to_owned());
896 progress.fail(format!(
897 "this process came up on {running}, not {to} - the upgrade may \
898 not have replaced the binary"
899 ));
900 }
901 write_progress_logged(home, &progress);
902}
903
904pub fn spawn(cfg: &Update, rt: &tokio::runtime::Handle) -> Option<Pending> {
906 if disabled_by_env() || cfg.mode == UpdateMode::Off {
907 return None;
908 }
909 let checker = Checker::new(cfg)?;
910 match cfg.mode {
911 UpdateMode::Off => None,
912 UpdateMode::Notify => {
913 if !checker.should_check() {
914 let latest = checker.cached_update()?;
915 return Some(Pending::Cached { checker, latest });
916 }
917 let inner = checker.inner.clone();
918 let handle = rt.spawn(async move { inner.check_and_save().await });
919 Some(Pending::Notify { checker, handle })
920 }
921 UpdateMode::Install => {
922 let inner = checker.inner.clone();
923 let handle = rt.spawn(async move { inner.auto_update().await });
924 Some(Pending::Install { handle })
925 }
926 }
927}
928
929pub async fn finalize(pending: Option<Pending>, budget: Duration) {
934 let Some(pending) = pending else {
935 return;
936 };
937 match pending {
938 Pending::Cached { checker, latest } => {
939 eprintln!("{}", checker.format_banner(&latest));
940 }
941 Pending::Notify { checker, handle } => {
942 if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
943 eprintln!("{}", checker.format_banner(&latest));
944 }
945 }
946 Pending::Install { handle } => {
947 if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
948 eprintln!("magi updated itself to {}", latest.tag_name);
949 }
950 }
951 }
952}
953
954#[cfg(test)]
955mod tests {
956 use super::*;
957
958 #[test]
963 fn effective_interval_floors_a_configured_interval_below_githubs_rate_limit() {
964 let cfg = Update {
965 mode: UpdateMode::Notify,
966 interval: Some("1s".to_owned()),
967 };
968 assert_eq!(
969 effective_interval(&cfg),
970 MIN_INTERVAL,
971 "an interval that would exceed GitHub's rate limit under continuous \
972 polling must be floored rather than honoured verbatim"
973 );
974
975 let sane = Update {
976 mode: UpdateMode::Notify,
977 interval: Some("2h".to_owned()),
978 };
979 assert_eq!(
980 effective_interval(&sane),
981 Duration::from_secs(2 * 60 * 60),
982 "an interval already above the floor must pass through unchanged"
983 );
984 }
985
986 #[test]
987 fn env_kill_switch_semantics() {
988 unsafe {
990 std::env::remove_var(NO_AUTOUPDATE_ENV);
991 }
992 assert!(!disabled_by_env());
993 for (value, disabled) in [
994 ("1", true),
995 ("true", true),
996 ("yes", true),
997 ("0", false),
998 ("false", false),
999 ("FALSE", false),
1000 ("", false),
1001 (" ", false),
1002 ] {
1003 unsafe {
1004 std::env::set_var(NO_AUTOUPDATE_ENV, value);
1005 }
1006 assert_eq!(
1007 disabled_by_env(),
1008 disabled,
1009 "MAGI_NO_AUTOUPDATE={value:?} should {} disable",
1010 if disabled { "" } else { "not" }
1011 );
1012 }
1013 unsafe {
1014 std::env::remove_var(NO_AUTOUPDATE_ENV);
1015 }
1016 }
1017
1018 #[test]
1019 fn off_mode_never_spawns() {
1020 let rt = tokio::runtime::Builder::new_current_thread()
1021 .enable_all()
1022 .build()
1023 .unwrap();
1024 let cfg = Update {
1025 mode: UpdateMode::Off,
1026 interval: None,
1027 };
1028 assert!(spawn(&cfg, rt.handle()).is_none());
1029 }
1030
1031 #[test]
1032 fn state_path_lives_under_the_cache_dir() {
1033 let path = state_path().expect("a cache dir on every supported platform");
1034 assert!(path.ends_with("magi/last_update_check.json"));
1035 let data = dirs::data_local_dir().unwrap_or_default();
1036 assert!(
1037 !path.starts_with(&data) || dirs::cache_dir() == dirs::data_local_dir(),
1038 "throttle state must not sit in the run history directory"
1039 );
1040 }
1041
1042 #[tokio::test]
1043 async fn finalize_of_nothing_is_a_no_op() {
1044 finalize(None, Duration::from_millis(1)).await;
1045 }
1046
1047 #[test]
1057 fn checking_is_off_for_every_caller_when_the_config_says_off() {
1058 assert!(
1059 Checker::new(&Update {
1060 mode: UpdateMode::Off,
1061 interval: None,
1062 })
1063 .is_none(),
1064 "an operator who writes mode = \"off\" means it"
1065 );
1066 for mode in [UpdateMode::Notify, UpdateMode::Install] {
1067 assert!(
1068 Checker::new(&Update {
1069 mode,
1070 interval: None,
1071 })
1072 .is_some(),
1073 "{mode:?} still asks the forge"
1074 );
1075 }
1076 }
1077
1078 #[test]
1089 fn cached_update_answers_from_disk_with_no_network_call() {
1090 let dir = tempfile::tempdir().expect("temp dir");
1091 let path = dir.path().join("state.json");
1092 let opts = kaishin::KaishinOptions::new("yukimemi", "magi", "magi", "0.1.0");
1093 let checker = Checker {
1094 inner: kaishin::Checker::new("magi", opts).state_path(path.clone()),
1095 };
1096
1097 assert!(
1098 checker.cached_update().is_none(),
1099 "no state file yet must read as \"unknown\", not an error"
1100 );
1101
1102 let state = kaishin::UpdateCheckState {
1103 last_checked_unix: 0,
1104 last_known_latest: Some("v9.9.9".to_owned()),
1105 last_known_url: Some("https://example.invalid/9.9.9".to_owned()),
1106 };
1107 kaishin::save_check_state(&path, &state).expect("seed the state file");
1108
1109 let latest = checker.cached_update().expect("a newer release was cached");
1110 assert_eq!(latest.tag_name, "v9.9.9");
1111 }
1112
1113 #[test]
1114 fn reconcile_after_restart_confirms_a_matching_version() {
1115 let home = tempfile::tempdir().expect("temp home");
1120 let mut progress = Progress::new(
1121 "0.1.0".to_owned(),
1122 format!("v{}", env!("CARGO_PKG_VERSION")),
1123 );
1124 progress.advance(Stage::Restarting);
1125 write_progress(home.path(), &progress).expect("seed progress");
1126
1127 reconcile_after_restart(home.path());
1128
1129 let after = read_progress(home.path()).expect("progress on disk");
1130 assert_eq!(
1131 after.stage,
1132 Stage::Done,
1133 "the successor is running exactly the release that was asked for, \
1134 `v` prefix and all"
1135 );
1136 }
1137
1138 #[test]
1139 fn reconcile_after_restart_flags_a_mismatched_version() {
1140 let home = tempfile::tempdir().expect("temp home");
1141 let mut progress = Progress::new("0.1.0".to_owned(), "v9.9.9".to_owned());
1142 progress.advance(Stage::Restarting);
1143 write_progress(home.path(), &progress).expect("seed progress");
1144
1145 reconcile_after_restart(home.path());
1146
1147 let after = read_progress(home.path()).expect("progress on disk");
1148 assert_eq!(after.stage, Stage::Failed);
1149 assert!(
1150 after.detail.is_some_and(|d| d.contains("9.9.9")),
1151 "the operator needs to know which release it did not come back on"
1152 );
1153 }
1154
1155 #[test]
1156 fn reconcile_after_restart_leaves_a_settled_record_alone() {
1157 let home = tempfile::tempdir().expect("temp home");
1158 let mut progress = Progress::new("0.1.0".to_owned(), "9.9.9".to_owned());
1159 progress.advance(Stage::Done);
1160 write_progress(home.path(), &progress).expect("seed progress");
1161
1162 reconcile_after_restart(home.path());
1163
1164 let after = read_progress(home.path()).expect("progress on disk");
1165 assert_eq!(
1166 after.stage,
1167 Stage::Done,
1168 "an already-settled record must not be rewritten by a later, unrelated start"
1169 );
1170 }
1171
1172 #[test]
1173 fn reconcile_after_restart_with_nothing_on_disk_is_a_quiet_no_op() {
1174 let home = tempfile::tempdir().expect("temp home");
1175 reconcile_after_restart(home.path());
1176 assert!(read_progress(home.path()).is_none());
1177 }
1178
1179 fn at(secs: i64) -> Timestamp {
1180 Timestamp::from_second(secs).expect("timestamp")
1181 }
1182
1183 fn staged(stage: Stage, since: i64) -> Progress {
1184 let mut p = Progress::new("0.1.0".to_owned(), "v0.2.0".to_owned());
1185 p.stage = stage;
1186 p.started_at = at(since);
1187 p.updated_at = at(since);
1188 p
1189 }
1190
1191 fn lease(entered: i64, beat: i64) -> HandoverLease {
1192 HandoverLease {
1193 entered_at: at(entered),
1194 beat_at: at(beat),
1195 parked_run: Some("r1".to_owned()),
1196 }
1197 }
1198
1199 #[test]
1200 fn a_handover_never_entered_is_stuck_and_says_only_what_is_known() {
1201 let p = staged(Stage::Replaced, 1000);
1202 assert!(stall(&p, None, at(1000 + STALL_AFTER_SECS)).is_none());
1203 let s = stall(&p, None, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
1204 assert_eq!(s.stage, Stage::Replaced);
1205 assert_eq!(s.kind, StallKind::NeverEntered);
1206 assert_eq!(s.age_secs, STALL_AFTER_SECS + 1);
1207 assert!(s.waiting_on.contains("hand_over"), "{}", s.waiting_on);
1208
1209 let p = staged(Stage::Restarting, 1000);
1210 let s = stall(&p, None, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
1211 assert_eq!(s.kind, StallKind::StoppedBeating);
1212
1213 for stage in [Stage::Done, Stage::Failed, Stage::Downloading] {
1214 assert!(stall(&staged(stage, 0), None, at(1_000_000)).is_none());
1215 }
1216 }
1217
1218 #[test]
1219 fn a_live_parking_wait_is_never_stuck_however_long_it_lasts() {
1220 let mut p = staged(Stage::Parking, 1000);
1221 p.started_at = at(900);
1222 let hours = 5 * 3600;
1223 let l = lease(1000, 1000 + hours);
1224 assert!(stall(&p, Some(&l), at(1000 + hours + 10)).is_none());
1225 let r = staged(Stage::Replaced, 1000);
1227 assert!(stall(&r, Some(&l), at(1000 + hours + 10)).is_none());
1228 let s = stall(&p, Some(&l), at(1000 + hours + LEASE_TTL_SECS + 1)).expect("stuck");
1230 assert_eq!(s.kind, StallKind::StoppedBeating);
1231 }
1232
1233 #[test]
1234 fn a_lease_from_an_earlier_upgrade_proves_nothing() {
1235 let p = staged(Stage::Replaced, 2000);
1236 let old = lease(10, 3000);
1237 assert!(live_lease(&p, Some(&old), at(3001)).is_none());
1238 }
1239
1240 #[test]
1241 fn a_stage_never_goes_backwards_but_a_new_upgrade_after_a_terminal_one_starts() {
1242 let parking = staged(Stage::Parking, 1000);
1243 for back in [Stage::Replaced, Stage::Downloading, Stage::Parking] {
1244 let mut cand = staged(back, 5000);
1245 cand.to = Some("v9.9.9".to_owned());
1246 let kept = monotonic(Some(&parking), &cand);
1247 assert_eq!(kept.stage, Stage::Parking);
1248 assert_eq!(kept.updated_at, at(1000));
1249 assert_eq!(kept.started_at, parking.started_at);
1250 assert_eq!(kept.to.as_deref(), Some("v9.9.9"), "data is refreshed");
1251 }
1252 assert_eq!(
1253 monotonic(Some(&parking), &staged(Stage::Restarting, 5000)).stage,
1254 Stage::Restarting
1255 );
1256 assert_eq!(
1257 monotonic(Some(&parking), &staged(Stage::Failed, 5000)).stage,
1258 Stage::Failed
1259 );
1260 let done = staged(Stage::Done, 1000);
1261 assert_eq!(
1262 monotonic(Some(&done), &staged(Stage::Downloading, 5000)).stage,
1263 Stage::Downloading
1264 );
1265 }
1266
1267 #[test]
1268 fn write_progress_refuses_a_regression_on_disk() {
1269 let home = tempfile::tempdir().expect("temp home");
1270 write_progress(home.path(), &staged(Stage::Parking, 1000)).expect("write");
1271 write_progress(home.path(), &staged(Stage::Replaced, 5000)).expect("write");
1272 let on_disk = read_progress(home.path()).expect("record");
1273 assert_eq!(on_disk.stage, Stage::Parking);
1274 assert_eq!(on_disk.updated_at, at(1000));
1275 }
1276
1277 #[test]
1278 fn a_failed_request_cannot_overwrite_a_live_handover() {
1279 let home = tempfile::tempdir().expect("temp home");
1280 write_progress(home.path(), &staged(Stage::Parking, 1000)).expect("write");
1281 assert!(!fail_progress(home.path(), "boom").expect("fail"));
1282 assert_eq!(read_progress(home.path()).unwrap().stage, Stage::Parking);
1283 write_progress(home.path(), &staged(Stage::Replaced, 1000)).ok();
1284 let fresh = tempfile::tempdir().expect("temp home");
1285 write_progress(fresh.path(), &staged(Stage::Downloading, 1000)).expect("write");
1286 assert!(fail_progress(fresh.path(), "boom").expect("fail"));
1287 assert_eq!(read_progress(fresh.path()).unwrap().stage, Stage::Failed);
1288 }
1289
1290 #[test]
1291 fn entering_parking_is_one_step_that_keeps_the_lease_newer_than_the_record() {
1292 let home = tempfile::tempdir().expect("temp home");
1293 write_progress(home.path(), &staged(Stage::Replaced, 1000)).expect("write");
1294 let (guard, recorded) = LeaseGuard::enter_parking(home.path());
1295 assert!(recorded);
1296 let p = read_progress(home.path()).expect("record");
1297 assert_eq!(p.stage, Stage::Parking);
1298 let l = read_lease(home.path()).expect("lease");
1299 assert!(l.entered_at >= p.started_at);
1300 assert!(!fail_progress(home.path(), "boom").expect("fail"));
1301 drop(guard);
1302 }
1303
1304 #[test]
1305 fn the_lease_guard_writes_beats_and_removes_the_lease() {
1306 let home = tempfile::tempdir().expect("temp home");
1307 {
1308 let mut guard = LeaseGuard::enter(home.path(), Some("r1".to_owned()));
1309 let first = read_lease(home.path()).expect("lease");
1310 assert_eq!(first.parked_run.as_deref(), Some("r1"));
1311 guard.beat();
1312 assert!(read_lease(home.path()).is_some());
1313 }
1314 assert!(read_lease(home.path()).is_none());
1315 }
1316
1317 #[test]
1318 fn a_clock_that_went_backwards_is_age_zero() {
1319 let p = staged(Stage::Replaced, 5000);
1320 assert_eq!(stage_age_secs(&p, at(100)), 0);
1321 assert!(stall(&p, None, at(100)).is_none());
1322 }
1323
1324 #[test]
1325 fn a_note_is_kept_beside_the_record_and_matches_only_its_own_stage() {
1326 let home = tempfile::tempdir().expect("temp home");
1327 let p = staged(Stage::Replaced, 1000);
1328 write_progress(home.path(), &p).expect("write");
1329 write_note(home.path(), &p, "stuck");
1330 assert_eq!(read_note(home.path(), &p).as_deref(), Some("stuck"));
1331 assert_eq!(read_progress(home.path()).unwrap().updated_at, at(1000));
1332 assert!(read_note(home.path(), &staged(Stage::Parking, 1000)).is_none());
1333 assert!(read_note(home.path(), &staged(Stage::Replaced, 2000)).is_none());
1334 }
1335
1336 #[test]
1337 fn the_watchdog_speaks_once_a_minute_and_resets_on_a_new_stage() {
1338 let mut dog = Watchdog::default();
1339 let p = staged(Stage::Replaced, 1000);
1340 assert!(dog.tick(&p, None, at(1060)).is_none(), "not stalled yet");
1341 let beat = dog.tick(&p, None, at(1200)).expect("stalled");
1342 assert!(beat.warn);
1343 assert!(dog.tick(&p, None, at(1230)).is_none(), "spoke 30 s ago");
1344 assert!(dog.tick(&p, None, at(1260)).is_some(), "a minute later");
1345
1346 let parking = staged(Stage::Parking, 1260);
1347 let l = lease(1260, 1270);
1348 let beat = dog
1349 .tick(&parking, Some(&l), at(1275))
1350 .expect("parking heartbeat at once");
1351 assert!(!beat.warn, "a live wait is not a warning");
1352 assert!(dog.tick(&parking, Some(&l), at(1300)).is_none());
1353 let l = lease(1260, 1260 + 4 * 3600);
1354 let later = dog
1355 .tick(&parking, Some(&l), at(1260 + 4 * 3600 + 5))
1356 .expect("heartbeat");
1357 assert!(!later.warn, "hours of waiting on a run is still not stuck");
1358 assert!(later.message.contains("r1"), "{}", later.message);
1359
1360 let done = staged(Stage::Done, 0);
1361 assert!(dog.tick(&done, None, at(9_999_999)).is_none());
1362 }
1363
1364 #[test]
1365 fn the_upgrade_log_appends_and_rotates_to_one_generation() {
1366 let dir = tempfile::tempdir().expect("temp dir");
1367 let path = dir.path().join("upgrade.log");
1368 append_bounded(&path, "one", 16).expect("append");
1369 append_bounded(&path, "two", 16).expect("append");
1370 assert_eq!(std::fs::read_to_string(&path).unwrap(), "one\ntwo\n");
1371 append_bounded(&path, "three-and-more", 16).expect("append");
1372 append_bounded(&path, "four", 16).expect("append");
1373 assert_eq!(std::fs::read_to_string(&path).unwrap(), "four\n");
1374 let old = dir.path().join("upgrade.log.1");
1375 assert!(
1376 std::fs::read_to_string(old)
1377 .unwrap()
1378 .contains("three-and-more")
1379 );
1380 }
1381}