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
263/// One upgrade's progress, persisted at [`progress_path`].
264///
265/// Kept on disk rather than in memory because the process that finishes an
266/// upgrade is never the one that started it: the successor is a fresh binary
267/// (see `web::spawn_successor`), and the only thing the two share is the
268/// disk. Beside the run history rather than under the cache dir alongside
269/// [`state_path`]: this is not throttle bookkeeping, it is the record of one
270/// upgrade the operator asked for, and - like a parked run - it is meant to
271/// outlive the process that wrote it.
272#[derive(Debug, Clone, Serialize, Deserialize)]
273pub struct Progress {
274    /// Where this upgrade has gotten to.
275    pub stage: Stage,
276    /// Version this upgrade started from.
277    pub from: String,
278    /// Version it is replacing itself with.
279    pub to: Option<String>,
280    /// The run [`Stage::Parking`] is waiting on, when one was in flight.
281    #[serde(default)]
282    pub parked_run: Option<String>,
283    /// When this upgrade was asked for.
284    pub started_at: Timestamp,
285    /// Last time `stage` changed.
286    pub updated_at: Timestamp,
287    /// Why [`Stage::Failed`] happened; `None` for every other stage.
288    #[serde(default)]
289    pub detail: Option<String>,
290}
291
292impl Progress {
293    /// A fresh record for an upgrade that is about to replace the binary.
294    #[must_use]
295    pub fn new(from: String, to: String) -> Self {
296        let now = Timestamp::now();
297        Self {
298            stage: Stage::Downloading,
299            from,
300            to: Some(to),
301            parked_run: None,
302            started_at: now,
303            updated_at: now,
304            detail: None,
305        }
306    }
307
308    /// Move to `stage`, stamping when it changed.
309    pub fn advance(&mut self, stage: Stage) {
310        self.stage = stage;
311        self.updated_at = Timestamp::now();
312        // A note belongs to the stage it was written for.
313        self.detail = None;
314    }
315
316    /// Stop at [`Stage::Failed`], with a reason a human can read.
317    pub fn fail(&mut self, detail: impl Into<String>) {
318        self.stage = Stage::Failed;
319        self.updated_at = Timestamp::now();
320        self.detail = Some(detail.into());
321    }
322}
323
324/// Where [`Progress`] is recorded: beside `daemon.json`, not under the cache
325/// dir - see [`Progress`]'s own doc for why the two are not the same place.
326#[must_use]
327pub fn progress_path(home: &Path) -> PathBuf {
328    home.join("upgrade.json")
329}
330
331/// Serialises the read-compare-rename of [`write_progress`], so two writers
332/// in this process cannot each compare against the same stale record.
333static WRITE_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
334
335/// What to persist when `candidate` is written over `current`.
336///
337/// An upgrade in progress never goes backwards: a candidate whose stage is
338/// not past the recorded one (a second `POST /api/upgrade` landing on a live
339/// handover wrote `replaced` over `parking`) keeps the recorded stage and
340/// timestamps and only refreshes `to`. A terminal record is over, so whatever
341/// comes next starts a new upgrade, and a candidate that ends the upgrade
342/// (`Done` / `Failed`) always goes through.
343#[must_use]
344pub fn monotonic(current: Option<&Progress>, candidate: &Progress) -> Progress {
345    match current {
346        Some(cur)
347            if !cur.stage.terminal()
348                && !candidate.stage.terminal()
349                && candidate.stage.rank() <= cur.stage.rank() =>
350        {
351            let mut kept = cur.clone();
352            if candidate.to.is_some() {
353                kept.to.clone_from(&candidate.to);
354            }
355            kept
356        }
357        _ => candidate.clone(),
358    }
359}
360
361/// Persist `progress`, atomically, never moving an upgrade in progress back
362/// to an earlier stage - see [`monotonic`].
363///
364/// Written to a sibling `.tmp` and renamed, the same reason
365/// `daemon::write_status_to` does it: `/api/health` reads this file on every
366/// poll and must never see a half-written one.
367pub fn write_progress(home: &Path, progress: &Progress) -> Result<()> {
368    let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
369    write_progress_locked(home, progress)
370}
371
372fn write_progress_locked(home: &Path, progress: &Progress) -> Result<()> {
373    let path = progress_path(home);
374    if let Some(parent) = path.parent() {
375        std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
376    }
377    let to_write = monotonic(read_progress(home).as_ref(), progress);
378    let body = serde_json::to_string_pretty(&to_write).context("serialize upgrade progress")?;
379    let tmp = path.with_extension("json.tmp");
380    std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
381    std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
382    Ok(())
383}
384
385/// Record `detail` as the failure of the upgrade request that is asking,
386/// unless a handover is already in flight (`parking` / `restarting`): that
387/// handover is not this request's to fail. The check and the write happen
388/// under the same lock as [`write_progress`], so a handover that records
389/// `parking` in between cannot be overwritten. Returns whether it was written.
390pub fn fail_progress(home: &Path, detail: &str) -> Result<bool> {
391    let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
392    let Some(mut progress) = read_progress(home) else {
393        return Ok(false);
394    };
395    if matches!(progress.stage, Stage::Parking | Stage::Restarting) {
396        return Ok(false);
397    }
398    progress.fail(detail);
399    write_progress_locked(home, &progress)?;
400    Ok(true)
401}
402
403/// The last upgrade this deck recorded, if it has ever started one.
404#[must_use]
405pub fn read_progress(home: &Path) -> Option<Progress> {
406    let body = std::fs::read_to_string(progress_path(home)).ok()?;
407    serde_json::from_str(&body).ok()
408}
409
410/// Where every handover step is appended, beside `upgrade.json`. Written by
411/// this module directly, so no supervisor redirection of stderr can orphan it.
412#[must_use]
413pub fn log_path(home: &Path) -> PathBuf {
414    home.join("upgrade.log")
415}
416
417/// The log is rotated to `upgrade.log.1` (one generation) past this size.
418pub const LOG_MAX_BYTES: u64 = 256 * 1024;
419
420/// A non-terminal `replaced` / `restarting` stage older than this is stuck.
421pub const STALL_AFTER_SECS: i64 = 120;
422
423/// A handover lease with no beat for this long is not alive. Same idiom as
424/// `ask::LEASE_TTL`; no pid is consulted.
425pub const LEASE_TTL_SECS: i64 = 90;
426
427/// How often the watchdog repeats itself for one stage.
428pub const HEARTBEAT_SECS: i64 = 60;
429
430/// How often the watchdog thread looks at `upgrade.json`.
431pub const WATCHDOG_POLL: Duration = Duration::from_secs(30);
432
433/// Append `line` to `path`, rotating to `<path>.1` first when the file has
434/// reached `max` bytes. Open-append-close every time, so a rename or a
435/// deleted file never leaves a stale handle.
436pub fn append_bounded(path: &Path, line: &str, max: u64) -> std::io::Result<()> {
437    use std::io::Write as _;
438    if std::fs::metadata(path).is_ok_and(|m| m.len() >= max) {
439        let mut old = path.as_os_str().to_owned();
440        old.push(".1");
441        std::fs::rename(path, PathBuf::from(old))?;
442    }
443    if let Some(parent) = path.parent() {
444        std::fs::create_dir_all(parent)?;
445    }
446    let mut file = std::fs::OpenOptions::new()
447        .create(true)
448        .append(true)
449        .open(path)?;
450    writeln!(file, "{line}")
451}
452
453/// One handover step: to `tracing` at INFO and to `<home>/upgrade.log`, with
454/// a UTC timestamp and this process's pid. Best-effort - a failed write is a
455/// warning, never a failed upgrade.
456pub fn log_step(home: &Path, msg: &str) {
457    tracing::info!("handover: {msg}");
458    log_line(home, "INFO", msg);
459}
460
461/// Like [`log_step`], at WARN.
462pub fn log_warn(home: &Path, msg: &str) {
463    tracing::warn!("handover: {msg}");
464    log_line(home, "WARN", msg);
465}
466
467fn log_line(home: &Path, level: &str, msg: &str) {
468    let line = format!(
469        "{} pid={} {level} {msg}",
470        Timestamp::now(),
471        std::process::id()
472    );
473    if let Err(e) = append_bounded(&log_path(home), &line, LOG_MAX_BYTES) {
474        tracing::warn!("could not append to {}: {e}", log_path(home).display());
475    }
476}
477
478/// [`write_progress`] that says so when it fails, instead of dropping the
479/// error: a progress file that silently stops moving is the symptom this
480/// module exists to explain.
481pub fn write_progress_logged(home: &Path, progress: &Progress) {
482    if let Err(e) = write_progress(home, progress) {
483        log_warn(home, &format!("could not write upgrade.json: {e:#}"));
484    }
485}
486
487/// Proof that `hand_over` is running: written when it is entered, beaten
488/// while it waits on the loop, removed when it leaves.
489#[derive(Debug, Clone, Serialize, Deserialize)]
490pub struct HandoverLease {
491    /// When `hand_over` was entered.
492    pub entered_at: Timestamp,
493    /// Last time it said it was alive.
494    pub beat_at: Timestamp,
495    /// The run it is waiting on, when one was in flight.
496    #[serde(default)]
497    pub parked_run: Option<String>,
498}
499
500impl HandoverLease {
501    /// Whether the last beat is recent enough at `now`.
502    #[must_use]
503    pub fn fresh(&self, now: Timestamp) -> bool {
504        now.as_second() - self.beat_at.as_second() <= LEASE_TTL_SECS
505    }
506}
507
508/// Where the [`HandoverLease`] lives, beside `upgrade.json`.
509#[must_use]
510pub fn lease_path(home: &Path) -> PathBuf {
511    home.join("upgrade.handover.json")
512}
513
514/// The lease on disk, if there is one.
515#[must_use]
516pub fn read_lease(home: &Path) -> Option<HandoverLease> {
517    let body = std::fs::read_to_string(lease_path(home)).ok()?;
518    serde_json::from_str(&body).ok()
519}
520
521fn write_lease(home: &Path, lease: &HandoverLease) {
522    let path = lease_path(home);
523    let tmp = path.with_extension("json.tmp");
524    let written = serde_json::to_string(lease)
525        .map_err(std::io::Error::other)
526        .and_then(|body| std::fs::write(&tmp, body))
527        .and_then(|()| std::fs::rename(&tmp, &path));
528    if let Err(e) = written {
529        log_warn(home, &format!("could not write the handover lease: {e}"));
530    }
531}
532
533/// Held by `hand_over` for as long as it runs; removes the lease on drop.
534#[derive(Debug)]
535pub struct LeaseGuard {
536    home: PathBuf,
537    lease: HandoverLease,
538}
539
540impl LeaseGuard {
541    /// Record that `hand_over` has been entered.
542    #[must_use]
543    pub fn enter(home: &Path, parked_run: Option<String>) -> Self {
544        let now = Timestamp::now();
545        let lease = HandoverLease {
546            entered_at: now,
547            beat_at: now,
548            parked_run,
549        };
550        write_lease(home, &lease);
551        Self {
552            home: home.to_owned(),
553            lease,
554        }
555    }
556
557    /// Enter the handover: write the lease and the `parking` stage as one step
558    /// under the progress lock, so a failing upgrade request ([`fail_progress`])
559    /// or a fresh one ([`write_progress`]) cannot land between the two and leave
560    /// a record that is newer than the lease. `None` for the stage means
561    /// `upgrade.json` was unreadable (the lease is still written).
562    #[must_use]
563    pub fn enter_parking(home: &Path) -> (Self, bool) {
564        let _guard = WRITE_LOCK.lock().unwrap_or_else(|e| e.into_inner());
565        let progress = read_progress(home);
566        let this = Self::enter(home, progress.as_ref().and_then(|p| p.parked_run.clone()));
567        let recorded = match progress {
568            Some(mut p) => {
569                p.advance(Stage::Parking);
570                if let Err(e) = write_progress_locked(home, &p) {
571                    log_warn(home, &format!("could not write upgrade.json: {e:#}"));
572                }
573                true
574            }
575            None => false,
576        };
577        (this, recorded)
578    }
579
580    /// Say it is still alive.
581    pub fn beat(&mut self) {
582        self.lease.beat_at = Timestamp::now();
583        write_lease(&self.home, &self.lease);
584    }
585}
586
587impl Drop for LeaseGuard {
588    fn drop(&mut self) {
589        let _ = std::fs::remove_file(lease_path(&self.home));
590    }
591}
592
593/// How a non-terminal stage came to be called stuck.
594#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
595#[serde(rename_all = "snake_case")]
596pub enum StallKind {
597    /// The handover was signalled but `hand_over` left no record of entering.
598    NeverEntered,
599    /// `hand_over` did enter (or the successor is starting) but nothing has
600    /// moved or beaten for longer than allowed.
601    StoppedBeating,
602}
603
604/// A non-terminal stage that has outlived what it should take.
605#[derive(Debug, Clone, PartialEq, Eq)]
606pub struct Stall {
607    /// The stage that is stuck.
608    pub stage: Stage,
609    /// Which kind of stuck, which decides what the operator is told.
610    pub kind: StallKind,
611    /// How long it has been without progress.
612    pub age_secs: i64,
613    /// What it is waiting on, in words.
614    pub waiting_on: String,
615}
616
617/// Seconds `progress` has been in its stage at `now`. A clock that went
618/// backwards counts as zero, never as a negative age.
619#[must_use]
620pub fn stage_age_secs(progress: &Progress, now: Timestamp) -> i64 {
621    (now.as_second() - progress.updated_at.as_second()).max(0)
622}
623
624/// The lease, when it proves `hand_over` of *this* upgrade is alive at `now`.
625#[must_use]
626pub fn live_lease<'a>(
627    progress: &Progress,
628    lease: Option<&'a HandoverLease>,
629    now: Timestamp,
630) -> Option<&'a HandoverLease> {
631    lease.filter(|l| {
632        !progress.stage.terminal() && l.fresh(now) && l.entered_at >= progress.started_at
633    })
634}
635
636/// What a stage is waiting on, in words.
637fn waiting_on(progress: &Progress, lease: Option<&HandoverLease>) -> String {
638    match progress.stage {
639        Stage::Replaced => "hand_over to start (HANDOVER was signalled; hand_over has left no \
640                            record that it was entered)"
641            .to_owned(),
642        Stage::Parking => {
643            let run = lease
644                .and_then(|l| l.parked_run.as_ref())
645                .or(progress.parked_run.as_ref());
646            match run {
647                Some(run) => format!("the loop to finish run {run} at its next node boundary"),
648                None => "the loop to stop (no run was recorded as in flight)".to_owned(),
649            }
650        }
651        Stage::Restarting => "spawn_successor returning and this process exiting".to_owned(),
652        Stage::Downloading => "the release download and binary replacement".to_owned(),
653        Stage::Done | Stage::Failed => String::new(),
654    }
655}
656
657/// Seconds `hand_over` has been alive and waiting, when `lease` proves it.
658#[must_use]
659pub fn waited_secs(lease: &HandoverLease, now: Timestamp) -> i64 {
660    (now.as_second() - lease.entered_at.as_second()).max(0)
661}
662
663/// Pure: whether `progress` is stuck at `now`.
664///
665/// A fresh lease means `hand_over` is alive and waiting on the loop, which is
666/// legitimate for as long as the run's node takes, so it is never stuck.
667#[must_use]
668pub fn stall(progress: &Progress, lease: Option<&HandoverLease>, now: Timestamp) -> Option<Stall> {
669    if progress.stage.terminal() || progress.stage == Stage::Downloading {
670        return None;
671    }
672    if live_lease(progress, lease, now).is_some() {
673        return None;
674    }
675    // A stale lease for this upgrade means hand_over was alive and went quiet.
676    let stale = lease.filter(|l| l.entered_at >= progress.started_at);
677    let (kind, since) = match (progress.stage, stale) {
678        (_, Some(l)) => (StallKind::StoppedBeating, l.beat_at),
679        (Stage::Replaced, None) => (StallKind::NeverEntered, progress.updated_at),
680        _ => (StallKind::StoppedBeating, progress.updated_at),
681    };
682    let age_secs = (now.as_second() - since.as_second()).max(0);
683    let limit = if stale.is_some() {
684        LEASE_TTL_SECS
685    } else {
686        STALL_AFTER_SECS
687    };
688    (age_secs > limit).then(|| Stall {
689        stage: progress.stage,
690        kind,
691        age_secs,
692        waiting_on: waiting_on(progress, lease),
693    })
694}
695
696/// What the watchdog wants said this tick.
697#[derive(Debug, Clone, PartialEq, Eq)]
698pub struct Beat {
699    /// The stage this beat is about.
700    pub stage: Stage,
701    /// Whether to say it as a warning (stuck) rather than a heartbeat.
702    pub warn: bool,
703    /// The line to log and to persist as `detail`.
704    pub message: String,
705}
706
707/// Decides when the watchdog speaks. Pure: the caller supplies `now`.
708#[derive(Debug, Default)]
709pub struct Watchdog {
710    last: Option<(Stage, Timestamp)>,
711}
712
713impl Watchdog {
714    /// Look at the record at `now`. `None` means stay quiet: terminal, too
715    /// early, or already spoken within [`HEARTBEAT_SECS`] for this stage.
716    pub fn tick(
717        &mut self,
718        progress: &Progress,
719        lease: Option<&HandoverLease>,
720        now: Timestamp,
721    ) -> Option<Beat> {
722        if progress.stage.terminal() {
723            self.last = None;
724            return None;
725        }
726        if self.last.is_some_and(|(stage, _)| stage != progress.stage) {
727            self.last = None;
728        }
729        let stalled = stall(progress, lease, now);
730        let alive = live_lease(progress, lease, now);
731        // A live `hand_over` is allowed to wait long, but says what it waits
732        // on; the other stages are silent until they are stalled.
733        if stalled.is_none() && alive.is_none() && progress.stage != Stage::Parking {
734            return None;
735        }
736        if let Some((_, at)) = self.last
737            && now.as_second() - at.as_second() < HEARTBEAT_SECS
738        {
739            return None;
740        }
741        self.last = Some((progress.stage, now));
742        let age = alive.map_or_else(|| stage_age_secs(progress, now), |l| waited_secs(l, now));
743        let (warn, message) = match &stalled {
744            Some(s) => (
745                true,
746                format!(
747                    "stuck in {:?} for {} min {} s without progress, waiting on {}",
748                    s.stage,
749                    s.age_secs / 60,
750                    s.age_secs % 60,
751                    s.waiting_on
752                ),
753            ),
754            None => (
755                false,
756                format!(
757                    "parking for {} min {} s, waiting on {}",
758                    age / 60,
759                    age % 60,
760                    waiting_on(progress, lease)
761                ),
762            ),
763        };
764        Some(Beat {
765            stage: progress.stage,
766            warn,
767            message,
768        })
769    }
770}
771
772/// The watchdog's latest message, kept apart from `upgrade.json` so a
773/// diagnostic write can never race a stage transition.
774#[derive(Debug, Serialize, Deserialize)]
775struct Note {
776    stage: Stage,
777    stage_since: Timestamp,
778    message: String,
779}
780
781fn note_path(home: &Path) -> PathBuf {
782    home.join("upgrade.note.json")
783}
784
785fn write_note(home: &Path, progress: &Progress, message: &str) {
786    let note = Note {
787        stage: progress.stage,
788        stage_since: progress.updated_at,
789        message: message.to_owned(),
790    };
791    let path = note_path(home);
792    let tmp = path.with_extension("json.tmp");
793    let written = serde_json::to_string(&note)
794        .map_err(std::io::Error::other)
795        .and_then(|body| std::fs::write(&tmp, body))
796        .and_then(|()| std::fs::rename(&tmp, &path));
797    if let Err(e) = written {
798        tracing::warn!("could not write {}: {e}", path.display());
799    }
800}
801
802/// The watchdog's message for exactly this stage of this upgrade, if any;
803/// a note from another stage or another upgrade is ignored.
804#[must_use]
805pub fn read_note(home: &Path, progress: &Progress) -> Option<String> {
806    let body = std::fs::read_to_string(note_path(home)).ok()?;
807    let note: Note = serde_json::from_str(&body).ok()?;
808    (note.stage == progress.stage && note.stage_since == progress.updated_at)
809        .then_some(note.message)
810}
811
812/// Run the watchdog on a thread of its own, so a blocked runtime, or a
813/// process half-way through dropping one, still speaks. It never ends; it is
814/// a daemon thread and dies with the process.
815pub fn spawn_watchdog(home: PathBuf) {
816    let spawned = std::thread::Builder::new()
817        .name("upgrade-watchdog".to_owned())
818        .spawn(move || {
819            let mut dog = Watchdog::default();
820            loop {
821                std::thread::sleep(WATCHDOG_POLL);
822                let Some(progress) = read_progress(&home) else {
823                    continue;
824                };
825                let lease = read_lease(&home);
826                let Some(beat) = dog.tick(&progress, lease.as_ref(), Timestamp::now()) else {
827                    continue;
828                };
829                if beat.warn {
830                    log_warn(&home, &beat.message);
831                } else {
832                    log_step(&home, &beat.message);
833                }
834                // Never rewrites upgrade.json: a read-modify-write here could
835                // overwrite a stage the handover or the successor saved in
836                // between. The note goes to its own file, which only this
837                // thread writes, and is matched to the record when read.
838                write_note(&home, &progress, &beat.message);
839            }
840        });
841    if let Err(e) = spawned {
842        tracing::warn!("could not start the upgrade watchdog: {e}");
843    }
844}
845
846/// Reconcile a leftover progress record on startup, before the server starts
847/// answering requests.
848///
849/// A non-terminal record on disk when a process starts can only mean one of
850/// two things: this *is* the successor `spawn_successor` started, or the
851/// predecessor died before finishing the handover (a crash, a reboot, an
852/// operator killing it by hand). Either way waiting longer will not resolve
853/// it - this process is already up - so it is settled immediately:
854/// [`Stage::Done`] when the running version matches what was asked for,
855/// [`Stage::Failed`] otherwise, so the operator is told rather than left
856/// watching a stage that will never move again.
857pub fn reconcile_after_restart(home: &Path) {
858    let Some(mut progress) = read_progress(home) else {
859        return;
860    };
861    // Whoever held it is not this process.
862    let _ = std::fs::remove_file(lease_path(home));
863    if progress.stage.terminal() {
864        return;
865    }
866    let running = env!("CARGO_PKG_VERSION");
867    // `progress.to` is `latest.tag_name` from the forge, which - like every
868    // tag in this repository - carries a `v` prefix `CARGO_PKG_VERSION` does
869    // not. kaishin's own `is_update_available` strips it before comparing;
870    // an exact-string match here would call a successful upgrade `Failed`
871    // every time, because "v0.5.2" is never equal to "0.5.2".
872    if progress
873        .to
874        .as_deref()
875        .is_some_and(|to| to.trim_start_matches('v') == running)
876    {
877        progress.advance(Stage::Done);
878    } else {
879        let to = progress
880            .to
881            .clone()
882            .unwrap_or_else(|| "the expected release".to_owned());
883        progress.fail(format!(
884            "this process came up on {running}, not {to} - the upgrade may \
885             not have replaced the binary"
886        ));
887    }
888    write_progress_logged(home, &progress);
889}
890
891/// Spawn the background check for `cfg`, unless it is switched off.
892pub fn spawn(cfg: &Update, rt: &tokio::runtime::Handle) -> Option<Pending> {
893    if disabled_by_env() || cfg.mode == UpdateMode::Off {
894        return None;
895    }
896    let checker = Checker::new(cfg)?;
897    match cfg.mode {
898        UpdateMode::Off => None,
899        UpdateMode::Notify => {
900            if !checker.should_check() {
901                let latest = checker.cached_update()?;
902                return Some(Pending::Cached { checker, latest });
903            }
904            let inner = checker.inner.clone();
905            let handle = rt.spawn(async move { inner.check_and_save().await });
906            Some(Pending::Notify { checker, handle })
907        }
908        UpdateMode::Install => {
909            let inner = checker.inner.clone();
910            let handle = rt.spawn(async move { inner.auto_update().await });
911            Some(Pending::Install { handle })
912        }
913    }
914}
915
916/// Drain a pending check and print at most one line.
917///
918/// Bounded on purpose: a slow network must never hold up the exit of a command
919/// that already did its work.
920pub async fn finalize(pending: Option<Pending>, budget: Duration) {
921    let Some(pending) = pending else {
922        return;
923    };
924    match pending {
925        Pending::Cached { checker, latest } => {
926            eprintln!("{}", checker.format_banner(&latest));
927        }
928        Pending::Notify { checker, handle } => {
929            if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
930                eprintln!("{}", checker.format_banner(&latest));
931            }
932        }
933        Pending::Install { handle } => {
934            if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
935                eprintln!("magi updated itself to {}", latest.tag_name);
936            }
937        }
938    }
939}
940
941#[cfg(test)]
942mod tests {
943    use super::*;
944
945    /// An interval configured below GitHub's unauthenticated 60 req/hour/IP
946    /// limit must be floored, or `web`'s background recheck loop - which,
947    /// unlike a one-off CLI invocation, keeps polling for as long as `magi
948    /// web` stays up - would repeat the same call far past that limit.
949    #[test]
950    fn effective_interval_floors_a_configured_interval_below_githubs_rate_limit() {
951        let cfg = Update {
952            mode: UpdateMode::Notify,
953            interval: Some("1s".to_owned()),
954        };
955        assert_eq!(
956            effective_interval(&cfg),
957            MIN_INTERVAL,
958            "an interval that would exceed GitHub's rate limit under continuous \
959             polling must be floored rather than honoured verbatim"
960        );
961
962        let sane = Update {
963            mode: UpdateMode::Notify,
964            interval: Some("2h".to_owned()),
965        };
966        assert_eq!(
967            effective_interval(&sane),
968            Duration::from_secs(2 * 60 * 60),
969            "an interval already above the floor must pass through unchanged"
970        );
971    }
972
973    #[test]
974    fn env_kill_switch_semantics() {
975        // SAFETY: single-threaded test, no other thread reads the variable.
976        unsafe {
977            std::env::remove_var(NO_AUTOUPDATE_ENV);
978        }
979        assert!(!disabled_by_env());
980        for (value, disabled) in [
981            ("1", true),
982            ("true", true),
983            ("yes", true),
984            ("0", false),
985            ("false", false),
986            ("FALSE", false),
987            ("", false),
988            ("  ", false),
989        ] {
990            unsafe {
991                std::env::set_var(NO_AUTOUPDATE_ENV, value);
992            }
993            assert_eq!(
994                disabled_by_env(),
995                disabled,
996                "MAGI_NO_AUTOUPDATE={value:?} should {} disable",
997                if disabled { "" } else { "not" }
998            );
999        }
1000        unsafe {
1001            std::env::remove_var(NO_AUTOUPDATE_ENV);
1002        }
1003    }
1004
1005    #[test]
1006    fn off_mode_never_spawns() {
1007        let rt = tokio::runtime::Builder::new_current_thread()
1008            .enable_all()
1009            .build()
1010            .unwrap();
1011        let cfg = Update {
1012            mode: UpdateMode::Off,
1013            interval: None,
1014        };
1015        assert!(spawn(&cfg, rt.handle()).is_none());
1016    }
1017
1018    #[test]
1019    fn state_path_lives_under_the_cache_dir() {
1020        let path = state_path().expect("a cache dir on every supported platform");
1021        assert!(path.ends_with("magi/last_update_check.json"));
1022        let data = dirs::data_local_dir().unwrap_or_default();
1023        assert!(
1024            !path.starts_with(&data) || dirs::cache_dir() == dirs::data_local_dir(),
1025            "throttle state must not sit in the run history directory"
1026        );
1027    }
1028
1029    #[tokio::test]
1030    async fn finalize_of_nothing_is_a_no_op() {
1031        finalize(None, Duration::from_millis(1)).await;
1032    }
1033
1034    /// `mode = "off"` means no checker, for every caller.
1035    ///
1036    /// It used to mean it only for the notify path: `Checker::new` returned
1037    /// `Some` unconditionally, so `POST /api/upgrade` went to the GitHub
1038    /// releases API on a deck configured never to check. A test that reached
1039    /// that route made a live, unauthenticated request, and GitHub's
1040    /// 60-per-hour-per-address limit then turned the suite red on one runner
1041    /// at a time - for as long as anyone kept re-running it, since each
1042    /// attempt spent another request.
1043    #[test]
1044    fn checking_is_off_for_every_caller_when_the_config_says_off() {
1045        assert!(
1046            Checker::new(&Update {
1047                mode: UpdateMode::Off,
1048                interval: None,
1049            })
1050            .is_none(),
1051            "an operator who writes mode = \"off\" means it"
1052        );
1053        for mode in [UpdateMode::Notify, UpdateMode::Install] {
1054            assert!(
1055                Checker::new(&Update {
1056                    mode,
1057                    interval: None,
1058                })
1059                .is_some(),
1060                "{mode:?} still asks the forge"
1061            );
1062        }
1063    }
1064
1065    /// `cached_update` never touches the network: a state file written the
1066    /// way `check_and_save` writes one is enough to answer, and no file at
1067    /// all answers "unknown" rather than blocking or erroring.
1068    ///
1069    /// Built from `kaishin::Checker` directly, with an explicit state path,
1070    /// rather than through [`Checker::new`]: that constructor always points
1071    /// at the real cache directory, which is right for production - every
1072    /// `magi` invocation on the machine shares one throttle file - but wrong
1073    /// for a test, which must never read or write the operator's actual
1074    /// state.
1075    #[test]
1076    fn cached_update_answers_from_disk_with_no_network_call() {
1077        let dir = tempfile::tempdir().expect("temp dir");
1078        let path = dir.path().join("state.json");
1079        let opts = kaishin::KaishinOptions::new("yukimemi", "magi", "magi", "0.1.0");
1080        let checker = Checker {
1081            inner: kaishin::Checker::new("magi", opts).state_path(path.clone()),
1082        };
1083
1084        assert!(
1085            checker.cached_update().is_none(),
1086            "no state file yet must read as \"unknown\", not an error"
1087        );
1088
1089        let state = kaishin::UpdateCheckState {
1090            last_checked_unix: 0,
1091            last_known_latest: Some("v9.9.9".to_owned()),
1092            last_known_url: Some("https://example.invalid/9.9.9".to_owned()),
1093        };
1094        kaishin::save_check_state(&path, &state).expect("seed the state file");
1095
1096        let latest = checker.cached_update().expect("a newer release was cached");
1097        assert_eq!(latest.tag_name, "v9.9.9");
1098    }
1099
1100    #[test]
1101    fn reconcile_after_restart_confirms_a_matching_version() {
1102        // `to` is `latest.tag_name` as the forge and this repository's own
1103        // tags spell it - with a `v` - which `CARGO_PKG_VERSION` never
1104        // carries. A test that leaves the `v` off would not have caught the
1105        // exact-string-equality bug this function used to have.
1106        let home = tempfile::tempdir().expect("temp home");
1107        let mut progress = Progress::new(
1108            "0.1.0".to_owned(),
1109            format!("v{}", env!("CARGO_PKG_VERSION")),
1110        );
1111        progress.advance(Stage::Restarting);
1112        write_progress(home.path(), &progress).expect("seed progress");
1113
1114        reconcile_after_restart(home.path());
1115
1116        let after = read_progress(home.path()).expect("progress on disk");
1117        assert_eq!(
1118            after.stage,
1119            Stage::Done,
1120            "the successor is running exactly the release that was asked for, \
1121             `v` prefix and all"
1122        );
1123    }
1124
1125    #[test]
1126    fn reconcile_after_restart_flags_a_mismatched_version() {
1127        let home = tempfile::tempdir().expect("temp home");
1128        let mut progress = Progress::new("0.1.0".to_owned(), "v9.9.9".to_owned());
1129        progress.advance(Stage::Restarting);
1130        write_progress(home.path(), &progress).expect("seed progress");
1131
1132        reconcile_after_restart(home.path());
1133
1134        let after = read_progress(home.path()).expect("progress on disk");
1135        assert_eq!(after.stage, Stage::Failed);
1136        assert!(
1137            after.detail.is_some_and(|d| d.contains("9.9.9")),
1138            "the operator needs to know which release it did not come back on"
1139        );
1140    }
1141
1142    #[test]
1143    fn reconcile_after_restart_leaves_a_settled_record_alone() {
1144        let home = tempfile::tempdir().expect("temp home");
1145        let mut progress = Progress::new("0.1.0".to_owned(), "9.9.9".to_owned());
1146        progress.advance(Stage::Done);
1147        write_progress(home.path(), &progress).expect("seed progress");
1148
1149        reconcile_after_restart(home.path());
1150
1151        let after = read_progress(home.path()).expect("progress on disk");
1152        assert_eq!(
1153            after.stage,
1154            Stage::Done,
1155            "an already-settled record must not be rewritten by a later, unrelated start"
1156        );
1157    }
1158
1159    #[test]
1160    fn reconcile_after_restart_with_nothing_on_disk_is_a_quiet_no_op() {
1161        let home = tempfile::tempdir().expect("temp home");
1162        reconcile_after_restart(home.path());
1163        assert!(read_progress(home.path()).is_none());
1164    }
1165
1166    fn at(secs: i64) -> Timestamp {
1167        Timestamp::from_second(secs).expect("timestamp")
1168    }
1169
1170    fn staged(stage: Stage, since: i64) -> Progress {
1171        let mut p = Progress::new("0.1.0".to_owned(), "v0.2.0".to_owned());
1172        p.stage = stage;
1173        p.started_at = at(since);
1174        p.updated_at = at(since);
1175        p
1176    }
1177
1178    fn lease(entered: i64, beat: i64) -> HandoverLease {
1179        HandoverLease {
1180            entered_at: at(entered),
1181            beat_at: at(beat),
1182            parked_run: Some("r1".to_owned()),
1183        }
1184    }
1185
1186    #[test]
1187    fn a_handover_never_entered_is_stuck_and_says_only_what_is_known() {
1188        let p = staged(Stage::Replaced, 1000);
1189        assert!(stall(&p, None, at(1000 + STALL_AFTER_SECS)).is_none());
1190        let s = stall(&p, None, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
1191        assert_eq!(s.stage, Stage::Replaced);
1192        assert_eq!(s.kind, StallKind::NeverEntered);
1193        assert_eq!(s.age_secs, STALL_AFTER_SECS + 1);
1194        assert!(s.waiting_on.contains("hand_over"), "{}", s.waiting_on);
1195
1196        let p = staged(Stage::Restarting, 1000);
1197        let s = stall(&p, None, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
1198        assert_eq!(s.kind, StallKind::StoppedBeating);
1199
1200        for stage in [Stage::Done, Stage::Failed, Stage::Downloading] {
1201            assert!(stall(&staged(stage, 0), None, at(1_000_000)).is_none());
1202        }
1203    }
1204
1205    #[test]
1206    fn a_live_parking_wait_is_never_stuck_however_long_it_lasts() {
1207        let mut p = staged(Stage::Parking, 1000);
1208        p.started_at = at(900);
1209        let hours = 5 * 3600;
1210        let l = lease(1000, 1000 + hours);
1211        assert!(stall(&p, Some(&l), at(1000 + hours + 10)).is_none());
1212        // Even a record that regressed to `replaced` is read through the lease.
1213        let r = staged(Stage::Replaced, 1000);
1214        assert!(stall(&r, Some(&l), at(1000 + hours + 10)).is_none());
1215        // Once the beat stops, it is stuck, and says so.
1216        let s = stall(&p, Some(&l), at(1000 + hours + LEASE_TTL_SECS + 1)).expect("stuck");
1217        assert_eq!(s.kind, StallKind::StoppedBeating);
1218    }
1219
1220    #[test]
1221    fn a_lease_from_an_earlier_upgrade_proves_nothing() {
1222        let p = staged(Stage::Replaced, 2000);
1223        let old = lease(10, 3000);
1224        assert!(live_lease(&p, Some(&old), at(3001)).is_none());
1225    }
1226
1227    #[test]
1228    fn a_stage_never_goes_backwards_but_a_new_upgrade_after_a_terminal_one_starts() {
1229        let parking = staged(Stage::Parking, 1000);
1230        for back in [Stage::Replaced, Stage::Downloading, Stage::Parking] {
1231            let mut cand = staged(back, 5000);
1232            cand.to = Some("v9.9.9".to_owned());
1233            let kept = monotonic(Some(&parking), &cand);
1234            assert_eq!(kept.stage, Stage::Parking);
1235            assert_eq!(kept.updated_at, at(1000));
1236            assert_eq!(kept.started_at, parking.started_at);
1237            assert_eq!(kept.to.as_deref(), Some("v9.9.9"), "data is refreshed");
1238        }
1239        assert_eq!(
1240            monotonic(Some(&parking), &staged(Stage::Restarting, 5000)).stage,
1241            Stage::Restarting
1242        );
1243        assert_eq!(
1244            monotonic(Some(&parking), &staged(Stage::Failed, 5000)).stage,
1245            Stage::Failed
1246        );
1247        let done = staged(Stage::Done, 1000);
1248        assert_eq!(
1249            monotonic(Some(&done), &staged(Stage::Downloading, 5000)).stage,
1250            Stage::Downloading
1251        );
1252    }
1253
1254    #[test]
1255    fn write_progress_refuses_a_regression_on_disk() {
1256        let home = tempfile::tempdir().expect("temp home");
1257        write_progress(home.path(), &staged(Stage::Parking, 1000)).expect("write");
1258        write_progress(home.path(), &staged(Stage::Replaced, 5000)).expect("write");
1259        let on_disk = read_progress(home.path()).expect("record");
1260        assert_eq!(on_disk.stage, Stage::Parking);
1261        assert_eq!(on_disk.updated_at, at(1000));
1262    }
1263
1264    #[test]
1265    fn a_failed_request_cannot_overwrite_a_live_handover() {
1266        let home = tempfile::tempdir().expect("temp home");
1267        write_progress(home.path(), &staged(Stage::Parking, 1000)).expect("write");
1268        assert!(!fail_progress(home.path(), "boom").expect("fail"));
1269        assert_eq!(read_progress(home.path()).unwrap().stage, Stage::Parking);
1270        write_progress(home.path(), &staged(Stage::Replaced, 1000)).ok();
1271        let fresh = tempfile::tempdir().expect("temp home");
1272        write_progress(fresh.path(), &staged(Stage::Downloading, 1000)).expect("write");
1273        assert!(fail_progress(fresh.path(), "boom").expect("fail"));
1274        assert_eq!(read_progress(fresh.path()).unwrap().stage, Stage::Failed);
1275    }
1276
1277    #[test]
1278    fn entering_parking_is_one_step_that_keeps_the_lease_newer_than_the_record() {
1279        let home = tempfile::tempdir().expect("temp home");
1280        write_progress(home.path(), &staged(Stage::Replaced, 1000)).expect("write");
1281        let (guard, recorded) = LeaseGuard::enter_parking(home.path());
1282        assert!(recorded);
1283        let p = read_progress(home.path()).expect("record");
1284        assert_eq!(p.stage, Stage::Parking);
1285        let l = read_lease(home.path()).expect("lease");
1286        assert!(l.entered_at >= p.started_at);
1287        assert!(!fail_progress(home.path(), "boom").expect("fail"));
1288        drop(guard);
1289    }
1290
1291    #[test]
1292    fn the_lease_guard_writes_beats_and_removes_the_lease() {
1293        let home = tempfile::tempdir().expect("temp home");
1294        {
1295            let mut guard = LeaseGuard::enter(home.path(), Some("r1".to_owned()));
1296            let first = read_lease(home.path()).expect("lease");
1297            assert_eq!(first.parked_run.as_deref(), Some("r1"));
1298            guard.beat();
1299            assert!(read_lease(home.path()).is_some());
1300        }
1301        assert!(read_lease(home.path()).is_none());
1302    }
1303
1304    #[test]
1305    fn a_clock_that_went_backwards_is_age_zero() {
1306        let p = staged(Stage::Replaced, 5000);
1307        assert_eq!(stage_age_secs(&p, at(100)), 0);
1308        assert!(stall(&p, None, at(100)).is_none());
1309    }
1310
1311    #[test]
1312    fn a_note_is_kept_beside_the_record_and_matches_only_its_own_stage() {
1313        let home = tempfile::tempdir().expect("temp home");
1314        let p = staged(Stage::Replaced, 1000);
1315        write_progress(home.path(), &p).expect("write");
1316        write_note(home.path(), &p, "stuck");
1317        assert_eq!(read_note(home.path(), &p).as_deref(), Some("stuck"));
1318        assert_eq!(read_progress(home.path()).unwrap().updated_at, at(1000));
1319        assert!(read_note(home.path(), &staged(Stage::Parking, 1000)).is_none());
1320        assert!(read_note(home.path(), &staged(Stage::Replaced, 2000)).is_none());
1321    }
1322
1323    #[test]
1324    fn the_watchdog_speaks_once_a_minute_and_resets_on_a_new_stage() {
1325        let mut dog = Watchdog::default();
1326        let p = staged(Stage::Replaced, 1000);
1327        assert!(dog.tick(&p, None, at(1060)).is_none(), "not stalled yet");
1328        let beat = dog.tick(&p, None, at(1200)).expect("stalled");
1329        assert!(beat.warn);
1330        assert!(dog.tick(&p, None, at(1230)).is_none(), "spoke 30 s ago");
1331        assert!(dog.tick(&p, None, at(1260)).is_some(), "a minute later");
1332
1333        let parking = staged(Stage::Parking, 1260);
1334        let l = lease(1260, 1270);
1335        let beat = dog
1336            .tick(&parking, Some(&l), at(1275))
1337            .expect("parking heartbeat at once");
1338        assert!(!beat.warn, "a live wait is not a warning");
1339        assert!(dog.tick(&parking, Some(&l), at(1300)).is_none());
1340        let l = lease(1260, 1260 + 4 * 3600);
1341        let later = dog
1342            .tick(&parking, Some(&l), at(1260 + 4 * 3600 + 5))
1343            .expect("heartbeat");
1344        assert!(!later.warn, "hours of waiting on a run is still not stuck");
1345        assert!(later.message.contains("r1"), "{}", later.message);
1346
1347        let done = staged(Stage::Done, 0);
1348        assert!(dog.tick(&done, None, at(9_999_999)).is_none());
1349    }
1350
1351    #[test]
1352    fn the_upgrade_log_appends_and_rotates_to_one_generation() {
1353        let dir = tempfile::tempdir().expect("temp dir");
1354        let path = dir.path().join("upgrade.log");
1355        append_bounded(&path, "one", 16).expect("append");
1356        append_bounded(&path, "two", 16).expect("append");
1357        assert_eq!(std::fs::read_to_string(&path).unwrap(), "one\ntwo\n");
1358        append_bounded(&path, "three-and-more", 16).expect("append");
1359        append_bounded(&path, "four", 16).expect("append");
1360        assert_eq!(std::fs::read_to_string(&path).unwrap(), "four\n");
1361        let old = dir.path().join("upgrade.log.1");
1362        assert!(
1363            std::fs::read_to_string(old)
1364                .unwrap()
1365                .contains("three-and-more")
1366        );
1367    }
1368}