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    /// Finished, one way or the other - nothing is still moving.
244    #[must_use]
245    pub fn terminal(self) -> bool {
246        matches!(self, Self::Done | Self::Failed)
247    }
248}
249
250/// One upgrade's progress, persisted at [`progress_path`].
251///
252/// Kept on disk rather than in memory because the process that finishes an
253/// upgrade is never the one that started it: the successor is a fresh binary
254/// (see `web::spawn_successor`), and the only thing the two share is the
255/// disk. Beside the run history rather than under the cache dir alongside
256/// [`state_path`]: this is not throttle bookkeeping, it is the record of one
257/// upgrade the operator asked for, and - like a parked run - it is meant to
258/// outlive the process that wrote it.
259#[derive(Debug, Clone, Serialize, Deserialize)]
260pub struct Progress {
261    /// Where this upgrade has gotten to.
262    pub stage: Stage,
263    /// Version this upgrade started from.
264    pub from: String,
265    /// Version it is replacing itself with.
266    pub to: Option<String>,
267    /// The run [`Stage::Parking`] is waiting on, when one was in flight.
268    #[serde(default)]
269    pub parked_run: Option<String>,
270    /// When this upgrade was asked for.
271    pub started_at: Timestamp,
272    /// Last time `stage` changed.
273    pub updated_at: Timestamp,
274    /// Why [`Stage::Failed`] happened; `None` for every other stage.
275    #[serde(default)]
276    pub detail: Option<String>,
277}
278
279impl Progress {
280    /// A fresh record for an upgrade that is about to replace the binary.
281    #[must_use]
282    pub fn new(from: String, to: String) -> Self {
283        let now = Timestamp::now();
284        Self {
285            stage: Stage::Downloading,
286            from,
287            to: Some(to),
288            parked_run: None,
289            started_at: now,
290            updated_at: now,
291            detail: None,
292        }
293    }
294
295    /// Move to `stage`, stamping when it changed.
296    pub fn advance(&mut self, stage: Stage) {
297        self.stage = stage;
298        self.updated_at = Timestamp::now();
299        // A note belongs to the stage it was written for.
300        self.detail = None;
301    }
302
303    /// Stop at [`Stage::Failed`], with a reason a human can read.
304    pub fn fail(&mut self, detail: impl Into<String>) {
305        self.stage = Stage::Failed;
306        self.updated_at = Timestamp::now();
307        self.detail = Some(detail.into());
308    }
309}
310
311/// Where [`Progress`] is recorded: beside `daemon.json`, not under the cache
312/// dir - see [`Progress`]'s own doc for why the two are not the same place.
313#[must_use]
314pub fn progress_path(home: &Path) -> PathBuf {
315    home.join("upgrade.json")
316}
317
318/// Persist `progress`, atomically.
319///
320/// Written to a sibling `.tmp` and renamed, the same reason
321/// `daemon::write_status_to` does it: `/api/health` reads this file on every
322/// poll and must never see a half-written one.
323pub fn write_progress(home: &Path, progress: &Progress) -> Result<()> {
324    let path = progress_path(home);
325    if let Some(parent) = path.parent() {
326        std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
327    }
328    let body = serde_json::to_string_pretty(progress).context("serialize upgrade progress")?;
329    let tmp = path.with_extension("json.tmp");
330    std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
331    std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
332    Ok(())
333}
334
335/// The last upgrade this deck recorded, if it has ever started one.
336#[must_use]
337pub fn read_progress(home: &Path) -> Option<Progress> {
338    let body = std::fs::read_to_string(progress_path(home)).ok()?;
339    serde_json::from_str(&body).ok()
340}
341
342/// Where every handover step is appended, beside `upgrade.json`. Written by
343/// this module directly, so no supervisor redirection of stderr can orphan it.
344#[must_use]
345pub fn log_path(home: &Path) -> PathBuf {
346    home.join("upgrade.log")
347}
348
349/// The log is rotated to `upgrade.log.1` (one generation) past this size.
350pub const LOG_MAX_BYTES: u64 = 256 * 1024;
351
352/// A non-terminal `replaced` / `restarting` stage older than this is stuck.
353pub const STALL_AFTER_SECS: i64 = 120;
354
355/// `parking` waits for the node in flight, up to an implement wave (an hour by
356/// default), so it is only called stuck past that plus a margin.
357pub const PARKING_STALL_AFTER_SECS: i64 = 70 * 60;
358
359/// How often the watchdog repeats itself for one stage.
360pub const HEARTBEAT_SECS: i64 = 60;
361
362/// How often the watchdog thread looks at `upgrade.json`.
363pub const WATCHDOG_POLL: Duration = Duration::from_secs(30);
364
365/// Append `line` to `path`, rotating to `<path>.1` first when the file has
366/// reached `max` bytes. Open-append-close every time, so a rename or a
367/// deleted file never leaves a stale handle.
368pub fn append_bounded(path: &Path, line: &str, max: u64) -> std::io::Result<()> {
369    use std::io::Write as _;
370    if std::fs::metadata(path).is_ok_and(|m| m.len() >= max) {
371        let mut old = path.as_os_str().to_owned();
372        old.push(".1");
373        std::fs::rename(path, PathBuf::from(old))?;
374    }
375    if let Some(parent) = path.parent() {
376        std::fs::create_dir_all(parent)?;
377    }
378    let mut file = std::fs::OpenOptions::new()
379        .create(true)
380        .append(true)
381        .open(path)?;
382    writeln!(file, "{line}")
383}
384
385/// One handover step: to `tracing` at INFO and to `<home>/upgrade.log`, with
386/// a UTC timestamp and this process's pid. Best-effort - a failed write is a
387/// warning, never a failed upgrade.
388pub fn log_step(home: &Path, msg: &str) {
389    tracing::info!("handover: {msg}");
390    log_line(home, "INFO", msg);
391}
392
393/// Like [`log_step`], at WARN.
394pub fn log_warn(home: &Path, msg: &str) {
395    tracing::warn!("handover: {msg}");
396    log_line(home, "WARN", msg);
397}
398
399fn log_line(home: &Path, level: &str, msg: &str) {
400    let line = format!(
401        "{} pid={} {level} {msg}",
402        Timestamp::now(),
403        std::process::id()
404    );
405    if let Err(e) = append_bounded(&log_path(home), &line, LOG_MAX_BYTES) {
406        tracing::warn!("could not append to {}: {e}", log_path(home).display());
407    }
408}
409
410/// [`write_progress`] that says so when it fails, instead of dropping the
411/// error: a progress file that silently stops moving is the symptom this
412/// module exists to explain.
413pub fn write_progress_logged(home: &Path, progress: &Progress) {
414    if let Err(e) = write_progress(home, progress) {
415        log_warn(home, &format!("could not write upgrade.json: {e:#}"));
416    }
417}
418
419/// A non-terminal stage that has outlived what it should take.
420#[derive(Debug, Clone, PartialEq, Eq)]
421pub struct Stall {
422    /// The stage that is stuck.
423    pub stage: Stage,
424    /// How long it has been in that stage.
425    pub age_secs: i64,
426    /// What it is waiting on, in words.
427    pub waiting_on: String,
428}
429
430/// Seconds `progress` has been in its stage at `now`. A clock that went
431/// backwards counts as zero, never as a negative age.
432#[must_use]
433pub fn stage_age_secs(progress: &Progress, now: Timestamp) -> i64 {
434    (now.as_second() - progress.updated_at.as_second()).max(0)
435}
436
437/// What a stage is waiting on, in words.
438fn waiting_on(progress: &Progress) -> String {
439    match progress.stage {
440        Stage::Replaced => "serve() observing the handover signal and calling hand_over \
441                            (hand_over has not recorded `parking`)"
442            .to_owned(),
443        Stage::Parking => match &progress.parked_run {
444            Some(run) => format!("the loop to finish run {run} at its next node boundary"),
445            None => "the loop to stop (no run was recorded as in flight)".to_owned(),
446        },
447        Stage::Restarting => "spawn_successor returning and this process exiting".to_owned(),
448        Stage::Downloading => "the release download and binary replacement".to_owned(),
449        Stage::Done | Stage::Failed => String::new(),
450    }
451}
452
453/// Pure: whether `progress` is stuck at `now`.
454#[must_use]
455pub fn stall(progress: &Progress, now: Timestamp) -> Option<Stall> {
456    let limit = match progress.stage {
457        Stage::Replaced | Stage::Restarting => STALL_AFTER_SECS,
458        Stage::Parking => PARKING_STALL_AFTER_SECS,
459        Stage::Downloading | Stage::Done | Stage::Failed => return None,
460    };
461    let age_secs = stage_age_secs(progress, now);
462    (age_secs > limit).then(|| Stall {
463        stage: progress.stage,
464        age_secs,
465        waiting_on: waiting_on(progress),
466    })
467}
468
469/// What the watchdog wants said this tick.
470#[derive(Debug, Clone, PartialEq, Eq)]
471pub struct Beat {
472    /// The stage this beat is about.
473    pub stage: Stage,
474    /// Whether to say it as a warning (stuck) rather than a heartbeat.
475    pub warn: bool,
476    /// The line to log and to persist as `detail`.
477    pub message: String,
478}
479
480/// Decides when the watchdog speaks. Pure: the caller supplies `now`.
481#[derive(Debug, Default)]
482pub struct Watchdog {
483    last: Option<(Stage, Timestamp)>,
484}
485
486impl Watchdog {
487    /// Look at the record at `now`. `None` means stay quiet: terminal, too
488    /// early, or already spoken within [`HEARTBEAT_SECS`] for this stage.
489    pub fn tick(&mut self, progress: &Progress, now: Timestamp) -> Option<Beat> {
490        if progress.stage.terminal() {
491            self.last = None;
492            return None;
493        }
494        if self.last.is_some_and(|(stage, _)| stage != progress.stage) {
495            self.last = None;
496        }
497        let stalled = stall(progress, now);
498        // `parking` is allowed to be long, but says what it waits on; the
499        // other stages are silent until they are stalled.
500        if stalled.is_none() && progress.stage != Stage::Parking {
501            return None;
502        }
503        if let Some((_, at)) = self.last
504            && now.as_second() - at.as_second() < HEARTBEAT_SECS
505        {
506            return None;
507        }
508        self.last = Some((progress.stage, now));
509        let age = stage_age_secs(progress, now);
510        let (warn, message) = match &stalled {
511            Some(s) => (
512                true,
513                format!(
514                    "stuck in {:?} for {} min {} s, waiting on {}",
515                    s.stage,
516                    s.age_secs / 60,
517                    s.age_secs % 60,
518                    s.waiting_on
519                ),
520            ),
521            None => (
522                false,
523                format!(
524                    "parking for {} min {} s, waiting on {}",
525                    age / 60,
526                    age % 60,
527                    waiting_on(progress)
528                ),
529            ),
530        };
531        Some(Beat {
532            stage: progress.stage,
533            warn,
534            message,
535        })
536    }
537}
538
539/// The watchdog's latest message, kept apart from `upgrade.json` so a
540/// diagnostic write can never race a stage transition.
541#[derive(Debug, Serialize, Deserialize)]
542struct Note {
543    stage: Stage,
544    stage_since: Timestamp,
545    message: String,
546}
547
548fn note_path(home: &Path) -> PathBuf {
549    home.join("upgrade.note.json")
550}
551
552fn write_note(home: &Path, progress: &Progress, message: &str) {
553    let note = Note {
554        stage: progress.stage,
555        stage_since: progress.updated_at,
556        message: message.to_owned(),
557    };
558    let path = note_path(home);
559    let tmp = path.with_extension("json.tmp");
560    let written = serde_json::to_string(&note)
561        .map_err(std::io::Error::other)
562        .and_then(|body| std::fs::write(&tmp, body))
563        .and_then(|()| std::fs::rename(&tmp, &path));
564    if let Err(e) = written {
565        tracing::warn!("could not write {}: {e}", path.display());
566    }
567}
568
569/// The watchdog's message for exactly this stage of this upgrade, if any;
570/// a note from another stage or another upgrade is ignored.
571#[must_use]
572pub fn read_note(home: &Path, progress: &Progress) -> Option<String> {
573    let body = std::fs::read_to_string(note_path(home)).ok()?;
574    let note: Note = serde_json::from_str(&body).ok()?;
575    (note.stage == progress.stage && note.stage_since == progress.updated_at)
576        .then_some(note.message)
577}
578
579/// Run the watchdog on a thread of its own, so a blocked runtime, or a
580/// process half-way through dropping one, still speaks. It never ends; it is
581/// a daemon thread and dies with the process.
582pub fn spawn_watchdog(home: PathBuf) {
583    let spawned = std::thread::Builder::new()
584        .name("upgrade-watchdog".to_owned())
585        .spawn(move || {
586            let mut dog = Watchdog::default();
587            loop {
588                std::thread::sleep(WATCHDOG_POLL);
589                let Some(progress) = read_progress(&home) else {
590                    continue;
591                };
592                let Some(beat) = dog.tick(&progress, Timestamp::now()) else {
593                    continue;
594                };
595                if beat.warn {
596                    log_warn(&home, &beat.message);
597                } else {
598                    log_step(&home, &beat.message);
599                }
600                // Never rewrites upgrade.json: a read-modify-write here could
601                // overwrite a stage the handover or the successor saved in
602                // between. The note goes to its own file, which only this
603                // thread writes, and is matched to the record when read.
604                write_note(&home, &progress, &beat.message);
605            }
606        });
607    if let Err(e) = spawned {
608        tracing::warn!("could not start the upgrade watchdog: {e}");
609    }
610}
611
612/// Reconcile a leftover progress record on startup, before the server starts
613/// answering requests.
614///
615/// A non-terminal record on disk when a process starts can only mean one of
616/// two things: this *is* the successor `spawn_successor` started, or the
617/// predecessor died before finishing the handover (a crash, a reboot, an
618/// operator killing it by hand). Either way waiting longer will not resolve
619/// it - this process is already up - so it is settled immediately:
620/// [`Stage::Done`] when the running version matches what was asked for,
621/// [`Stage::Failed`] otherwise, so the operator is told rather than left
622/// watching a stage that will never move again.
623pub fn reconcile_after_restart(home: &Path) {
624    let Some(mut progress) = read_progress(home) else {
625        return;
626    };
627    if progress.stage.terminal() {
628        return;
629    }
630    let running = env!("CARGO_PKG_VERSION");
631    // `progress.to` is `latest.tag_name` from the forge, which - like every
632    // tag in this repository - carries a `v` prefix `CARGO_PKG_VERSION` does
633    // not. kaishin's own `is_update_available` strips it before comparing;
634    // an exact-string match here would call a successful upgrade `Failed`
635    // every time, because "v0.5.2" is never equal to "0.5.2".
636    if progress
637        .to
638        .as_deref()
639        .is_some_and(|to| to.trim_start_matches('v') == running)
640    {
641        progress.advance(Stage::Done);
642    } else {
643        let to = progress
644            .to
645            .clone()
646            .unwrap_or_else(|| "the expected release".to_owned());
647        progress.fail(format!(
648            "this process came up on {running}, not {to} - the upgrade may \
649             not have replaced the binary"
650        ));
651    }
652    write_progress_logged(home, &progress);
653}
654
655/// Spawn the background check for `cfg`, unless it is switched off.
656pub fn spawn(cfg: &Update, rt: &tokio::runtime::Handle) -> Option<Pending> {
657    if disabled_by_env() || cfg.mode == UpdateMode::Off {
658        return None;
659    }
660    let checker = Checker::new(cfg)?;
661    match cfg.mode {
662        UpdateMode::Off => None,
663        UpdateMode::Notify => {
664            if !checker.should_check() {
665                let latest = checker.cached_update()?;
666                return Some(Pending::Cached { checker, latest });
667            }
668            let inner = checker.inner.clone();
669            let handle = rt.spawn(async move { inner.check_and_save().await });
670            Some(Pending::Notify { checker, handle })
671        }
672        UpdateMode::Install => {
673            let inner = checker.inner.clone();
674            let handle = rt.spawn(async move { inner.auto_update().await });
675            Some(Pending::Install { handle })
676        }
677    }
678}
679
680/// Drain a pending check and print at most one line.
681///
682/// Bounded on purpose: a slow network must never hold up the exit of a command
683/// that already did its work.
684pub async fn finalize(pending: Option<Pending>, budget: Duration) {
685    let Some(pending) = pending else {
686        return;
687    };
688    match pending {
689        Pending::Cached { checker, latest } => {
690            eprintln!("{}", checker.format_banner(&latest));
691        }
692        Pending::Notify { checker, handle } => {
693            if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
694                eprintln!("{}", checker.format_banner(&latest));
695            }
696        }
697        Pending::Install { handle } => {
698            if let Ok(Ok(Ok(Some(latest)))) = tokio::time::timeout(budget, handle).await {
699                eprintln!("magi updated itself to {}", latest.tag_name);
700            }
701        }
702    }
703}
704
705#[cfg(test)]
706mod tests {
707    use super::*;
708
709    /// An interval configured below GitHub's unauthenticated 60 req/hour/IP
710    /// limit must be floored, or `web`'s background recheck loop - which,
711    /// unlike a one-off CLI invocation, keeps polling for as long as `magi
712    /// web` stays up - would repeat the same call far past that limit.
713    #[test]
714    fn effective_interval_floors_a_configured_interval_below_githubs_rate_limit() {
715        let cfg = Update {
716            mode: UpdateMode::Notify,
717            interval: Some("1s".to_owned()),
718        };
719        assert_eq!(
720            effective_interval(&cfg),
721            MIN_INTERVAL,
722            "an interval that would exceed GitHub's rate limit under continuous \
723             polling must be floored rather than honoured verbatim"
724        );
725
726        let sane = Update {
727            mode: UpdateMode::Notify,
728            interval: Some("2h".to_owned()),
729        };
730        assert_eq!(
731            effective_interval(&sane),
732            Duration::from_secs(2 * 60 * 60),
733            "an interval already above the floor must pass through unchanged"
734        );
735    }
736
737    #[test]
738    fn env_kill_switch_semantics() {
739        // SAFETY: single-threaded test, no other thread reads the variable.
740        unsafe {
741            std::env::remove_var(NO_AUTOUPDATE_ENV);
742        }
743        assert!(!disabled_by_env());
744        for (value, disabled) in [
745            ("1", true),
746            ("true", true),
747            ("yes", true),
748            ("0", false),
749            ("false", false),
750            ("FALSE", false),
751            ("", false),
752            ("  ", false),
753        ] {
754            unsafe {
755                std::env::set_var(NO_AUTOUPDATE_ENV, value);
756            }
757            assert_eq!(
758                disabled_by_env(),
759                disabled,
760                "MAGI_NO_AUTOUPDATE={value:?} should {} disable",
761                if disabled { "" } else { "not" }
762            );
763        }
764        unsafe {
765            std::env::remove_var(NO_AUTOUPDATE_ENV);
766        }
767    }
768
769    #[test]
770    fn off_mode_never_spawns() {
771        let rt = tokio::runtime::Builder::new_current_thread()
772            .enable_all()
773            .build()
774            .unwrap();
775        let cfg = Update {
776            mode: UpdateMode::Off,
777            interval: None,
778        };
779        assert!(spawn(&cfg, rt.handle()).is_none());
780    }
781
782    #[test]
783    fn state_path_lives_under_the_cache_dir() {
784        let path = state_path().expect("a cache dir on every supported platform");
785        assert!(path.ends_with("magi/last_update_check.json"));
786        let data = dirs::data_local_dir().unwrap_or_default();
787        assert!(
788            !path.starts_with(&data) || dirs::cache_dir() == dirs::data_local_dir(),
789            "throttle state must not sit in the run history directory"
790        );
791    }
792
793    #[tokio::test]
794    async fn finalize_of_nothing_is_a_no_op() {
795        finalize(None, Duration::from_millis(1)).await;
796    }
797
798    /// `mode = "off"` means no checker, for every caller.
799    ///
800    /// It used to mean it only for the notify path: `Checker::new` returned
801    /// `Some` unconditionally, so `POST /api/upgrade` went to the GitHub
802    /// releases API on a deck configured never to check. A test that reached
803    /// that route made a live, unauthenticated request, and GitHub's
804    /// 60-per-hour-per-address limit then turned the suite red on one runner
805    /// at a time - for as long as anyone kept re-running it, since each
806    /// attempt spent another request.
807    #[test]
808    fn checking_is_off_for_every_caller_when_the_config_says_off() {
809        assert!(
810            Checker::new(&Update {
811                mode: UpdateMode::Off,
812                interval: None,
813            })
814            .is_none(),
815            "an operator who writes mode = \"off\" means it"
816        );
817        for mode in [UpdateMode::Notify, UpdateMode::Install] {
818            assert!(
819                Checker::new(&Update {
820                    mode,
821                    interval: None,
822                })
823                .is_some(),
824                "{mode:?} still asks the forge"
825            );
826        }
827    }
828
829    /// `cached_update` never touches the network: a state file written the
830    /// way `check_and_save` writes one is enough to answer, and no file at
831    /// all answers "unknown" rather than blocking or erroring.
832    ///
833    /// Built from `kaishin::Checker` directly, with an explicit state path,
834    /// rather than through [`Checker::new`]: that constructor always points
835    /// at the real cache directory, which is right for production - every
836    /// `magi` invocation on the machine shares one throttle file - but wrong
837    /// for a test, which must never read or write the operator's actual
838    /// state.
839    #[test]
840    fn cached_update_answers_from_disk_with_no_network_call() {
841        let dir = tempfile::tempdir().expect("temp dir");
842        let path = dir.path().join("state.json");
843        let opts = kaishin::KaishinOptions::new("yukimemi", "magi", "magi", "0.1.0");
844        let checker = Checker {
845            inner: kaishin::Checker::new("magi", opts).state_path(path.clone()),
846        };
847
848        assert!(
849            checker.cached_update().is_none(),
850            "no state file yet must read as \"unknown\", not an error"
851        );
852
853        let state = kaishin::UpdateCheckState {
854            last_checked_unix: 0,
855            last_known_latest: Some("v9.9.9".to_owned()),
856            last_known_url: Some("https://example.invalid/9.9.9".to_owned()),
857        };
858        kaishin::save_check_state(&path, &state).expect("seed the state file");
859
860        let latest = checker.cached_update().expect("a newer release was cached");
861        assert_eq!(latest.tag_name, "v9.9.9");
862    }
863
864    #[test]
865    fn reconcile_after_restart_confirms_a_matching_version() {
866        // `to` is `latest.tag_name` as the forge and this repository's own
867        // tags spell it - with a `v` - which `CARGO_PKG_VERSION` never
868        // carries. A test that leaves the `v` off would not have caught the
869        // exact-string-equality bug this function used to have.
870        let home = tempfile::tempdir().expect("temp home");
871        let mut progress = Progress::new(
872            "0.1.0".to_owned(),
873            format!("v{}", env!("CARGO_PKG_VERSION")),
874        );
875        progress.advance(Stage::Restarting);
876        write_progress(home.path(), &progress).expect("seed progress");
877
878        reconcile_after_restart(home.path());
879
880        let after = read_progress(home.path()).expect("progress on disk");
881        assert_eq!(
882            after.stage,
883            Stage::Done,
884            "the successor is running exactly the release that was asked for, \
885             `v` prefix and all"
886        );
887    }
888
889    #[test]
890    fn reconcile_after_restart_flags_a_mismatched_version() {
891        let home = tempfile::tempdir().expect("temp home");
892        let mut progress = Progress::new("0.1.0".to_owned(), "v9.9.9".to_owned());
893        progress.advance(Stage::Restarting);
894        write_progress(home.path(), &progress).expect("seed progress");
895
896        reconcile_after_restart(home.path());
897
898        let after = read_progress(home.path()).expect("progress on disk");
899        assert_eq!(after.stage, Stage::Failed);
900        assert!(
901            after.detail.is_some_and(|d| d.contains("9.9.9")),
902            "the operator needs to know which release it did not come back on"
903        );
904    }
905
906    #[test]
907    fn reconcile_after_restart_leaves_a_settled_record_alone() {
908        let home = tempfile::tempdir().expect("temp home");
909        let mut progress = Progress::new("0.1.0".to_owned(), "9.9.9".to_owned());
910        progress.advance(Stage::Done);
911        write_progress(home.path(), &progress).expect("seed progress");
912
913        reconcile_after_restart(home.path());
914
915        let after = read_progress(home.path()).expect("progress on disk");
916        assert_eq!(
917            after.stage,
918            Stage::Done,
919            "an already-settled record must not be rewritten by a later, unrelated start"
920        );
921    }
922
923    #[test]
924    fn reconcile_after_restart_with_nothing_on_disk_is_a_quiet_no_op() {
925        let home = tempfile::tempdir().expect("temp home");
926        reconcile_after_restart(home.path());
927        assert!(read_progress(home.path()).is_none());
928    }
929
930    fn at(secs: i64) -> Timestamp {
931        Timestamp::from_second(secs).expect("timestamp")
932    }
933
934    fn staged(stage: Stage, since: i64) -> Progress {
935        let mut p = Progress::new("0.1.0".to_owned(), "v0.2.0".to_owned());
936        p.stage = stage;
937        p.updated_at = at(since);
938        p
939    }
940
941    #[test]
942    fn stall_has_a_threshold_per_stage_and_is_silent_when_terminal() {
943        let p = staged(Stage::Replaced, 1000);
944        assert!(stall(&p, at(1000 + STALL_AFTER_SECS)).is_none());
945        let s = stall(&p, at(1000 + STALL_AFTER_SECS + 1)).expect("stalled");
946        assert_eq!(s.stage, Stage::Replaced);
947        assert_eq!(s.age_secs, STALL_AFTER_SECS + 1);
948        assert!(s.waiting_on.contains("hand_over"), "{}", s.waiting_on);
949
950        let p = staged(Stage::Restarting, 1000);
951        assert!(stall(&p, at(1000 + STALL_AFTER_SECS + 1)).is_some());
952
953        let p = staged(Stage::Parking, 1000);
954        assert!(
955            stall(&p, at(1000 + 3600)).is_none(),
956            "an hour of parking is normal"
957        );
958        assert!(stall(&p, at(1000 + PARKING_STALL_AFTER_SECS + 1)).is_some());
959
960        for stage in [Stage::Done, Stage::Failed, Stage::Downloading] {
961            assert!(stall(&staged(stage, 0), at(1_000_000)).is_none());
962        }
963    }
964
965    #[test]
966    fn a_clock_that_went_backwards_is_age_zero() {
967        let p = staged(Stage::Replaced, 5000);
968        assert_eq!(stage_age_secs(&p, at(100)), 0);
969        assert!(stall(&p, at(100)).is_none());
970    }
971
972    #[test]
973    fn a_note_is_kept_beside_the_record_and_matches_only_its_own_stage() {
974        let home = tempfile::tempdir().expect("temp home");
975        let p = staged(Stage::Replaced, 1000);
976        write_progress(home.path(), &p).expect("write");
977        write_note(home.path(), &p, "stuck");
978        assert_eq!(read_note(home.path(), &p).as_deref(), Some("stuck"));
979        assert_eq!(read_progress(home.path()).unwrap().updated_at, at(1000));
980        assert!(read_note(home.path(), &staged(Stage::Parking, 1000)).is_none());
981        assert!(read_note(home.path(), &staged(Stage::Replaced, 2000)).is_none());
982    }
983
984    #[test]
985    fn the_watchdog_speaks_once_a_minute_and_resets_on_a_new_stage() {
986        let mut dog = Watchdog::default();
987        let p = staged(Stage::Replaced, 1000);
988        assert!(dog.tick(&p, at(1060)).is_none(), "not stalled yet");
989        let beat = dog.tick(&p, at(1200)).expect("stalled");
990        assert!(beat.warn);
991        assert!(dog.tick(&p, at(1230)).is_none(), "spoke 30 s ago");
992        assert!(dog.tick(&p, at(1260)).is_some(), "a minute later");
993
994        let parking = staged(Stage::Parking, 1260);
995        let beat = dog
996            .tick(&parking, at(1270))
997            .expect("parking heartbeat at once");
998        assert!(!beat.warn, "a short park is not a warning");
999        assert!(dog.tick(&parking, at(1300)).is_none());
1000        let late = dog
1001            .tick(&parking, at(1260 + PARKING_STALL_AFTER_SECS + 1))
1002            .expect("past the parking ceiling");
1003        assert!(late.warn);
1004
1005        let done = staged(Stage::Done, 0);
1006        assert!(dog.tick(&done, at(9_999_999)).is_none());
1007    }
1008
1009    #[test]
1010    fn the_upgrade_log_appends_and_rotates_to_one_generation() {
1011        let dir = tempfile::tempdir().expect("temp dir");
1012        let path = dir.path().join("upgrade.log");
1013        append_bounded(&path, "one", 16).expect("append");
1014        append_bounded(&path, "two", 16).expect("append");
1015        assert_eq!(std::fs::read_to_string(&path).unwrap(), "one\ntwo\n");
1016        append_bounded(&path, "three-and-more", 16).expect("append");
1017        append_bounded(&path, "four", 16).expect("append");
1018        assert_eq!(std::fs::read_to_string(&path).unwrap(), "four\n");
1019        let old = dir.path().join("upgrade.log.1");
1020        assert!(
1021            std::fs::read_to_string(old)
1022                .unwrap()
1023                .contains("three-and-more")
1024        );
1025    }
1026}