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