Skip to main content

magi/
updater.rs

1//! Self-update, via `kaishin`.
2//!
3//! A magi run takes minutes of agent latency, so a background release check
4//! costs nothing measurable: it is spawned on the same tokio runtime as the
5//! command, overlaps it, and is drained with a bounded wait at shutdown. It
6//! never delays the graph.
7use 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
16/// Env kill-switch. Any non-empty value other than `0` / `false` disables the
17/// background check, and it is read before the config so a broken `magi.toml`
18/// cannot force a network call.
19pub const NO_AUTOUPDATE_ENV: &str = "MAGI_NO_AUTOUPDATE";
20
21/// Default interval between checks.
22pub fn default_interval() -> Duration {
23    kaishin::default_interval()
24}
25
26/// Floor under [`effective_interval`], matching GitHub's unauthenticated rate
27/// limit for the releases API (60 requests/hour/address).
28///
29/// A one-off CLI invocation honouring a shorter configured interval could
30/// only ever make one call per process, so it was never at risk of tripping
31/// that limit on its own. The background recheck loop in `web.rs` is
32/// different: it polls for as long as `magi web` stays up, so an interval
33/// configured well below an hour would have it repeat the same call for as
34/// long as the deck runs - a shared floor here is what keeps it, and every
35/// other caller of [`Checker::new`], inside the budget regardless.
36const MIN_INTERVAL: Duration = Duration::from_secs(60);
37
38/// The interval `cfg` configures, floored at [`MIN_INTERVAL`], or
39/// [`default_interval`] when unset or unparsable.
40///
41/// Shared by [`Checker::new`], which throttles the network call itself
42/// against it, and the background recheck loop in `web.rs`, which uses it to
43/// decide how often to even ask - it cannot track an interval it never sees.
44pub 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
53/// Is the background check switched off by the environment?
54pub 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
64/// GitHub owner.
65const OWNER: &str = "yukimemi";
66/// GitHub repository — *not* `CARGO_PKG_NAME`, which is the published package.
67const REPO: &str = "magi";
68/// Binary inside the release asset.
69const BIN: &str = "magi";
70/// Published package name, for kaishin's `cargo install` fallback.
71const CRATE: &str = "magi-cli";
72
73/// This binary's own repository name, for `--repo .` auto-discovery
74/// (`web::serve`, `main::resolve_repo`) to match a checkout against without
75/// a second, driftable copy of [`REPO`] anywhere else.
76pub fn repo_name() -> &'static str {
77    REPO
78}
79
80/// kaishin options.
81///
82/// All four names are spelled out because three of them differ from
83/// `CARGO_PKG_NAME`: the package is `magi-cli` (the short name is a squatted
84/// placeholder on crates.io) while the repo, the binary and the library are
85/// `magi`. Deriving any of these from `CARGO_PKG_NAME` would send the updater
86/// looking for a `yukimemi/magi-cli` repository that does not exist.
87fn options() -> kaishin::KaishinOptions {
88    kaishin::KaishinOptions::new(OWNER, REPO, BIN, env!("CARGO_PKG_VERSION")).crate_name(CRATE)
89}
90
91/// Throttle bookkeeping is transient, so it belongs in the cache dir rather
92/// than beside the run history in the data dir.
93fn state_path() -> Option<PathBuf> {
94    dirs::cache_dir().map(|d| d.join("magi").join("last_update_check.json"))
95}
96
97/// `magi self-update`.
98pub 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
106/// A background update check, resolved at shutdown.
107pub enum Pending {
108    /// A previous run already found a newer release; just print the banner.
109    Cached {
110        /// For [`Checker::format_banner`].
111        checker: Checker,
112        /// The release found earlier.
113        latest: kaishin::LatestRelease,
114    },
115    /// A notify-mode check is in flight.
116    Notify {
117        /// For [`Checker::format_banner`].
118        checker: Checker,
119        /// The spawned task.
120        handle: tokio::task::JoinHandle<Result<Option<kaishin::LatestRelease>>>,
121    },
122    /// An install-mode update is in flight.
123    Install {
124        /// The spawned task.
125        handle: tokio::task::JoinHandle<Result<Option<kaishin::LatestRelease>>>,
126    },
127}
128
129/// Throttled release checker.
130#[derive(Clone)]
131pub struct Checker {
132    inner: kaishin::Checker,
133}
134
135impl Checker {
136    /// Build a checker honouring `cfg`, or `None` when checking is off.
137    ///
138    /// The `Option` had no `None` arm: every caller that asked for a checker
139    /// got one, so `[update] mode = "off"` was honoured by the *notify* path
140    /// alone (see [`cached_update`], which matches on the mode itself) and
141    /// ignored everywhere else. `POST /api/upgrade` therefore called the
142    /// GitHub releases API on a deck configured never to check - and so did
143    /// every unit test that reached that route, unauthenticated, against
144    /// GitHub's 60-per-hour-per-address limit.
145    ///
146    /// An operator who writes `mode = "off"` means it. The button is still
147    /// theirs to press; what it may not do is go to the network behind a
148    /// configuration that says not to.
149    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    /// Is a check due?
163    pub fn should_check(&self) -> bool {
164        self.inner.should_check()
165    }
166
167    /// Ask the forge now: is there a release newer than this build?
168    ///
169    /// Unlike [`Checker::cached_update`] this method does not consult
170    /// [`Checker::should_check`] itself - it has two callers, and they throttle
171    /// differently. `POST /api/upgrade` calls it unconditionally, because the
172    /// caller there is an operator who just pressed a button and is owed an
173    /// answer about the state of the world rather than about the last time
174    /// magi looked. `magi web`'s background recheck (`web::run_update_recheck`)
175    /// calls [`Checker::should_check`] itself first and only reaches here when
176    /// it says yes, which is what keeps that task's network use to at most
177    /// once per `[update] interval` no matter how often it polls.
178    pub async fn newer_release(&self) -> Result<Option<kaishin::LatestRelease>> {
179        self.inner.check_and_save().await
180    }
181
182    /// A newer release already known from a previous run.
183    pub fn cached_update(&self) -> Option<kaishin::LatestRelease> {
184        self.inner.cached_update()
185    }
186
187    /// One-line "a newer version exists" banner.
188    pub fn format_banner(&self, latest: &kaishin::LatestRelease) -> String {
189        self.inner.format_banner(latest)
190    }
191
192    /// A checker over an explicit state path and interval, for a test that
193    /// must control throttle timing without touching the operator's real
194    /// cache directory - see [`state_path`] for why sharing it would be
195    /// unsafe.
196    #[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/// How far a self-upgrade this deck set in motion has gotten.
208///
209/// `POST /api/upgrade` answers `202` and returns immediately - see its own
210/// doc for why - so [`Progress`] is the only way a phone that asked for an
211/// upgrade learns anything about it afterwards.
212#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
213#[serde(rename_all = "snake_case")]
214pub enum Stage {
215    /// Downloading the release asset and replacing the binary. This bundles
216    /// what would otherwise be two stages: kaishin re-confirms the release
217    /// and downloads it inside one `await` with no hook to split, so from
218    /// here a phone cannot tell "still checking" from "still downloading" -
219    /// only that nothing has been swapped in yet. The confirmation that ran
220    /// *before* this stage started is already known to the phone: it is what
221    /// the `to` version in the `202` answered with.
222    Downloading,
223    /// The new binary is in place and [`crate::web`]'s handover has been
224    /// signalled, but has not acted yet.
225    Replaced,
226    /// The handover is waiting for the run in flight, if any, to reach its
227    /// next node boundary. The deck answers throughout this - it is not the
228    /// unreachable gap, see `web::hand_over`.
229    Parking,
230    /// The listener has been released and the successor is starting. This is
231    /// the one genuinely unreachable moment, and it is meant to be
232    /// sub-second - see `web::bind_waiting`.
233    Restarting,
234    /// A successor came up and confirmed it is running the release this
235    /// upgrade asked for.
236    Done,
237    /// The upgrade did not reach [`Stage::Done`]. `detail` on [`Progress`]
238    /// says why.
239    Failed,
240}
241
242impl Stage {
243    /// Position in the upgrade's progression. Only meaningful for the
244    /// non-terminal stages; a terminal one ends the upgrade.
245    #[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    /// Finished, one way or the other - nothing is still moving.
257    #[must_use]
258    pub fn terminal(self) -> bool {
259        matches!(self, Self::Done | Self::Failed)
260    }
261
262    /// The wire spelling, the same one serde writes to `upgrade.json`.
263    #[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/// One upgrade's progress, persisted at [`progress_path`].
277///
278/// Kept on disk rather than in memory because the process that finishes an
279/// upgrade is never the one that started it: the successor is a fresh binary
280/// (see `web::spawn_successor`), and the only thing the two share is the
281/// disk. Beside the run history rather than under the cache dir alongside
282/// [`state_path`]: this is not throttle bookkeeping, it is the record of one
283/// upgrade the operator asked for, and - like a parked run - it is meant to
284/// outlive the process that wrote it.
285#[derive(Debug, Clone, Serialize, Deserialize)]
286pub struct Progress {
287    /// Where this upgrade has gotten to.
288    pub stage: Stage,
289    /// Version this upgrade started from.
290    pub from: String,
291    /// Version it is replacing itself with.
292    pub to: Option<String>,
293    /// The run [`Stage::Parking`] is waiting on, when one was in flight.
294    #[serde(default)]
295    pub parked_run: Option<String>,
296    /// When this upgrade was asked for.
297    pub started_at: Timestamp,
298    /// Last time `stage` changed.
299    pub updated_at: Timestamp,
300    /// Why [`Stage::Failed`] happened; `None` for every other stage.
301    #[serde(default)]
302    pub detail: Option<String>,
303}
304
305impl Progress {
306    /// A fresh record for an upgrade that is about to replace the binary.
307    #[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    /// Move to `stage`, stamping when it changed.
322    pub fn advance(&mut self, stage: Stage) {
323        self.stage = stage;
324        self.updated_at = Timestamp::now();
325        // A note belongs to the stage it was written for.
326        self.detail = None;
327    }
328
329    /// Stop at [`Stage::Failed`], with a reason a human can read.
330    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/// Where [`Progress`] is recorded: beside `daemon.json`, not under the cache
338/// dir - see [`Progress`]'s own doc for why the two are not the same place.
339#[must_use]
340pub fn progress_path(home: &Path) -> PathBuf {
341    home.join("upgrade.json")
342}
343
344/// Serialises the read-compare-rename of [`write_progress`], so two writers
345/// in this process cannot each compare against the same stale record.
346static WRITE_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
347
348/// What to persist when `candidate` is written over `current`.
349///
350/// An upgrade in progress never goes backwards: a candidate whose stage is
351/// not past the recorded one (a second `POST /api/upgrade` landing on a live
352/// handover wrote `replaced` over `parking`) keeps the recorded stage and
353/// timestamps and only refreshes `to`. A terminal record is over, so whatever
354/// comes next starts a new upgrade, and a candidate that ends the upgrade
355/// (`Done` / `Failed`) always goes through.
356#[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
374/// Persist `progress`, atomically, never moving an upgrade in progress back
375/// to an earlier stage - see [`monotonic`].
376///
377/// Written to a sibling `.tmp` and renamed, the same reason
378/// `daemon::write_status_to` does it: `/api/health` reads this file on every
379/// poll and must never see a half-written one.
380pub 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
398/// Record `detail` as the failure of the upgrade request that is asking,
399/// unless a handover is already in flight (`parking` / `restarting`): that
400/// handover is not this request's to fail. The check and the write happen
401/// under the same lock as [`write_progress`], so a handover that records
402/// `parking` in between cannot be overwritten. Returns whether it was written.
403pub 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/// The last upgrade this deck recorded, if it has ever started one.
417#[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/// Where every handover step is appended, beside `upgrade.json`. Written by
424/// this module directly, so no supervisor redirection of stderr can orphan it.
425#[must_use]
426pub fn log_path(home: &Path) -> PathBuf {
427    home.join("upgrade.log")
428}
429
430/// The log is rotated to `upgrade.log.1` (one generation) past this size.
431pub const LOG_MAX_BYTES: u64 = 256 * 1024;
432
433/// A non-terminal `replaced` / `restarting` stage older than this is stuck.
434pub const STALL_AFTER_SECS: i64 = 120;
435
436/// A handover lease with no beat for this long is not alive. Same idiom as
437/// `ask::LEASE_TTL`; no pid is consulted.
438pub const LEASE_TTL_SECS: i64 = 90;
439
440/// How often the watchdog repeats itself for one stage.
441pub const HEARTBEAT_SECS: i64 = 60;
442
443/// How often the watchdog thread looks at `upgrade.json`.
444pub const WATCHDOG_POLL: Duration = Duration::from_secs(30);
445
446/// Append `line` to `path`, rotating to `<path>.1` first when the file has
447/// reached `max` bytes. Open-append-close every time, so a rename or a
448/// deleted file never leaves a stale handle.
449pub 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
466/// One handover step: to `tracing` at INFO and to `<home>/upgrade.log`, with
467/// a UTC timestamp and this process's pid. Best-effort - a failed write is a
468/// warning, never a failed upgrade.
469pub fn log_step(home: &Path, msg: &str) {
470    tracing::info!("handover: {msg}");
471    log_line(home, "INFO", msg);
472}
473
474/// Like [`log_step`], at WARN.
475pub 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
491/// [`write_progress`] that says so when it fails, instead of dropping the
492/// error: a progress file that silently stops moving is the symptom this
493/// module exists to explain.
494pub 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/// Proof that `hand_over` is running: written when it is entered, beaten
501/// while it waits on the loop, removed when it leaves.
502#[derive(Debug, Clone, Serialize, Deserialize)]
503pub struct HandoverLease {
504    /// When `hand_over` was entered.
505    pub entered_at: Timestamp,
506    /// Last time it said it was alive.
507    pub beat_at: Timestamp,
508    /// The run it is waiting on, when one was in flight.
509    #[serde(default)]
510    pub parked_run: Option<String>,
511}
512
513impl HandoverLease {
514    /// Whether the last beat is recent enough at `now`.
515    #[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/// Where the [`HandoverLease`] lives, beside `upgrade.json`.
522#[must_use]
523pub fn lease_path(home: &Path) -> PathBuf {
524    home.join("upgrade.handover.json")
525}
526
527/// The lease on disk, if there is one.
528#[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/// Held by `hand_over` for as long as it runs; removes the lease on drop.
547#[derive(Debug)]
548pub struct LeaseGuard {
549    home: PathBuf,
550    lease: HandoverLease,
551}
552
553impl LeaseGuard {
554    /// Record that `hand_over` has been entered.
555    #[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    /// Enter the handover: write the lease and the `parking` stage as one step
571    /// under the progress lock, so a failing upgrade request ([`fail_progress`])
572    /// or a fresh one ([`write_progress`]) cannot land between the two and leave
573    /// a record that is newer than the lease. `None` for the stage means
574    /// `upgrade.json` was unreadable (the lease is still written).
575    #[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    /// Say it is still alive.
594    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/// How a non-terminal stage came to be called stuck.
607#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
608#[serde(rename_all = "snake_case")]
609pub enum StallKind {
610    /// The handover was signalled but `hand_over` left no record of entering.
611    NeverEntered,
612    /// `hand_over` did enter (or the successor is starting) but nothing has
613    /// moved or beaten for longer than allowed.
614    StoppedBeating,
615}
616
617/// A non-terminal stage that has outlived what it should take.
618#[derive(Debug, Clone, PartialEq, Eq)]
619pub struct Stall {
620    /// The stage that is stuck.
621    pub stage: Stage,
622    /// Which kind of stuck, which decides what the operator is told.
623    pub kind: StallKind,
624    /// How long it has been without progress.
625    pub age_secs: i64,
626    /// What it is waiting on, in words.
627    pub waiting_on: String,
628}
629
630/// Seconds `progress` has been in its stage at `now`. A clock that went
631/// backwards counts as zero, never as a negative age.
632#[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/// The lease, when it proves `hand_over` of *this* upgrade is alive at `now`.
638#[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
649/// What a stage is waiting on, in words.
650fn 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/// Seconds `hand_over` has been alive and waiting, when `lease` proves it.
671#[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/// Pure: whether `progress` is stuck at `now`.
677///
678/// A fresh lease means `hand_over` is alive and waiting on the loop, which is
679/// legitimate for as long as the run's node takes, so it is never stuck.
680#[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    // A stale lease for this upgrade means hand_over was alive and went quiet.
689    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/// What the watchdog wants said this tick.
710#[derive(Debug, Clone, PartialEq, Eq)]
711pub struct Beat {
712    /// The stage this beat is about.
713    pub stage: Stage,
714    /// Whether to say it as a warning (stuck) rather than a heartbeat.
715    pub warn: bool,
716    /// The line to log and to persist as `detail`.
717    pub message: String,
718}
719
720/// Decides when the watchdog speaks. Pure: the caller supplies `now`.
721#[derive(Debug, Default)]
722pub struct Watchdog {
723    last: Option<(Stage, Timestamp)>,
724}
725
726impl Watchdog {
727    /// Look at the record at `now`. `None` means stay quiet: terminal, too
728    /// early, or already spoken within [`HEARTBEAT_SECS`] for this stage.
729    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        // A live `hand_over` is allowed to wait long, but says what it waits
745        // on; the other stages are silent until they are stalled.
746        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/// The watchdog's latest message, kept apart from `upgrade.json` so a
786/// diagnostic write can never race a stage transition.
787#[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(&note)
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/// The watchdog's message for exactly this stage of this upgrade, if any;
816/// a note from another stage or another upgrade is ignored.
817#[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
825/// Run the watchdog on a thread of its own, so a blocked runtime, or a
826/// process half-way through dropping one, still speaks. It never ends; it is
827/// a daemon thread and dies with the process.
828pub 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                // Never rewrites upgrade.json: a read-modify-write here could
848                // overwrite a stage the handover or the successor saved in
849                // between. The note goes to its own file, which only this
850                // thread writes, and is matched to the record when read.
851                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
859/// Reconcile a leftover progress record on startup, before the server starts
860/// answering requests.
861///
862/// A non-terminal record on disk when a process starts can only mean one of
863/// two things: this *is* the successor `spawn_successor` started, or the
864/// predecessor died before finishing the handover (a crash, a reboot, an
865/// operator killing it by hand). Either way waiting longer will not resolve
866/// it - this process is already up - so it is settled immediately:
867/// [`Stage::Done`] when the running version matches what was asked for,
868/// [`Stage::Failed`] otherwise, so the operator is told rather than left
869/// watching a stage that will never move again.
870pub fn reconcile_after_restart(home: &Path) {
871    let Some(mut progress) = read_progress(home) else {
872        return;
873    };
874    // Whoever held it is not this process.
875    let _ = std::fs::remove_file(lease_path(home));
876    if progress.stage.terminal() {
877        return;
878    }
879    let running = env!("CARGO_PKG_VERSION");
880    // `progress.to` is `latest.tag_name` from the forge, which - like every
881    // tag in this repository - carries a `v` prefix `CARGO_PKG_VERSION` does
882    // not. kaishin's own `is_update_available` strips it before comparing;
883    // an exact-string match here would call a successful upgrade `Failed`
884    // every time, because "v0.5.2" is never equal to "0.5.2".
885    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
904/// Spawn the background check for `cfg`, unless it is switched off.
905pub 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
929/// Drain a pending check and print at most one line.
930///
931/// Bounded on purpose: a slow network must never hold up the exit of a command
932/// that already did its work.
933pub 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    /// An interval configured below GitHub's unauthenticated 60 req/hour/IP
959    /// limit must be floored, or `web`'s background recheck loop - which,
960    /// unlike a one-off CLI invocation, keeps polling for as long as `magi
961    /// web` stays up - would repeat the same call far past that limit.
962    #[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        // SAFETY: single-threaded test, no other thread reads the variable.
989        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    /// `mode = "off"` means no checker, for every caller.
1048    ///
1049    /// It used to mean it only for the notify path: `Checker::new` returned
1050    /// `Some` unconditionally, so `POST /api/upgrade` went to the GitHub
1051    /// releases API on a deck configured never to check. A test that reached
1052    /// that route made a live, unauthenticated request, and GitHub's
1053    /// 60-per-hour-per-address limit then turned the suite red on one runner
1054    /// at a time - for as long as anyone kept re-running it, since each
1055    /// attempt spent another request.
1056    #[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    /// `cached_update` never touches the network: a state file written the
1079    /// way `check_and_save` writes one is enough to answer, and no file at
1080    /// all answers "unknown" rather than blocking or erroring.
1081    ///
1082    /// Built from `kaishin::Checker` directly, with an explicit state path,
1083    /// rather than through [`Checker::new`]: that constructor always points
1084    /// at the real cache directory, which is right for production - every
1085    /// `magi` invocation on the machine shares one throttle file - but wrong
1086    /// for a test, which must never read or write the operator's actual
1087    /// state.
1088    #[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        // `to` is `latest.tag_name` as the forge and this repository's own
1116        // tags spell it - with a `v` - which `CARGO_PKG_VERSION` never
1117        // carries. A test that leaves the `v` off would not have caught the
1118        // exact-string-equality bug this function used to have.
1119        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        // Even a record that regressed to `replaced` is read through the lease.
1226        let r = staged(Stage::Replaced, 1000);
1227        assert!(stall(&r, Some(&l), at(1000 + hours + 10)).is_none());
1228        // Once the beat stops, it is stuck, and says so.
1229        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}