use std::sync::Arc;
use std::time::{Duration, Instant};
use clap::Args;
use serde_json::json;
use crate::cli::output::OutputConfig;
use crate::cli::ui::{self, Caps, Mode};
use crate::config;
use crate::error::OlError;
use crate::install_state;
use crate::update::{self, ApplyStage, CheckResult, Severity};
#[derive(Args, Clone, Debug, Default)]
pub struct UpdateArgs {
#[arg(long)]
pub check: bool,
#[arg(long, short = 'y')]
pub yes: bool,
#[arg(long, hide = true)]
pub apply: bool,
#[arg(long = "force-cargo")]
pub force_cargo: bool,
}
impl UpdateArgs {
fn applies(&self) -> bool {
self.apply || (self.yes && !self.check)
}
}
struct Ctx {
current: String,
registry: String,
port: u16,
egress: crate::egress::EgressConfig,
auto: AutoUpdate,
}
impl Ctx {
fn load() -> Self {
let cfg = config::Config::load(None, None, false).ok();
Self {
current: env!("CARGO_PKG_VERSION").to_string(),
registry: cfg
.as_ref()
.map(|c| c.update.registry_origin.clone())
.unwrap_or_else(|| "https://registry.npmjs.org".to_string()),
port: cfg.as_ref().map(|c| c.port).unwrap_or(7443),
egress: cfg
.as_ref()
.map(|c| c.egress.clone())
.unwrap_or_else(crate::egress::EgressConfig::direct),
auto: AutoUpdate {
on: cfg.as_ref().is_none_or(|c| c.update.auto_update),
by_env: std::env::var_os("OPENLATCH_AUTO_UPDATE").is_some(),
last_check: install_state::InstallState::load_or_default().last_check_at,
},
}
}
}
pub fn run(args: &UpdateArgs, output: &OutputConfig) -> i32 {
let runtime = match tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
{
Ok(rt) => rt,
Err(e) => {
output.print_info(&format!("failed to build local tokio runtime: {e}"));
return 1;
}
};
runtime.block_on(async {
let ctx = Ctx::load();
match Caps::detect(output) {
Some(caps) => run_rail(args, output, &ctx, caps).await,
None => run_json(args, output, &ctx).await,
}
})
}
const POLL: Duration = Duration::from_millis(250);
const TICK: Duration = Duration::from_millis(100);
const BACK_WITHIN: Duration = Duration::from_secs(60);
const STALL: Duration = Duration::from_secs(120);
const DRAIN: Duration = Duration::from_secs(5);
const STAGES: usize = 6;
#[derive(Debug, Clone)]
struct Release {
from: String,
to: String,
severity: Severity,
size: Option<u64>,
}
async fn run_rail(args: &UpdateArgs, output: &OutputConfig, ctx: &Ctx, caps: Caps) -> i32 {
let applies = args.applies();
if applies && refuses_cargo(args) {
let _guard = ui::install(output, "update", 1);
ui::intro("Updating OpenLatch");
return stop(&Stop::new(cargo_refusal(CARGO_SUGGESTION), 5));
}
let plain = caps.mode == Mode::Plain;
let asks = !applies && !args.check && !plain && crate::cli::prompt::interactive(output, false);
let mut guard = if plain {
None
} else {
open(output, STAGES + usize::from(asks))
};
let check = update::check(&ctx.current, &ctx.registry, &ctx.egress).await;
let available = matches!(check, CheckResult::Available { .. });
let daemon = if available && (applies || asks) {
probe_daemon(ctx.port).await
} else {
DaemonState::NotRunning
};
if guard.is_none() {
let total = match (available && applies, &daemon) {
(false, _) => 1,
(true, DaemonState::RunningAndReachable { .. }) => 6,
(true, _) => 4,
};
guard = open(output, total);
}
let _guard = guard;
let now = chrono::Utc::now();
let release = match close_check(check, &ctx.current, &ctx.registry) {
Ok(release) => release,
Err(Outcome::UpToDate { version, .. }) => {
let outcome = up_to_date(version, IDENTITY, ctx.port).await;
return finish(&outcome, &ctx.auto, now);
}
Err(outcome) => return finish(&outcome, &ctx.auto, now),
};
if !applies && !asks {
let outcome = Outcome::Available {
release,
check_only: args.check,
};
return finish(&outcome, &ctx.auto, now);
}
if !applies {
if refuses_cargo(args) {
return stop(&Stop::new(cargo_refusal(CARGO_SUGGESTION), 5));
}
let restarts = matches!(daemon, DaemonState::RunningAndReachable { .. });
if !ask(&release, restarts) {
return finish(&Outcome::Declined { release }, &ctx.auto, now);
}
}
let applied = match daemon {
DaemonState::RunningAndReachable { port, token } => {
via_daemon(release, port, &token, args.force_cargo, Timing::REAL).await
}
DaemonState::NotRunning => in_process(release, ctx, args.force_cargo).await,
DaemonState::RunningButUnauthenticated => Err(Stop::new(unauthenticated(), 6)),
};
match applied {
Ok(outcome) => finish(&outcome, &ctx.auto, chrono::Utc::now()),
Err(s) => stop(&s),
}
}
async fn up_to_date(version: String, identity: &str, port: u16) -> Outcome {
let stale = daemon_version(port)
.await
.filter(|running| !is_version(running, identity));
Outcome::UpToDate { version, stale }
}
fn open(output: &OutputConfig, total: usize) -> Option<ui::RailGuard> {
let guard = ui::install(output, "update", total);
begin();
guard
}
fn begin() {
ui::intro("Updating OpenLatch");
ui::stage_named("Check", "Checking for updates", "Checked for updates");
}
fn plain() -> bool {
ui::caps().is_some_and(|c| c.mode == Mode::Plain)
}
fn margin(text: &str) {
ui::outro(text);
}
fn close_check(check: CheckResult, current: &str, registry: &str) -> Result<Release, Outcome> {
match check {
CheckResult::UpToDate { current } => {
ui::done(&format!("already on the latest ({current})"));
Err(Outcome::UpToDate {
version: current,
stale: None,
})
}
CheckResult::Failed { reason } => {
let unreachable = reason.contains("send:") || reason.starts_with("http client");
ui::failed(if unreachable {
"couldn't reach the update server"
} else {
"the update server had no usable release"
});
ui::detail(&format!("Registry {registry}: {reason}"));
Err(Outcome::CheckFailed {
from: current.to_string(),
host: host(registry),
unreachable,
})
}
CheckResult::Available {
current,
latest,
severity,
size,
..
} => {
let result = format!("{current} → {latest}");
match severity {
Severity::Critical => {
ui::done_toned(&format!("{result} · security fix"), ui::Tone::Fail)
}
Severity::Normal => ui::done(&result),
}
ui::detail(&format!(
"Registry {registry} · latest {latest} · severity {}",
severity.as_str()
));
Ok(Release {
from: current,
to: latest,
severity,
size,
})
}
}
}
fn host(origin: &str) -> String {
let rest = origin.split_once("://").map_or(origin, |(_, rest)| rest);
rest.split('/').next().unwrap_or(rest).to_string()
}
fn question(release: &Release, restarts: bool) -> (String, String, Option<String>) {
let title = format!("Update to {}", release.to);
let question = match release.severity {
Severity::Critical => format!("{} fixes a security issue. Update now?", release.to),
Severity::Normal => format!("Update to {} now?", release.to),
};
let size = release.size.map(|s| format!("About {} MB", mb(s)));
let line = match (size, restarts) {
(Some(size), true) => Some(format!(
"{size} · the background service restarts for a few seconds."
)),
(Some(size), false) => Some(format!("{size}.")),
(None, true) => Some("The background service restarts for a few seconds.".to_string()),
(None, false) => None,
};
(title, question, line)
}
fn ask(release: &Release, restarts: bool) -> bool {
let (title, question, line) = question(release, restarts);
ui::stage_named("Update", &title, &title);
let lines: Vec<&str> = line.as_deref().into_iter().collect();
let select = ui::Select {
title: &title,
question: &question,
lines: &lines,
yes: "Update now",
no: "Not now",
default_yes: true,
timeout: None,
};
let yes = matches!(ui::select(&select), Some(ui::SelectOutcome::Answered(true)));
answered(yes);
yes
}
fn answered(yes: bool) {
ui::done(if yes { "yes" } else { "not now" });
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum Step {
Download,
Verify,
Install,
Restart,
Confirm,
}
impl Step {
fn of(stage: ApplyStage) -> Self {
match stage {
ApplyStage::Check | ApplyStage::Download => Self::Download,
ApplyStage::Extract | ApplyStage::Verify | ApplyStage::Sanity | ApplyStage::Reader => {
Self::Verify
}
ApplyStage::Swap => Self::Install,
ApplyStage::Drain | ApplyStage::Restart => Self::Restart,
ApplyStage::Healthz => Self::Confirm,
}
}
fn next(self) -> Option<Self> {
match self {
Self::Download => Some(Self::Verify),
Self::Verify => Some(Self::Install),
Self::Install => Some(Self::Restart),
Self::Restart => Some(Self::Confirm),
Self::Confirm => None,
}
}
fn name(self) -> &'static str {
match self {
Self::Download => "Download",
Self::Verify => "Verify",
Self::Install => "Install",
Self::Restart => "Restart",
Self::Confirm => "Confirm",
}
}
fn key(self) -> &'static str {
match self {
Self::Download => "download",
Self::Verify => "verify",
Self::Install => "install",
Self::Restart => "restart",
Self::Confirm => "confirm",
}
}
fn done(self) -> &'static str {
match self {
Self::Download => "Downloaded",
Self::Verify => "Verified",
Self::Install => "Installed",
Self::Restart => "Restarted the background service",
Self::Confirm => "Confirmed",
}
}
}
fn apply_stage(wire: &str) -> Option<ApplyStage> {
use ApplyStage::*;
[
Check, Download, Extract, Verify, Sanity, Reader, Swap, Drain, Restart, Healthz,
]
.into_iter()
.find(|s| s.as_str() == wire)
}
#[derive(Debug, Default)]
struct Rate {
last: Option<(Instant, u64)>,
bytes_per_sec: Option<f64>,
}
impl Rate {
fn sample(&mut self, now: Instant, bytes: u64) {
let Some((at, before)) = self.last else {
self.last = Some((now, bytes));
return;
};
let dt = now.saturating_duration_since(at).as_secs_f64();
if dt < 0.2 || bytes < before {
return;
}
let instant = (bytes - before) as f64 / dt;
self.bytes_per_sec = Some(match self.bytes_per_sec {
Some(r) => 0.7 * r + 0.3 * instant,
None => instant,
});
self.last = Some((now, bytes));
}
}
fn download_tail(
done: u64,
total: Option<u64>,
bytes_per_sec: Option<f64>,
) -> (Option<f64>, String) {
let total = total.filter(|t| *t > 0);
let mut parts = vec![match total {
Some(t) => format!("{} / {} MB", mb(done), mb(t)),
None => format!("{} MB", mb(done)),
}];
let rate = bytes_per_sec.filter(|r| *r > 0.0);
if let Some(r) = rate {
parts.push(format!("{:.1} MB/s", r / MIB));
}
if let (Some(t), Some(r)) = (total, rate) {
let left = (t.saturating_sub(done) as f64 / r).ceil() as u64;
parts.push(if left < 60 {
format!("{left}s left")
} else {
format!("{}m left", left.div_ceil(60))
});
}
(total.map(|t| done as f64 / t as f64), parts.join(" · "))
}
const MIB: f64 = 1_048_576.0;
fn mb(bytes: u64) -> String {
format!("{:.1}", bytes as f64 / MIB)
}
fn secs(d: Duration) -> String {
format!("{:.1}s", d.as_secs_f64())
}
struct Apply {
release: Release,
step: Option<Step>,
since: Instant,
total: Option<u64>,
rate: Rate,
drain_at: Option<Instant>,
restart_at: Option<Instant>,
}
impl Apply {
fn new(release: Release, now: Instant) -> Self {
Self {
release,
step: None,
since: now,
total: None,
rate: Rate::default(),
drain_at: None,
restart_at: None,
}
}
fn restarting(&mut self, at: Instant) {
self.restart_at.get_or_insert(at);
}
fn reach(&mut self, step: Step, now: Instant) {
loop {
let next = match self.step {
Some(s) if s >= step => return,
Some(s) => {
self.close(now);
match s.next() {
Some(n) => n,
None => return,
}
}
None => Step::Download,
};
self.open(next, now);
}
}
fn open(&mut self, step: Step, now: Instant) {
let running = match step {
Step::Download => match (plain(), self.release.size) {
(true, Some(size)) => format!("Downloading {} MB", mb(size)),
_ => format!("Downloading {}", self.release.to),
},
Step::Verify => "Verifying the signature".to_string(),
Step::Install => format!("Installing {}", self.release.to),
Step::Restart => "Restarting the background service".to_string(),
Step::Confirm => "Confirming the new version".to_string(),
};
ui::stage_named(step.name(), &running, step.done());
self.step = Some(step);
self.since = match step {
Step::Restart => self.restart_at.map_or(now, |at| at.min(now)),
_ => now,
};
if matches!(step, Step::Download | Step::Restart) {
ui::announce();
}
}
fn close(&mut self, now: Instant) {
let Some(step) = self.step.take() else {
return;
};
let took = now.saturating_duration_since(self.since);
let result = match step {
Step::Download => match self.total.or(self.release.size) {
Some(total) => format!("{} MB in {}", mb(total), secs(took)),
None => format!("in {}", secs(took)),
},
Step::Verify => "signed by OpenLatch".to_string(),
Step::Install => format!("{} kept as a backup", self.release.from),
Step::Restart => format!("back in {}", secs(took)),
Step::Confirm => format!("running {}", self.release.to),
};
ui::done(&result);
}
fn bytes(&mut self, done: Option<u64>, total: Option<u64>, now: Instant) {
let Some(done) = done.filter(|_| self.step == Some(Step::Download)) else {
return;
};
self.total = total.or(self.total);
self.rate.sample(now, done);
let (frac, tail) = download_tail(done, self.total, self.rate.bytes_per_sec);
ui::meter(frac, &tail);
}
fn observe(&mut self, progress: &update::ApplyProgress, now: Instant) {
if let Some(stage) = progress.stage() {
self.reach(Step::of(stage), now);
}
self.bytes(progress.bytes_done(), progress.bytes_total(), now);
}
fn draining(&mut self, now: Instant) {
self.restarting(now);
let at = *self.drain_at.get_or_insert(now);
if self.step == Some(Step::Restart) {
let left = DRAIN.saturating_sub(now.saturating_duration_since(at));
let left = left.as_secs_f64().ceil().max(1.0) as u64;
ui::meter(None, &format!("draining in-flight events · {left}s"));
}
}
fn waiting(&self) {
if self.step == Some(Step::Restart) {
ui::meter(None, &format!("waiting for {} to answer", self.release.to));
}
}
fn finish(&mut self, now: Instant) {
self.close(now);
}
fn fail(&mut self, at: FailedAt, now: Instant) -> Outcome {
let why = at.why(&self.release, self.step);
self.reach(why.step, now);
ui::fail(&OlError::new(why.code, why.line));
self.step = None;
Outcome::Failed {
release: self.release.clone(),
at,
step: why.step,
}
}
fn rolled_back(&mut self) -> Outcome {
let Release { from, to, .. } = &self.release;
ui::fail(&OlError::new(
crate::error::ERR_DAEMON_START_FAILED,
format!("{to} didn't come back"),
));
self.step = None;
ui::stage_named(
"Roll back",
&format!("Rolling back to {from}"),
"Rolled back",
);
ui::done(&format!("{from} restored and running"));
Outcome::RolledBack {
release: self.release.clone(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FailedAt {
Stage(ApplyStage),
NotBack,
Stalled,
}
struct Why {
step: Step,
code: &'static str,
line: String,
headline: String,
what: String,
short: String,
unchanged: bool,
}
impl FailedAt {
fn why(self, release: &Release, running: Option<Step>) -> Why {
use crate::error::{
ERR_DAEMON_START_FAILED as START, ERR_READER_RETIREMENT_BLOCKED as READER,
ERR_UPDATE_DAEMON_UNREACHABLE as GONE, ERR_UPDATE_VERIFY_FAILED as VERIFY,
};
let to = &release.to;
let (step, code, line, headline, what, unchanged) = match self {
Self::Stage(ApplyStage::Check) => (
Step::Download,
VERIFY,
"Couldn't reach the update server".to_string(),
"the download didn't start".to_string(),
"The update server stopped answering.".to_string(),
true,
),
Self::Stage(ApplyStage::Download) => (
Step::Download,
VERIFY,
"The download didn't finish".to_string(),
"the download didn't finish".to_string(),
"The download from the update server didn't finish.".to_string(),
true,
),
Self::Stage(ApplyStage::Extract) => (
Step::Verify,
VERIFY,
"The download is damaged".to_string(),
"the download didn't verify".to_string(),
"The downloaded package is damaged.".to_string(),
true,
),
Self::Stage(ApplyStage::Verify) => (
Step::Verify,
VERIFY,
"The signature doesn't match OpenLatch's release key".to_string(),
"the download didn't verify".to_string(),
"The signature doesn't match OpenLatch's key.".to_string(),
true,
),
Self::Stage(ApplyStage::Sanity) => (
Step::Verify,
VERIFY,
format!("{to} failed its self-test"),
"the download didn't verify".to_string(),
format!("The downloaded {to} failed its self-test."),
true,
),
Self::Stage(ApplyStage::Reader) => (
Step::Verify,
READER,
format!("{to} can't read this machine's policies"),
format!("{to} can't read this machine's policies"),
format!("{to} can't read the policies this machine enforces now."),
true,
),
Self::Stage(ApplyStage::Swap) => (
Step::Install,
VERIFY,
"The installed files couldn't be replaced".to_string(),
format!("{to} couldn't be installed"),
"The installed files couldn't be replaced.".to_string(),
false,
),
Self::Stage(ApplyStage::Drain | ApplyStage::Restart) => (
Step::Restart,
START,
"The background service couldn't restart".to_string(),
"the background service didn't restart".to_string(),
format!("{to} is installed, but the background service couldn't restart into it."),
false,
),
Self::Stage(ApplyStage::Healthz) => (
Step::Confirm,
START,
format!("{to} didn't pass its health check"),
format!("{to} didn't pass its health check"),
format!("{to} started but didn't pass its own health check."),
false,
),
Self::NotBack => (
running.unwrap_or(Step::Restart),
START,
format!("{to} didn't come back"),
format!("{to} didn't come back"),
format!(
"{to} was installed, but the background service didn't answer within 60 s."
),
false,
),
Self::Stalled => (
running.unwrap_or(Step::Download),
GONE,
"The background service stopped reporting progress".to_string(),
"the update stalled".to_string(),
"The background service stopped reporting progress. The update may still finish."
.to_string(),
false,
),
};
let short = match self {
Self::Stage(ApplyStage::Verify) => "The signature doesn't match".to_string(),
_ => what.trim_end_matches('.').to_string(),
};
let step = match self {
Self::NotBack | Self::Stalled => step,
_ => running.map_or(step, |r| r.max(step)),
};
Why {
step,
code,
line,
headline,
what,
short,
unchanged,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Service {
Running,
NotRunning,
Stale(String),
}
#[derive(Debug, Clone)]
enum Outcome {
UpToDate {
version: String,
stale: Option<String>,
},
Available {
release: Release,
check_only: bool,
},
Declined {
release: Release,
},
CheckFailed {
from: String,
host: String,
unreachable: bool,
},
Updated {
release: Release,
service: Service,
},
Failed {
release: Release,
at: FailedAt,
step: Step,
},
RolledBack {
release: Release,
},
}
impl Outcome {
fn kind(&self) -> ui::Update {
match self {
Self::UpToDate { .. } => ui::Update::UpToDate,
Self::Available { .. } => ui::Update::Available,
Self::Declined { .. } => ui::Update::Declined,
Self::CheckFailed { .. } => ui::Update::CheckFailed,
Self::Updated { .. } => ui::Update::Updated,
Self::Failed { .. } => ui::Update::Failed,
Self::RolledBack { .. } => ui::Update::RolledBack,
}
}
fn verdict(&self) -> ui::Verdict {
match self {
Self::Updated { service, .. } if *service != Service::Running => {
ui::Verdict::NeedsAttention
}
Self::UpToDate { stale: Some(_), .. } => ui::Verdict::NeedsAttention,
_ => self.kind().verdict(),
}
}
fn exit_code(&self) -> i32 {
match self.verdict() {
ui::Verdict::Live => 0,
ui::Verdict::NeedsAttention | ui::Verdict::NoAgent => crate::cli::report::EXIT_DEGRADED,
ui::Verdict::NotLive => 1,
}
}
fn outro(&self) -> &'static str {
match self {
Self::UpToDate { .. } | Self::Updated { .. } => "Done.",
Self::Available { .. } | Self::Declined { .. } => "Nothing changed.",
_ => "Stopped.",
}
}
fn card(
&self,
auto: &AutoUpdate,
now: chrono::DateTime<chrono::Utc>,
log: Option<&str>,
) -> (ui::Card, Vec<String>) {
use ui::Row;
let (title, rows, mut warnings) =
match self {
Self::UpToDate { version, stale } => {
let latest =
Row::new("Version", version.as_str()).with_tail("· the latest release");
match stale {
None => (
ui::Update::UpToDate.title(version, version),
vec![latest],
vec![],
),
Some(running) => (
format!("OpenLatch {version} is installed"),
vec![
latest,
Row::new("Service", format!("still running {running}")),
],
vec![stale_warning(running)],
),
}
}
Self::Available {
release,
check_only,
} => (
ui::Update::Available.title(&release.from, &release.to),
vec![
Row::new("Running", release.from.as_str()),
Row::new(
"When ready",
if *check_only {
"Run `openlatch update`."
} else {
"Run `openlatch update --yes`."
},
),
],
vec![],
),
Self::Declined { release } => (
ui::Update::Declined.title(&release.from, &release.to),
vec![
Row::new("Available", release.to.as_str())
.with_tail(severity_tail(release, "")),
Row::new("When ready", "Run `openlatch update` again."),
],
vec![],
),
Self::CheckFailed { from, .. } => (
ui::Update::CheckFailed.title(from, from),
vec![
Row::new("Still running", from.as_str())
.with_tail("· nothing changed, your agents stay protected"),
Row::new("Why", self.check_why()),
Row::new(
"Try",
"Check the connection or proxy, then run `openlatch update`",
),
],
vec![],
),
Self::Updated { release, service } => {
let updated = Row::new("Updated", format!("{} → {}", release.from, release.to))
.with_tail(severity_tail(release, " applied"));
match service {
Service::Running => (
ui::Update::Updated.title(&release.from, &release.to),
vec![
updated,
Row::new("Your agents", "stay protected · nothing to restart"),
],
vec![],
),
Service::NotRunning => (
format!("OpenLatch {} is installed", release.to),
vec![updated],
vec!["The background service isn't running · start it: `openlatch start`"
.to_string()],
),
Service::Stale(running) => (
format!("OpenLatch {} is installed", release.to),
vec![updated, Row::new("Service", format!("still running {running}"))],
vec![stale_warning(running)],
),
}
}
Self::Failed { release, at, step } => {
let why = at.why(release, Some(*step));
let mut rows = Vec::new();
if why.unchanged {
rows.push(
Row::new("Still running", release.from.as_str())
.with_tail("· nothing was changed"),
);
}
rows.push(Row::new("What happened", why.what));
rows.push(Row::more(why.code));
rows.push(Row::new("Try", try_this(why.unchanged)));
(
format!(
"{}: {}",
ui::Update::Failed.title(&release.from, &release.to),
why.headline
),
rows,
vec![],
)
}
Self::RolledBack { release } => (
ui::Update::RolledBack.title(&release.from, &release.to),
vec![
Row::new("Running", release.from.as_str())
.with_tail("· restored automatically, agents protected"),
Row::new(
"What happened",
format!("The {} service stopped right after starting.", release.to),
),
Row::more(crate::error::ERR_DAEMON_START_FAILED),
Row::new("Try", try_this(false)),
],
vec![],
),
};
let mut footer = Vec::new();
match auto.line(now) {
(true, line) => warnings.push(line),
(false, line) => footer.push(line),
}
if let Some(log) = log {
footer.push(format!("Log {log}"));
}
let card = ui::Card {
verdict: self.verdict(),
title,
rows,
sections: vec![],
footer,
};
(card, warnings)
}
fn check_why(&self) -> String {
match self {
Self::CheckFailed {
host,
unreachable: true,
..
} => format!("{host} didn't answer"),
Self::CheckFailed { host, .. } => format!("{host} didn't return a usable release"),
_ => String::new(),
}
}
fn plain_note(&self) -> Option<String> {
match self {
Self::Available {
check_only: true, ..
} => Some("To install it, run: openlatch update --yes".to_string()),
Self::Available { .. } => Some(
"Not applied: no terminal to ask. To update unattended, run: openlatch update --yes"
.to_string(),
),
_ => None,
}
}
fn sentence(&self) -> Option<String> {
match self {
Self::CheckFailed { from, .. } => Some(format!(
"Still running {from} · nothing changed. {}.",
self.check_why()
)),
Self::Updated {
release,
service: Service::NotRunning,
} => Some(format!(
"Installed {} · the background service isn't running. Start it: openlatch start",
release.to
)),
Self::Updated {
release: Release { to: version, .. },
service: Service::Stale(running),
}
| Self::UpToDate {
version,
stale: Some(running),
} => Some(format!(
"Installed {version} · the background service is still on {running}. Run: openlatch restart"
)),
Self::Failed { release, at, step } => {
let why = at.why(release, Some(*step));
Some(if why.unchanged {
format!(
"Still running {} · nothing was changed. {} ({}).",
release.from, why.short, why.code
)
} else {
format!("{} ({}).", why.short, why.code)
})
}
Self::RolledBack { release } => Some(format!(
"Running {} again · {} stopped right after starting and was rolled back ({}).",
release.from,
release.to,
crate::error::ERR_DAEMON_START_FAILED
)),
_ => None,
}
}
fn fields(&self, log: Option<&str>) -> Vec<(&'static str, String)> {
let mut f: Vec<(&'static str, String)> = match self {
Self::UpToDate { version, stale } => {
let mut f = vec![("version", version.clone())];
if stale.is_some() {
f.push(("hint", "openlatch restart".to_string()));
}
f
}
Self::Available { release, .. } => vec![
("from", release.from.clone()),
("to", release.to.clone()),
("severity", release.severity.as_str().to_string()),
("hint", "openlatch update --yes".to_string()),
],
Self::Declined { release } => {
vec![("from", release.from.clone()), ("to", release.to.clone())]
}
Self::CheckFailed {
from, unreachable, ..
} => vec![
("version", from.clone()),
(
"reason",
if *unreachable {
"registry-unreachable"
} else {
"registry-error"
}
.to_string(),
),
],
Self::Updated { release, service } => {
let mut f = vec![("from", release.from.clone()), ("to", release.to.clone())];
match service {
Service::Running => {}
Service::NotRunning => f.push(("hint", "openlatch start".to_string())),
Service::Stale(_) => f.push(("hint", "openlatch restart".to_string())),
}
f
}
Self::Failed { release, at, step } => {
let why = at.why(release, Some(*step));
let mut f = vec![
("step", why.step.key().to_string()),
("code", why.code.to_string()),
];
if why.unchanged {
f.push(("version", release.from.clone()));
}
f
}
Self::RolledBack { release } => vec![
("from", release.to.clone()),
("to", release.from.clone()),
("code", crate::error::ERR_DAEMON_START_FAILED.to_string()),
],
};
if let Some(log) = log {
f.push(("log", log.to_string()));
}
f
}
}
fn stale_warning(running: &str) -> String {
format!("The background service is still on {running} · restart it: `openlatch restart`")
}
fn severity_tail(release: &Release, suffix: &str) -> String {
match release.severity {
Severity::Critical => format!("· security fix{suffix}"),
Severity::Normal => String::new(),
}
}
fn try_this(unchanged: bool) -> &'static str {
if unchanged {
"Run `openlatch update` again. If it fails twice, share the log below with IT or OpenLatch."
} else {
"Run `openlatch doctor`, then share the log below."
}
}
#[derive(Debug, Clone, Default)]
struct AutoUpdate {
on: bool,
by_env: bool,
last_check: Option<String>,
}
impl AutoUpdate {
fn line(&self, now: chrono::DateTime<chrono::Utc>) -> (bool, String) {
if !self.on {
let remedy = if self.by_env {
"`OPENLATCH_AUTO_UPDATE=true`"
} else {
"`[update] auto_update = true`"
};
return (true, format!("Automatic updates: off · turn on: {remedy}"));
}
let checked = self
.last_check
.as_deref()
.and_then(|t| chrono::DateTime::parse_from_rfc3339(t).ok())
.map(|t| ago(now.signed_duration_since(t)));
match checked {
Some(ago) => (false, format!("Automatic updates: on · last checked {ago}")),
None => (false, "Automatic updates: on".to_string()),
}
}
fn plain_line(&self, now: chrono::DateTime<chrono::Utc>) -> String {
let (warn, line) = self.line(now);
let line = ui::style(&line, false);
if warn {
format!("! {line}")
} else {
line
}
}
}
fn ago(elapsed: chrono::Duration) -> String {
let s = elapsed.num_seconds().max(0);
match s {
0..60 => "just now".to_string(),
60..3600 => format!("{}m ago", s / 60),
3600..86_400 => format!("{}h ago", s / 3600),
_ => format!("{}d ago", s / 86_400),
}
}
fn finish(outcome: &Outcome, auto: &AutoUpdate, now: chrono::DateTime<chrono::Utc>) -> i32 {
let exit = outcome.exit_code();
let log = ui::log_path()
.map(|p| p.display().to_string())
.filter(|_| outcome.verdict() == ui::Verdict::NotLive);
if plain() {
if let Some(note) = outcome.plain_note() {
margin(&format!(" {note}"));
}
if let Some(sentence) = outcome.sentence() {
margin(&sentence);
}
margin(&auto.plain_line(now));
} else {
ui::outro(outcome.outro());
let (card, warnings) = outcome.card(auto, now, log.as_deref());
ui::card_with_warnings(&card, &warnings);
}
let fields = outcome.fields(log.as_deref());
let fields: Vec<(&str, &str)> = fields.iter().map(|(k, v)| (*k, v.as_str())).collect();
ui::result_token(outcome.kind().token(), exit, &fields);
exit
}
struct Stop {
error: OlError,
exit: i32,
}
impl Stop {
fn new(error: OlError, exit: i32) -> Self {
Self { error, exit }
}
}
fn stop(s: &Stop) -> i32 {
ui::fail(&s.error);
ui::outro("Stopped.");
ui::result_token("failed", s.exit, &[("code", s.error.code)]);
s.exit
}
#[derive(Debug, Clone, Copy)]
struct Timing {
poll: Duration,
request: Duration,
back_within: Duration,
stall: Duration,
misses: u32,
}
impl Timing {
const REAL: Self = Self {
poll: POLL,
request: Duration::from_secs(2),
back_within: BACK_WITHIN,
stall: STALL,
misses: 8,
};
}
enum Follow {
Done(Outcome),
Restarting {
dropped: bool,
},
}
enum Back {
Done(Outcome),
Resume,
}
#[derive(Default)]
struct Watch {
seen: (Option<String>, Option<u64>),
changed: Option<Instant>,
swapped: bool,
misses: u32,
}
async fn via_daemon(
release: Release,
port: u16,
token: &str,
force_cargo: bool,
timing: Timing,
) -> Result<Outcome, Stop> {
let client = local_client(Duration::from_secs(15)).map_err(|e| {
Stop::new(
OlError::new(
crate::error::ERR_UPDATE_DAEMON_UNREACHABLE,
format!("failed to build HTTP client: {e}"),
),
1,
)
})?;
let resp = client
.post(format!("http://127.0.0.1:{port}/admin/update"))
.bearer_auth(token)
.json(&json!({"force_cargo_install": force_cargo}))
.send()
.await
.map_err(|e| Stop::new(rpc_failed(&e), 6))?;
let status = resp.status();
let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
if status != reqwest::StatusCode::ACCEPTED {
return match rejection(status, &body, &release.from) {
Rejection::UpToDate(version) => Ok(Outcome::UpToDate {
version,
stale: None,
}),
Rejection::Refused(error, exit) => Err(Stop::new(error, exit)),
};
}
let said = |key: &str, or: &str| {
body.get(key)
.and_then(|v| v.as_str())
.unwrap_or(or)
.to_string()
};
let release = Release {
from: said("from", &release.from),
to: said("to", &release.to),
..release
};
let mut apply = Apply::new(release, Instant::now());
apply.reach(Step::Download, Instant::now());
let mut watch = Watch::default();
loop {
let dropped = match follow(&client, port, token, &mut apply, &mut watch, timing).await {
Follow::Done(outcome) => return Ok(outcome),
Follow::Restarting { dropped } => dropped,
};
match come_back(&client, port, token, &mut apply, &watch, dropped, timing).await {
Back::Done(outcome) => return Ok(outcome),
Back::Resume => {
watch.misses = 0;
if !watch.swapped {
apply.restart_at = None;
}
}
}
}
}
async fn follow(
client: &reqwest::Client,
port: u16,
token: &str,
apply: &mut Apply,
watch: &mut Watch,
timing: Timing,
) -> Follow {
loop {
let now = Instant::now();
let changed = *watch.changed.get_or_insert(now);
if now.saturating_duration_since(changed) > timing.stall {
return Follow::Done(apply.fail(FailedAt::Stalled, now));
}
let Some(poll) = poll_status(client, port, token, timing.request).await else {
if health_version(client, port, timing.request).await.is_some() {
watch.misses = 0;
} else {
watch.misses += 1;
if watch.swapped || watch.misses >= timing.misses {
apply.restarting(Instant::now());
return Follow::Restarting { dropped: true };
}
}
tokio::time::sleep(timing.poll).await;
continue;
};
watch.misses = 0;
let now = Instant::now();
let key = (poll.stage.clone(), poll.bytes_done);
if key != watch.seen {
watch.seen = key;
watch.changed = Some(now);
}
let stage = poll.stage.as_deref().and_then(apply_stage);
if matches!(
stage,
Some(ApplyStage::Swap | ApplyStage::Drain | ApplyStage::Restart)
) {
apply.restarting(now);
}
if let Some(stage) = stage {
apply.reach(Step::of(stage), now);
}
watch.swapped |= matches!(
stage,
Some(ApplyStage::Swap | ApplyStage::Drain | ApplyStage::Restart | ApplyStage::Healthz)
);
apply.bytes(poll.bytes_done, poll.bytes_total, now);
if matches!(stage, Some(ApplyStage::Drain | ApplyStage::Restart)) {
apply.draining(now);
}
match poll.status.as_deref() {
Some("completed") => {
if !watch.swapped {
let serving = health_version(client, port, timing.request).await;
if serving.is_some_and(|v| serves_from(&v, &apply.release, IDENTITY)) {
ui::done("already on the latest");
return Follow::Done(Outcome::UpToDate {
version: apply.release.from.clone(),
stale: None,
});
}
}
return Follow::Restarting { dropped: false };
}
Some("idle") => {
apply.restarting(now);
return Follow::Restarting { dropped: true };
}
Some("failed") => {
if let Some(error) = &poll.error {
ui::detail(error);
}
let at = FailedAt::Stage(stage.unwrap_or(ApplyStage::Download));
return Follow::Done(apply.fail(at, now));
}
_ => {}
}
tokio::time::sleep(timing.poll).await;
}
}
async fn come_back(
client: &reqwest::Client,
port: u16,
token: &str,
apply: &mut Apply,
watch: &Watch,
mut dropped: bool,
timing: Timing,
) -> Back {
if watch.swapped {
apply.reach(Step::Restart, Instant::now());
}
let deadline = Instant::now() + timing.back_within;
let mut answered = false;
while Instant::now() < deadline {
let now = Instant::now();
match health_version(client, port, timing.request).await {
None => {
dropped = true;
apply.restarting(now);
apply.waiting();
}
Some(v) if is_version(&v, &apply.release.to) => {
apply.reach(Step::Confirm, now);
answered = true;
let settled = poll_status(client, port, token, timing.request).await;
if !matches!(
settled.and_then(|p| p.status).as_deref(),
Some("in_progress" | "failed")
) {
apply.finish(Instant::now());
return Back::Done(Outcome::Updated {
release: apply.release.clone(),
service: Service::Running,
});
}
}
Some(v)
if dropped
&& serves_from(&v, &apply.release, IDENTITY)
&& !is_version(&v, &apply.release.to) =>
{
let poll = poll_status(client, port, token, timing.request).await;
let applying = poll.as_ref().is_some_and(|p| {
p.status.as_deref() == Some("in_progress")
&& !matches!(
p.stage.as_deref().and_then(apply_stage),
Some(ApplyStage::Drain | ApplyStage::Restart | ApplyStage::Healthz)
)
});
if applying {
return Back::Resume;
}
return Back::Done(apply.rolled_back());
}
Some(_) => apply.draining(now),
}
tokio::time::sleep(timing.poll).await;
}
let at = if answered {
FailedAt::Stage(ApplyStage::Healthz)
} else {
FailedAt::NotBack
};
Back::Done(apply.fail(at, Instant::now()))
}
async fn in_process(release: Release, ctx: &Ctx, force_cargo: bool) -> Result<Outcome, Stop> {
let progress = Arc::new(update::ApplyProgress::default());
let opts = update::ApplyOptions {
current_version: ctx.current.clone(),
registry_origin: ctx.registry.clone(),
download_timeout: Duration::from_secs(60),
force_cargo_install: force_cargo,
mode: update::ApplyMode::InProcess,
egress: ctx.egress.clone(),
progress: progress.clone(),
};
let mut apply = Apply::new(release, Instant::now());
apply.reach(Step::Download, Instant::now());
let run = update::apply_local(opts);
tokio::pin!(run);
let mut tick = tokio::time::interval(TICK);
let result = loop {
tokio::select! {
result = &mut run => break result,
_ = tick.tick() => apply.observe(&progress, Instant::now()),
}
};
let now = Instant::now();
match result {
update::ApplyResult::Applied { from, to, .. } => {
apply.reach(Step::Install, now);
apply.finish(now);
let service = match daemon_version(ctx.port).await {
None => Service::NotRunning,
Some(v) if is_version(&v, &to) => Service::Running,
Some(v) => Service::Stale(v),
};
let release = Release {
from,
to,
..apply.release
};
Ok(Outcome::Updated { release, service })
}
update::ApplyResult::UpToDate { current } => {
ui::done("already on the latest");
Ok(Outcome::UpToDate {
version: current,
stale: None,
})
}
update::ApplyResult::RefusedCargoInstall { suggestion } => {
Err(Stop::new(cargo_refusal(&suggestion), 5))
}
update::ApplyResult::Failed { stage, reason } => {
ui::detail(&format!("{}: {reason}", stage.as_str()));
Ok(apply.fail(FailedAt::Stage(stage), now))
}
}
}
fn print_json(output: &OutputConfig, value: &serde_json::Value) {
if output.format == crate::cli::output::OutputFormat::Json {
output.print_json(value);
}
}
async fn run_json(args: &UpdateArgs, output: &OutputConfig, ctx: &Ctx) -> i32 {
if !args.applies() {
return run_check(output, &ctx.current, &ctx.registry, &ctx.egress).await;
}
if refuses_cargo(args) {
output.print_error(&cargo_refusal(CARGO_SUGGESTION));
return 5;
}
match probe_daemon(ctx.port).await {
DaemonState::RunningAndReachable { port, token } => {
apply_via_daemon_rpc(output, &ctx.current, port, &token, args.force_cargo).await
}
DaemonState::NotRunning => apply_in_process(output, ctx, args).await,
DaemonState::RunningButUnauthenticated => {
output.print_error(&unauthenticated());
6
}
}
}
async fn run_check(
output: &OutputConfig,
current: &str,
registry: &str,
egress: &crate::egress::EgressConfig,
) -> i32 {
let result = update::check(current, registry, egress).await;
match result {
update::CheckResult::UpToDate { current } => {
print_json(output, &json!({"current": current, "latest": null}));
0
}
update::CheckResult::Available {
current,
latest,
severity,
..
} => {
print_json(
output,
&json!({
"current": current,
"latest": latest,
"severity": severity.as_str(),
}),
);
0
}
update::CheckResult::Failed { reason } => {
output.print_info(&format!("update check failed: {reason}"));
0
}
}
}
async fn apply_in_process(output: &OutputConfig, ctx: &Ctx, args: &UpdateArgs) -> i32 {
let opts = update::ApplyOptions {
current_version: ctx.current.clone(),
registry_origin: ctx.registry.clone(),
download_timeout: Duration::from_secs(60),
force_cargo_install: args.force_cargo,
mode: update::ApplyMode::InProcess,
egress: ctx.egress.clone(),
progress: Default::default(),
};
match update::apply_local(opts).await {
update::ApplyResult::Applied { from, to, .. } => {
output.print_info(&format!("Updated {from} → {to} (no daemon was running)"));
let stale = warn_if_daemon_still_stale(output, ctx.port, &to).await;
print_json(
output,
&json!({
"from": from, "to": to, "applied": true, "restart_required": stale,
}),
);
0
}
update::ApplyResult::UpToDate { current } => {
output.print_info(&format!("Already on the latest version ({current})"));
print_json(output, &json!({"current": current, "idempotent": true}));
0
}
update::ApplyResult::RefusedCargoInstall { suggestion } => {
output.print_error(&cargo_refusal(&suggestion));
5
}
update::ApplyResult::Failed { stage, reason } => {
let err = OlError::new(
crate::error::ERR_UPDATE_VERIFY_FAILED,
format!("auto-update failed at stage `{}`: {reason}", stage.as_str()),
);
output.print_error(&err);
1
}
}
}
async fn apply_via_daemon_rpc(
output: &OutputConfig,
current: &str,
port: u16,
token: &str,
force_cargo: bool,
) -> i32 {
let client = match local_client(Duration::from_secs(15)) {
Ok(c) => c,
Err(e) => {
output.print_info(&format!("failed to build HTTP client: {e}"));
return 1;
}
};
let admin_url = format!("http://127.0.0.1:{port}/admin/update");
let status_url = format!("{admin_url}/status");
let post_resp = match client
.post(&admin_url)
.bearer_auth(token)
.json(&json!({"force_cargo_install": force_cargo}))
.send()
.await
{
Ok(r) => r,
Err(e) => {
output.print_error(&rpc_failed(&e));
return 6;
}
};
let status = post_resp.status();
if status != reqwest::StatusCode::ACCEPTED {
let body: serde_json::Value = post_resp.json().await.unwrap_or(serde_json::Value::Null);
return match rejection(status, &body, current) {
Rejection::UpToDate(cur) => {
output.print_info(&format!("Already on the latest version ({cur})"));
print_json(output, &json!({"current": cur, "idempotent": true}));
0
}
Rejection::Refused(err, exit) => {
output.print_error(&err);
exit
}
};
}
let post_body: serde_json::Value = post_resp.json().await.unwrap_or(serde_json::Value::Null);
let from = post_body
.get("from")
.and_then(|v| v.as_str())
.unwrap_or(current)
.to_string();
let to = post_body
.get("to")
.and_then(|v| v.as_str())
.unwrap_or("?")
.to_string();
output.print_info(&format!("Updating {from} → {to}…"));
let started = std::time::Instant::now();
let max_wait = Duration::from_secs(120);
let poll_interval = Duration::from_secs(1);
loop {
if started.elapsed() > max_wait {
output.print_info("daemon long-poll timed out — the update may still be in progress");
return 1;
}
let resp = client.get(&status_url).bearer_auth(token).send().await;
match resp {
Ok(r) if r.status().is_success() => {
let body: StatusPoll = r.json().await.unwrap_or_default();
match body.status.as_deref().unwrap_or("in_progress") {
"completed" => {
output.print_info(&format!("Updated {from} → {to}"));
let stale = warn_if_daemon_still_stale(output, port, &to).await;
print_json(
output,
&json!({
"from": from, "to": to, "applied": true, "restart_required": stale,
}),
);
return 0;
}
"failed" => {
let stage = body.stage.as_deref().unwrap_or("unknown");
let reason = body.error.as_deref().unwrap_or("");
let code = if matches!(stage, "drain" | "restart" | "healthz") {
crate::error::ERR_DAEMON_START_FAILED
} else {
crate::error::ERR_UPDATE_VERIFY_FAILED
};
let err = OlError::new(
code,
format!("update failed at stage `{stage}`: {reason}"),
);
output.print_error(&err);
return 1;
}
_ => {
tokio::time::sleep(poll_interval).await;
}
}
}
Ok(_) | Err(_) => {
if let Some(new_version) = wait_for_daemon_version(port, &to).await {
output.print_info(&format!("Updated {from} → {new_version}"));
print_json(
output,
&json!({
"from": from,
"to": new_version,
"applied": true,
}),
);
return 0;
}
let err = OlError::new(
crate::error::ERR_DAEMON_START_FAILED,
"daemon did not return after the update — check `openlatch status`",
);
output.print_error(&err);
return 1;
}
}
}
}
async fn warn_if_daemon_still_stale(output: &OutputConfig, port: u16, expected: &str) -> bool {
match daemon_version(port).await {
Some(running) if running != expected => {
output.print_info(&format!(
"The running daemon is still serving {running} — run `openlatch restart` to \
put {expected} in effect."
));
crate::cli::report::record_exit_code(crate::cli::report::EXIT_DEGRADED);
true
}
_ => false,
}
}
async fn wait_for_daemon_version(port: u16, expected: &str) -> Option<String> {
let client = local_client(Duration::from_secs(2)).ok()?;
let deadline = std::time::Instant::now() + Duration::from_secs(60);
while std::time::Instant::now() < deadline {
if let Some(v) = health_version(&client, port, Duration::from_secs(2)).await {
if is_version(&v, expected) {
return Some(v);
}
}
tokio::time::sleep(Duration::from_secs(2)).await;
}
None
}
#[derive(Debug, Default, serde::Deserialize)]
#[serde(default)]
struct StatusPoll {
status: Option<String>,
stage: Option<String>,
error: Option<String>,
bytes_done: Option<u64>,
bytes_total: Option<u64>,
}
const CONNECT: Duration = Duration::from_millis(300);
fn local_client(timeout: Duration) -> reqwest::Result<reqwest::Client> {
crate::egress::client_builder()
.connect_timeout(CONNECT)
.timeout(timeout)
.use_rustls_tls()
.build()
}
async fn poll_status(
client: &reqwest::Client,
port: u16,
token: &str,
timeout: Duration,
) -> Option<StatusPoll> {
let resp = client
.get(format!("http://127.0.0.1:{port}/admin/update/status"))
.bearer_auth(token)
.timeout(timeout)
.send()
.await
.ok()
.filter(|r| r.status().is_success())?;
Some(resp.json().await.unwrap_or_default())
}
async fn health_version(client: &reqwest::Client, port: u16, timeout: Duration) -> Option<String> {
let body: serde_json::Value = client
.get(format!("http://127.0.0.1:{port}/health"))
.timeout(timeout)
.send()
.await
.ok()
.filter(|r| r.status().is_success())?
.json()
.await
.ok()?;
body.get("version")?.as_str().map(str::to_string)
}
async fn daemon_version(port: u16) -> Option<String> {
let timeout = Duration::from_secs(2);
health_version(&local_client(timeout).ok()?, port, timeout).await
}
const IDENTITY: &str = env!("OPENLATCH_VERSION");
fn serves_from(reported: &str, release: &Release, identity: &str) -> bool {
is_version(reported, &release.from) || is_version(reported, identity)
}
fn is_version(reported: &str, expected: &str) -> bool {
reported
.strip_prefix(expected)
.is_some_and(|rest| rest.is_empty() || rest.starts_with('+'))
}
fn read_token() -> Option<String> {
let token = std::fs::read_to_string(config::openlatch_dir().join("daemon.token")).ok()?;
Some(token.trim().to_string()).filter(|t| !t.is_empty())
}
const CARGO_SUGGESTION: &str = "Run: cargo install --force --locked openlatch-client";
fn refuses_cargo(args: &UpdateArgs) -> bool {
!args.force_cargo
&& matches!(
install_state::detect_install_method(),
install_state::InstallMethod::CargoInstall
)
}
fn cargo_refusal(suggestion: &str) -> OlError {
OlError::new(
crate::error::ERR_UPDATE_REFUSED_CARGO_INSTALL,
"this binary was installed via `cargo install` — auto-update would not take effect",
)
.with_suggestion(suggestion)
}
fn unauthenticated() -> OlError {
OlError::new(
crate::error::ERR_UPDATE_DAEMON_UNREACHABLE,
"daemon is running but the bearer token in ~/.openlatch/daemon.token is missing or wrong",
)
.with_suggestion(
"Run `openlatch init --reconfig` to regenerate, or stop the daemon and re-run.",
)
}
fn rpc_failed(e: &reqwest::Error) -> OlError {
OlError::new(
crate::error::ERR_UPDATE_DAEMON_UNREACHABLE,
format!("daemon RPC failed: {e}"),
)
.with_suggestion("Is the daemon still running? Try `openlatch status`.")
}
enum Rejection {
UpToDate(String),
Refused(OlError, i32),
}
fn rejection(status: reqwest::StatusCode, body: &serde_json::Value, current: &str) -> Rejection {
let message = body
.pointer("/error/message")
.and_then(|v| v.as_str())
.unwrap_or("daemon rejected update");
let suggestion = body
.pointer("/error/suggestion")
.and_then(|v| v.as_str())
.map(String::from);
match status.as_u16() {
409 => {
if body.get("idempotent").and_then(|v| v.as_bool()) == Some(true) {
let cur = body
.get("current")
.and_then(|v| v.as_str())
.unwrap_or(current);
return Rejection::UpToDate(cur.to_string());
}
let mut err = OlError::new(
crate::error::ERR_UPDATE_REFUSED_CARGO_INSTALL,
message.to_string(),
);
if let Some(s) = suggestion {
err = err.with_suggestion(s);
}
Rejection::Refused(err, 5)
}
503 => Rejection::Refused(
OlError::new(crate::error::ERR_DAEMON_START_FAILED, message.to_string()),
1,
),
_ => Rejection::Refused(
OlError::new(
crate::error::ERR_UPDATE_VERIFY_FAILED,
format!("daemon rejected update (HTTP {status}): {message}"),
),
1,
),
}
}
enum DaemonState {
RunningAndReachable { port: u16, token: String },
NotRunning,
RunningButUnauthenticated,
}
async fn probe_daemon(port: u16) -> DaemonState {
let Ok(client) = local_client(Duration::from_secs(2)) else {
return DaemonState::NotRunning;
};
let health_url = format!("http://127.0.0.1:{port}/health");
if client.get(&health_url).send().await.is_err() {
return DaemonState::NotRunning;
}
match read_token() {
Some(token) => DaemonState::RunningAndReachable { port, token },
None => DaemonState::RunningButUnauthenticated,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn status_poll_tolerates_a_daemon_without_byte_counts() {
let older: StatusPoll = serde_json::from_value(json!({
"status": "in_progress", "stage": "download", "from": "0.1.0", "to": "0.2.0",
"started_at": null, "ended_at": null, "error": null,
}))
.expect("an older daemon's body parses");
assert_eq!(older.status.as_deref(), Some("in_progress"));
assert_eq!(older.bytes_done, None);
assert_eq!(older.bytes_total, None);
let newer: StatusPoll = serde_json::from_value(json!({
"status": "in_progress", "stage": "some_future_stage",
"bytes_done": 10, "bytes_total": 20,
}))
.expect("a newer daemon's body parses");
assert_eq!(newer.stage.as_deref(), Some("some_future_stage"));
assert_eq!((newer.bytes_done, newer.bytes_total), (Some(10), Some(20)));
}
}
#[cfg(test)]
mod rail_tests {
use super::*;
use crate::cli::ui::{Rail, INSTALL_LOCK};
const SIZE: u64 = 8_806_432;
fn release(severity: Severity) -> Release {
Release {
from: "0.5.7".into(),
to: "0.5.8".into(),
severity,
size: Some(SIZE),
}
}
fn available(severity: Severity) -> CheckResult {
CheckResult::Available {
current: "0.5.7".into(),
latest: "0.5.8".into(),
severity,
tarball_url: String::new(),
tarball_integrity: String::new(),
size: Some(SIZE),
}
}
fn now() -> chrono::DateTime<chrono::Utc> {
"2026-09-23T12:00:00Z".parse().expect("a timestamp")
}
fn auto_on() -> AutoUpdate {
AutoUpdate {
on: true,
by_env: false,
last_check: Some("2026-09-23T08:00:00Z".into()),
}
}
fn drive(mode: Mode, total: usize, script: impl FnOnce(Instant) -> i32) -> (Vec<String>, i32) {
let _l = INSTALL_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (rail, screen) = Rail::buffer(Caps::new(mode, false, false, 80), total);
let exit = {
let _guard = ui::install_rail(rail).expect("installs");
begin();
script(Instant::now())
};
let lines = screen.lock().unwrap().clone();
if mode != Mode::Live {
assert!(lines.iter().all(|l| l.is_ascii()), "{mode:?}: {lines:#?}");
}
(lines, exit)
}
fn ms(t0: Instant, ms: u64) -> Instant {
t0 + Duration::from_millis(ms)
}
fn check_and_download(t0: Instant, severity: Severity, ask: bool) -> Apply {
let release = close_check(available(severity), "0.5.7", "https://registry.npmjs.org")
.expect("a release");
if ask {
let (title, _, _) = question(&release, true);
ui::stage_named("Update", &title, &title);
answered(true);
}
let mut apply = Apply::new(release, t0);
apply.reach(Step::Download, t0);
apply.bytes(Some(0), Some(SIZE), t0);
apply.bytes(Some(SIZE / 2), Some(SIZE), ms(t0, 1000));
apply.bytes(Some(SIZE), Some(SIZE), ms(t0, 2900));
apply.reach(Step::of(ApplyStage::Extract), ms(t0, 2900));
apply
}
fn install_and_restart(apply: &mut Apply, t0: Instant) {
apply.reach(Step::of(ApplyStage::Swap), ms(t0, 4000));
apply.reach(Step::of(ApplyStage::Drain), ms(t0, 5100));
apply.draining(ms(t0, 5100));
apply.waiting();
}
fn updated(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 6, |t0| {
let mut apply = check_and_download(t0, Severity::Normal, mode != Mode::Plain);
install_and_restart(&mut apply, t0);
apply.reach(Step::of(ApplyStage::Healthz), ms(t0, 10_500));
apply.finish(ms(t0, 11_300));
let outcome = Outcome::Updated {
release: apply.release.clone(),
service: Service::Running,
};
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn updated_in_every_mode() {
let (live, exit) = updated(Mode::Live);
assert_eq!(exit, 0);
assert_eq!(
live[..18],
[
"┌ Updating OpenLatch",
"│",
"◇ Checked for updates · 0.5.7 → 0.5.8",
"│",
"◇ Update to 0.5.8 · yes",
"│",
"◇ Downloaded · 8.4 MB in 2.9s",
"│",
"◇ Verified · signed by OpenLatch",
"│",
"◇ Installed · 0.5.7 kept as a backup",
"│",
"◇ Restarted the background service · back in 5.4s",
"│",
"◇ Confirmed · running 0.5.8",
"│",
"└ Done.",
"",
]
);
assert!(live[18].starts_with(" ╭─"), "{live:#?}");
assert!(live[20].starts_with(" │ ✓ OpenLatch 0.5.8 is running"));
let card = live.join("\n");
assert!(card.contains(" Updated 0.5.7 → 0.5.8"));
assert!(card.contains(" Your agents stay protected · nothing to restart"));
assert_eq!(
live.last().map(String::as_str),
Some(" Automatic updates: on · last checked 4h ago")
);
let (append, _) = updated(Mode::Append);
assert_eq!(
append[..9],
[
"+ Updating OpenLatch",
"|",
"> Checking for updates...",
"+ Checked for updates - 0.5.7 -> 0.5.8",
"|",
"> Update to 0.5.8...",
"+ Update to 0.5.8 - yes",
"|",
"> Downloading 0.5.8...",
]
);
assert_eq!(append[9], "| 25%");
assert_eq!(append[10], "| 50% 4.2 / 8.4 MB - 4.2 MB/s - 1s left");
assert_eq!(append[11], "| 75%");
assert!(
append[12].starts_with("| 100% 8.4 / 8.4 MB - "),
"{append:#?}"
);
assert_eq!(append[13], "+ Downloaded - 8.4 MB in 2.9s");
assert!(append.contains(&"> Restarting the background service...".to_string()));
assert!(append.contains(&"+ Confirmed - running 0.5.8".to_string()));
assert!(append.contains(&"+ Done.".to_string()));
assert!(append
.iter()
.any(|l| l.starts_with(" | + OpenLatch 0.5.8 is running")));
let (plain, exit) = updated(Mode::Plain);
assert_eq!(exit, 0);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/6] Checked for updates - 0.5.7 -> 0.5.8",
"[2/6] Downloading 8.4 MB...",
"[2/6] Downloaded - 8.4 MB in 2.9s",
"[3/6] Verified - signed by OpenLatch",
"[4/6] Installed - 0.5.7 kept as a backup",
"[5/6] Restarting the background service...",
"[5/6] Restarted the background service - back in 5.4s",
"[6/6] Confirmed - running 0.5.8",
"Automatic updates: on - last checked 4h ago",
"result=updated from=0.5.7 to=0.5.8 exit=0",
]
);
}
fn up_to_date(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 1, |_| {
let check = CheckResult::UpToDate {
current: "0.5.8".into(),
};
let outcome = close_check(check, "0.5.8", "https://registry.npmjs.org")
.expect_err("nothing to install");
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn up_to_date_in_every_mode() {
let (live, exit) = up_to_date(Mode::Live);
assert_eq!(exit, 0);
assert_eq!(
live[..6],
[
"┌ Updating OpenLatch",
"│",
"◇ Checked for updates · already on the latest (0.5.8)",
"│",
"└ Done.",
"",
]
);
assert!(live[8].starts_with(" │ ✓ OpenLatch 0.5.8 is up to date"));
assert!(live
.iter()
.any(|l| l.contains(" Version 0.5.8 · the latest release")));
let (append, _) = up_to_date(Mode::Append);
assert!(append.contains(&"+ Checked for updates - already on the latest (0.5.8)".into()));
assert!(append
.iter()
.any(|l| l.starts_with(" | + OpenLatch 0.5.8 is up to date")));
let (plain, _) = up_to_date(Mode::Plain);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/1] Checked for updates - already on the latest (0.5.8)",
"Automatic updates: on - last checked 4h ago",
"result=up_to_date version=0.5.8 exit=0",
]
);
}
fn available_not_applied(mode: Mode, check_only: bool) -> (Vec<String>, i32) {
drive(mode, 1, |_| {
let check = available(Severity::Normal);
let release =
close_check(check, "0.5.7", "https://registry.npmjs.org").expect("a release");
let outcome = Outcome::Available {
release,
check_only,
};
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn available_without_a_terminal_in_every_mode() {
let (live, exit) = available_not_applied(Mode::Live, false);
assert_eq!(exit, 7);
assert!(live.contains(&"└ Nothing changed.".to_string()));
assert!(live
.iter()
.any(|l| l.starts_with(" │ ! 0.5.8 is available")));
assert!(live
.iter()
.any(|l| l.contains("When ready Run openlatch update --yes.")));
let (append, _) = available_not_applied(Mode::Append, false);
assert!(append
.iter()
.any(|l| l.starts_with(" | ! 0.5.8 is available")));
let (plain, exit) = available_not_applied(Mode::Plain, false);
assert_eq!(exit, 7);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/1] Checked for updates - 0.5.7 -> 0.5.8",
" Not applied: no terminal to ask. To update unattended, run: openlatch update --yes",
"Automatic updates: on - last checked 4h ago",
r#"result=available from=0.5.7 to=0.5.8 severity=normal hint="openlatch update --yes" exit=7"#,
]
);
let (plain, _) = available_not_applied(Mode::Plain, true);
assert_eq!(plain[2], " To install it, run: openlatch update --yes");
let (live, _) = available_not_applied(Mode::Live, true);
assert!(live
.iter()
.any(|l| l.contains("When ready Run openlatch update.")));
}
fn declined(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 2, |_| {
let check = available(Severity::Critical);
let release =
close_check(check, "0.5.7", "https://registry.npmjs.org").expect("a release");
let (title, _, _) = question(&release, true);
ui::stage_named("Update", &title, &title);
answered(false);
finish(&Outcome::Declined { release }, &auto_on(), now())
})
}
#[test]
fn declined_in_every_mode() {
let (live, exit) = declined(Mode::Live);
assert_eq!(exit, 7);
assert_eq!(
live[..8],
[
"┌ Updating OpenLatch",
"│",
"◇ Checked for updates · 0.5.7 → 0.5.8 · security fix",
"│",
"◇ Update to 0.5.8 · not now",
"│",
"└ Nothing changed.",
"",
]
);
assert!(live[10].starts_with(" │ ! Not updated · still on 0.5.7"));
let card = live.join("\n");
assert!(card.contains(" Available 0.5.8 · security fix"));
assert!(card.contains(" When ready Run openlatch update again."));
let (append, _) = declined(Mode::Append);
assert!(append.contains(&"+ Update to 0.5.8 - not now".to_string()));
assert!(append
.iter()
.any(|l| l.starts_with(" | ! Not updated - still on 0.5.7")));
let (plain, _) = declined(Mode::Plain);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/2] Checked for updates - 0.5.7 -> 0.5.8 - security fix",
"[2/2] Update to 0.5.8 - not now",
"Automatic updates: on - last checked 4h ago",
"result=declined from=0.5.7 to=0.5.8 exit=7",
]
);
}
fn check_failed(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 1, |_| {
let check = CheckResult::Failed {
reason: "manifest fetch: send: operation timed out".into(),
};
let outcome = close_check(check, "0.5.7", "https://registry.npmjs.org")
.expect_err("nothing to install");
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn check_failed_in_every_mode() {
let (live, exit) = check_failed(Mode::Live);
assert_eq!(exit, 7);
assert_eq!(
live[..6],
[
"┌ Updating OpenLatch",
"│",
"✗ Check failed · couldn't reach the update server",
"│",
"└ Stopped.",
"",
]
);
assert!(live[8].starts_with(" │ ! Couldn't check for updates"));
let card = live.join("\n");
assert!(card.contains(" Still running 0.5.7 "), "{card}");
assert!(card.contains(" · nothing changed, your agents stay protected"));
assert!(card.contains(" Why registry.npmjs.org didn't answer"));
assert!(card.contains(" Try Check the connection or proxy, then run"));
let (append, _) = check_failed(Mode::Append);
assert!(append.contains(&"x Check failed - couldn't reach the update server".into()));
let (plain, _) = check_failed(Mode::Plain);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/1] Check failed - couldn't reach the update server",
"Still running 0.5.7 - nothing changed. registry.npmjs.org didn't answer.",
"Automatic updates: on - last checked 4h ago",
"result=check_failed version=0.5.7 reason=registry-unreachable exit=7",
]
);
}
fn failed_verify(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 6, |t0| {
let mut apply = check_and_download(t0, Severity::Normal, false);
let outcome = apply.fail(FailedAt::Stage(ApplyStage::Verify), ms(t0, 3500));
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn failed_verify_in_every_mode() {
let (live, exit) = failed_verify(Mode::Live);
assert_eq!(exit, 1);
assert_eq!(
live[..10],
[
"┌ Updating OpenLatch",
"│",
"◇ Checked for updates · 0.5.7 → 0.5.8",
"│",
"◇ Downloaded · 8.4 MB in 2.9s",
"│",
"✗ Verify failed · the signature doesn't match OpenLatch's release key",
"│",
"└ Stopped.",
"",
]
);
assert!(live[12].starts_with(" │ ✗ Update stopped: the download didn't verify"));
let card = live.join("\n");
assert!(card.contains(" Still running 0.5.7 · nothing was changed"));
assert!(card.contains(" What happened The signature doesn't match OpenLatch's key."));
assert!(card.contains(" OL-1504"));
assert!(card.contains(" Try Run openlatch update again. If it fails twice,"));
let (append, _) = failed_verify(Mode::Append);
assert!(append.contains(
&"x Verify failed - the signature doesn't match OpenLatch's release key".into()
));
let (plain, _) = failed_verify(Mode::Plain);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/6] Checked for updates - 0.5.7 -> 0.5.8",
"[2/6] Downloading 8.4 MB...",
"[2/6] Downloaded - 8.4 MB in 2.9s",
"[3/6] Verifying the signature - failed",
" The signature doesn't match OpenLatch's release key (OL-1504)",
"Still running 0.5.7 - nothing was changed. The signature doesn't match (OL-1504).",
"Automatic updates: on - last checked 4h ago",
"result=failed step=verify code=OL-1504 version=0.5.7 exit=1",
]
);
}
fn rolled_back(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 6, |t0| {
let mut apply = check_and_download(t0, Severity::Normal, false);
install_and_restart(&mut apply, t0);
let outcome = apply.rolled_back();
finish(&outcome, &auto_on(), now())
})
}
#[test]
fn rolled_back_in_every_mode() {
let (live, exit) = rolled_back(Mode::Live);
assert_eq!(exit, 1);
let at = live
.iter()
.position(|l| l == "✗ Restart failed · 0.5.8 didn't come back")
.unwrap_or_else(|| panic!("{live:#?}"));
assert_eq!(
live[at + 1..at + 6],
[
"│",
"◇ Rolled back · 0.5.7 restored and running",
"│",
"└ Stopped.",
"",
]
);
assert!(live[at + 8].starts_with(" │ ✗ 0.5.8 didn't start, so 0.5.7 is back"));
let card = live.join("\n");
assert!(
card.contains(" Running 0.5.7 · restored automatically, agents protected")
);
assert!(card.contains(" What happened The 0.5.8 service stopped right after starting."));
assert!(card.contains(" OL-1502"));
assert!(card.contains(" Try Run openlatch doctor, then share the log below."));
let (append, _) = rolled_back(Mode::Append);
assert!(append.contains(&"x Restart failed - 0.5.8 didn't come back".into()));
assert!(append.contains(&"+ Rolled back - 0.5.7 restored and running".into()));
let (plain, _) = rolled_back(Mode::Plain);
assert_eq!(
plain[4..],
[
"[3/6] Verified - signed by OpenLatch",
"[4/6] Installed - 0.5.7 kept as a backup",
"[5/6] Restarting the background service...",
"[5/6] Restarting the background service - failed",
" 0.5.8 didn't come back (OL-1502)",
"[6/6] Rolled back - 0.5.7 restored and running",
"Running 0.5.7 again - 0.5.8 stopped right after starting and was rolled back (OL-1502).",
"Automatic updates: on - last checked 4h ago",
"result=rolled_back from=0.5.8 to=0.5.7 code=OL-1502 exit=1",
]
);
}
fn no_daemon_updated(mode: Mode) -> (Vec<String>, i32) {
drive(mode, 4, |t0| {
let mut apply = check_and_download(t0, Severity::Normal, false);
apply.reach(Step::of(ApplyStage::Swap), ms(t0, 4000));
apply.finish(ms(t0, 4200));
let outcome = Outcome::Updated {
release: apply.release.clone(),
service: Service::NotRunning,
};
let auto = AutoUpdate {
on: false,
..auto_on()
};
finish(&outcome, &auto, now())
})
}
#[test]
fn no_daemon_updated_in_every_mode() {
let (live, exit) = no_daemon_updated(Mode::Live);
assert_eq!(exit, 7);
assert!(live.contains(&"◇ Installed · 0.5.7 kept as a backup".to_string()));
assert!(!live.iter().any(|l| l.contains("Restart")), "{live:#?}");
assert!(live.contains(&"└ Done.".to_string()));
assert!(live
.iter()
.any(|l| l.starts_with(" │ ! OpenLatch 0.5.8 is installed")));
let n = live.len();
assert_eq!(
live[n - 2..],
[
" ! The background service isn't running · start it: openlatch start",
" ! Automatic updates: off · turn on: [update] auto_update = true",
]
);
let (append, _) = no_daemon_updated(Mode::Append);
assert!(append.contains(&"+ Installed - 0.5.7 kept as a backup".into()));
assert!(append
.iter()
.any(|l| l == " ! The background service isn't running - start it: openlatch start"));
let (plain, exit) = no_daemon_updated(Mode::Plain);
assert_eq!(exit, 7);
assert_eq!(
plain,
[
"Updating OpenLatch",
"[1/4] Checked for updates - 0.5.7 -> 0.5.8",
"[2/4] Downloading 8.4 MB...",
"[2/4] Downloaded - 8.4 MB in 2.9s",
"[3/4] Verified - signed by OpenLatch",
"[4/4] Installed - 0.5.7 kept as a backup",
"Installed 0.5.8 - the background service isn't running. Start it: openlatch start",
"! Automatic updates: off - turn on: [update] auto_update = true",
r#"result=updated from=0.5.7 to=0.5.8 hint="openlatch start" exit=7"#,
]
);
}
#[test]
fn a_daemon_without_byte_counts_draws_no_meter() {
let (append, _) = drive(Mode::Append, 6, |t0| {
let mut release = release(Severity::Normal);
release.size = None;
let mut apply = Apply::new(release, t0);
apply.reach(Step::Download, t0);
apply.bytes(None, None, ms(t0, 500));
apply.reach(Step::Verify, ms(t0, 2000));
0
});
assert_eq!(
append[4..],
[
"|",
"> Downloading 0.5.8...",
"+ Downloaded - in 2.0s",
"|",
"> Verifying the signature...",
]
);
}
}
#[cfg(test)]
mod flow_tests {
use super::*;
use clap::Parser;
fn release(severity: Severity, size: Option<u64>) -> Release {
Release {
from: "0.5.7".into(),
to: "0.5.8".into(),
severity,
size,
}
}
#[test]
fn pipeline_stages_map_onto_the_rail() {
use ApplyStage::*;
let cases = [
(Check, Step::Download),
(Download, Step::Download),
(Extract, Step::Verify),
(Verify, Step::Verify),
(Sanity, Step::Verify),
(Reader, Step::Verify),
(Swap, Step::Install),
(Drain, Step::Restart),
(Restart, Step::Restart),
(Healthz, Step::Confirm),
];
for (stage, step) in cases {
assert_eq!(Step::of(stage), step, "{stage:?}");
assert_eq!(apply_stage(stage.as_str()), Some(stage));
}
assert_eq!(apply_stage("some_future_stage"), None);
}
#[test]
fn the_question_follows_severity_size_and_daemon() {
let (title, q, line) = question(&release(Severity::Normal, Some(8_806_432)), true);
assert_eq!(title, "Update to 0.5.8");
assert_eq!(q, "Update to 0.5.8 now?");
assert_eq!(
line.as_deref(),
Some("About 8.4 MB · the background service restarts for a few seconds.")
);
let (_, q, line) = question(&release(Severity::Critical, Some(8_806_432)), false);
assert_eq!(q, "0.5.8 fixes a security issue. Update now?");
assert_eq!(line.as_deref(), Some("About 8.4 MB."));
let (_, _, line) = question(&release(Severity::Normal, None), false);
assert_eq!(line, None);
}
#[test]
fn exit_codes_follow_the_health_convention() {
let r = || release(Severity::Normal, None);
let failed = |at| Outcome::Failed {
release: r(),
at,
step: Step::Verify,
};
let updated = |service| Outcome::Updated {
release: r(),
service,
};
let cases = [
(
Outcome::UpToDate {
version: "0.5.8".into(),
stale: None,
},
0,
),
(
Outcome::UpToDate {
version: "0.5.8".into(),
stale: Some("0.5.7".into()),
},
7,
),
(updated(Service::Running), 0),
(updated(Service::NotRunning), 7),
(updated(Service::Stale("0.5.7".into())), 7),
(
Outcome::Available {
release: r(),
check_only: true,
},
7,
),
(Outcome::Declined { release: r() }, 7),
(
Outcome::CheckFailed {
from: "0.5.7".into(),
host: "registry.npmjs.org".into(),
unreachable: true,
},
7,
),
(failed(FailedAt::Stage(ApplyStage::Verify)), 1),
(failed(FailedAt::NotBack), 1),
(Outcome::RolledBack { release: r() }, 1),
];
for (outcome, exit) in cases {
assert_eq!(outcome.exit_code(), exit, "{outcome:?}");
}
}
#[test]
fn failures_name_their_stage_and_code() {
let r = release(Severity::Normal, None);
let why = |at: FailedAt| at.why(&r, None);
let verify = why(FailedAt::Stage(ApplyStage::Verify));
assert_eq!(
(verify.step, verify.code, verify.unchanged),
(Step::Verify, "OL-1504", true)
);
let reader = why(FailedAt::Stage(ApplyStage::Reader));
assert_eq!((reader.step, reader.code), (Step::Verify, "OL-1216"));
let swap = why(FailedAt::Stage(ApplyStage::Swap));
assert_eq!((swap.step, swap.unchanged), (Step::Install, false));
let gone = why(FailedAt::NotBack);
assert_eq!((gone.step, gone.code), (Step::Restart, "OL-1502"));
assert_eq!(gone.line, "0.5.8 didn't come back");
let late = FailedAt::Stage(ApplyStage::Download).why(&r, Some(Step::Install));
assert_eq!(late.step, Step::Install);
assert_eq!(
FailedAt::Stalled.why(&r, Some(Step::Verify)).step,
Step::Verify
);
}
#[test]
fn the_meter_reads_bytes_speed_and_time_left() {
let mib = 1_048_576;
assert_eq!(
download_tail(5 * mib, Some(8 * mib), Some(3.0 * mib as f64)),
(Some(0.625), "5.0 / 8.0 MB · 3.0 MB/s · 1s left".to_string())
);
assert_eq!(
download_tail(mib, Some(300 * mib), Some(mib as f64)),
(
Some(1.0 / 300.0),
"1.0 / 300.0 MB · 1.0 MB/s · 5m left".to_string()
)
);
assert_eq!(download_tail(mib, None, None), (None, "1.0 MB".to_string()));
assert_eq!(
download_tail(0, Some(0), None),
(None, "0.0 MB".to_string())
);
let t0 = Instant::now();
let mut rate = Rate::default();
rate.sample(t0, 0);
rate.sample(t0 + Duration::from_millis(100), 10);
assert_eq!(rate.bytes_per_sec, None, "too soon to tell");
rate.sample(t0 + Duration::from_secs(1), 1000);
assert_eq!(rate.bytes_per_sec, Some(1000.0));
rate.sample(t0 + Duration::from_secs(2), 1500);
assert_eq!(rate.bytes_per_sec, Some(0.7 * 1000.0 + 0.3 * 500.0));
}
#[test]
fn automatic_updates_line() {
let now: chrono::DateTime<chrono::Utc> = "2026-09-23T12:00:00Z".parse().unwrap();
let on = AutoUpdate {
on: true,
by_env: false,
last_check: None,
};
assert_eq!(on.line(now), (false, "Automatic updates: on".to_string()));
let checked = AutoUpdate {
last_check: Some("2026-09-23T11:48:00Z".into()),
..on.clone()
};
assert_eq!(
checked.line(now).1,
"Automatic updates: on · last checked 12m ago"
);
let off = AutoUpdate {
on: false,
..on.clone()
};
assert_eq!(
off.line(now),
(
true,
"Automatic updates: off · turn on: `[update] auto_update = true`".to_string()
)
);
let off_by_env = AutoUpdate {
by_env: true,
..off
};
assert_eq!(
off_by_env.plain_line(now),
"! Automatic updates: off · turn on: OPENLATCH_AUTO_UPDATE=true"
);
assert_eq!(ago(chrono::Duration::seconds(5)), "just now");
assert_eq!(ago(chrono::Duration::hours(50)), "2d ago");
}
#[test]
fn daemon_rejections_keep_their_exit_codes() {
use reqwest::StatusCode;
let up = rejection(
StatusCode::CONFLICT,
&json!({"idempotent": true, "current": "0.5.8"}),
"0.5.7",
);
assert!(matches!(up, Rejection::UpToDate(v) if v == "0.5.8"));
let exit = |status, body| match rejection(status, &body, "0.5.7") {
Rejection::Refused(e, exit) => (e.code, exit),
Rejection::UpToDate(_) => ("", 0),
};
assert_eq!(exit(StatusCode::CONFLICT, json!({})), ("OL-1505", 5));
assert_eq!(
exit(StatusCode::SERVICE_UNAVAILABLE, json!({})),
("OL-1502", 1)
);
assert_eq!(exit(StatusCode::BAD_GATEWAY, json!({})), ("OL-1504", 1));
}
#[test]
fn apply_is_a_hidden_alias_of_yes() {
let parse = |args: &[&str]| {
let argv = ["openlatch", "update"].iter().chain(args).copied();
let cli = crate::cli::Cli::try_parse_from(argv).expect("parses");
match cli.command {
Some(crate::cli::Commands::Update(a)) => a,
_ => panic!("not update"),
}
};
assert!(!parse(&[]).applies());
assert!(parse(&["--yes"]).applies());
assert!(parse(&["-y"]).applies());
assert!(parse(&["--apply"]).applies());
assert!(parse(&["--apply", "--yes"]).applies(), "existing scripts");
assert!(!parse(&["--check"]).applies());
assert!(!parse(&["--check", "--yes"]).applies());
let mut cli = <crate::cli::Cli as clap::CommandFactory>::command();
let help = cli
.find_subcommand_mut("update")
.expect("update")
.render_long_help()
.to_string();
assert!(
help.contains("--yes") && !help.contains("--apply"),
"{help}"
);
}
#[test]
fn versions_match_exactly_up_to_a_build_suffix() {
assert!(is_version("0.5.8", "0.5.8"));
assert!(is_version("0.5.8+abc123", "0.5.8"));
assert!(
!is_version("0.5.8-rc.1", "0.5.8"),
"a pre-release is not the release"
);
assert!(!is_version("0.5.9-dev.12+g73ec01c", "0.5.9"));
assert!(!is_version("0.5.10", "0.5.1"), "0.5.1 → 0.5.10");
assert!(!is_version("0.5.1", "0.5.10"));
assert!(!is_version("0.5.7", "0.5.8"));
}
#[test]
fn host_of_the_registry() {
assert_eq!(host("https://registry.npmjs.org"), "registry.npmjs.org");
assert_eq!(host("http://127.0.0.1:4873/npm/"), "127.0.0.1:4873");
assert_eq!(host("registry.local"), "registry.local");
}
}
#[cfg(test)]
mod daemon_tests {
use super::*;
use crate::cli::ui::{Rail, INSTALL_LOCK};
use std::collections::{HashMap, VecDeque};
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::Mutex;
#[derive(Clone)]
enum Reply {
Json(u16, serde_json::Value),
Down,
}
use Reply::Down;
fn ok(body: serde_json::Value) -> Reply {
Reply::Json(200, body)
}
fn health(version: &str) -> Reply {
ok(json!({"status": "ok", "version": version}))
}
fn applying(stage: &str) -> Reply {
ok(json!({"status": "in_progress", "stage": stage, "bytes_done": 10, "bytes_total": 20}))
}
fn accepted(from: &str) -> Reply {
Reply::Json(202, json!({"started": true, "from": from, "to": "0.5.8"}))
}
const FAST: Timing = Timing {
poll: Duration::from_millis(5),
request: Duration::from_secs(2),
back_within: Duration::from_secs(3),
stall: Duration::from_secs(5),
misses: 3,
};
fn serve(routes: Vec<(&'static str, Vec<Reply>)>) -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").expect("binds");
let port = listener.local_addr().expect("addr").port();
let scripts: Mutex<HashMap<&str, VecDeque<Reply>>> = Mutex::new(
routes
.into_iter()
.map(|(path, replies)| (path, replies.into()))
.collect(),
);
std::thread::spawn(move || {
for stream in listener.incoming() {
let Ok(mut stream) = stream else { continue };
let path = read_request(&mut stream);
let reply = {
let mut scripts = scripts.lock().unwrap();
let script = scripts.get_mut(path.as_str());
match script {
Some(q) if q.len() > 1 => q.pop_front(),
Some(q) => q.front().cloned(),
None => None,
}
};
match reply {
Some(Reply::Json(code, body)) => {
let body = body.to_string();
let _ = write!(
stream,
"HTTP/1.1 {code} X\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
);
}
Some(Down) | None => {}
}
}
});
port
}
fn read_request(stream: &mut TcpStream) -> String {
let _ = stream.set_read_timeout(Some(Duration::from_secs(2)));
let mut buf = Vec::new();
let mut chunk = [0u8; 1024];
let end = loop {
if let Some(end) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
break end + 4;
}
match stream.read(&mut chunk) {
Ok(0) | Err(_) => return String::new(),
Ok(n) => buf.extend_from_slice(&chunk[..n]),
}
};
let head = String::from_utf8_lossy(&buf[..end]).to_string();
let length = head
.lines()
.find_map(|l| {
let (k, v) = l.split_once(':')?;
k.eq_ignore_ascii_case("content-length")
.then(|| v.trim().parse::<usize>().ok())?
})
.unwrap_or(0);
while buf.len() < end + length {
match stream.read(&mut chunk) {
Ok(0) | Err(_) => break,
Ok(n) => buf.extend_from_slice(&chunk[..n]),
}
}
head.split_whitespace().nth(1).unwrap_or("").to_string()
}
fn release() -> Release {
Release {
from: "0.5.7".into(),
to: "0.5.8".into(),
severity: Severity::Normal,
size: Some(20),
}
}
fn on_rail<T>(f: impl std::future::Future<Output = T>) -> (T, Vec<String>) {
let _l = INSTALL_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (rail, screen) = Rail::buffer(Caps::new(Mode::Plain, false, false, 80), 6);
let out = {
let _guard = ui::install_rail(rail).expect("installs");
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(f)
};
let lines = screen.lock().unwrap().clone();
(out, lines)
}
fn apply_through(routes: Vec<(&'static str, Vec<Reply>)>) -> (Outcome, Vec<String>) {
let port = serve(routes);
let (out, lines) = on_rail(via_daemon(release(), port, "token", false, FAST));
match out {
Ok(outcome) => (outcome, lines),
Err(stop) => panic!("stopped: {} {lines:#?}", stop.error.message),
}
}
#[test]
fn a_missed_poll_mid_download_is_not_a_restart() {
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.6")]),
(
"/admin/update/status",
vec![
applying("download"),
Down,
applying("download"),
applying("verify"),
applying("swap"),
applying("drain"),
Down,
ok(json!({"status": "completed"})),
],
),
("/health", vec![Down, Down, Down, health("0.5.8")]),
]);
let Outcome::Updated { release, service } = outcome else {
panic!("{outcome:?} {lines:#?}");
};
assert_eq!(service, Service::Running);
assert_eq!(release.from, "0.5.6", "the daemon's own version");
assert!(
lines.iter().any(|l| l == "[5/6] Confirmed - running 0.5.8"),
"{lines:#?}"
);
assert!(
!lines.iter().any(|l| l.contains("Rolled back")),
"{lines:#?}"
);
}
#[test]
fn the_old_version_still_applying_is_resumed_not_rolled_back() {
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.7")]),
(
"/admin/update/status",
vec![
applying("download"),
Down,
Down,
Down,
applying("download"),
applying("swap"),
applying("drain"),
Down,
ok(json!({"status": "completed"})),
],
),
(
"/health",
vec![
Down,
Down,
Down,
health("0.5.7"),
Down,
Down,
health("0.5.8"),
],
),
]);
assert!(
matches!(outcome, Outcome::Updated { .. }),
"{outcome:?} {lines:#?}"
);
assert!(
!lines.iter().any(|l| l.contains("Rolled back")),
"{lines:#?}"
);
}
#[test]
fn idle_on_the_old_version_after_the_restart_is_a_rollback() {
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.7")]),
(
"/admin/update/status",
vec![
applying("download"),
applying("swap"),
applying("drain"),
Down,
ok(json!({"status": "idle"})),
],
),
("/health", vec![Down, Down, health("0.5.7")]),
]);
assert!(
matches!(outcome, Outcome::RolledBack { .. }),
"{outcome:?} {lines:#?}"
);
assert!(lines.contains(&" 0.5.8 didn't come back (OL-1502)".to_string()));
assert!(lines
.iter()
.any(|l| l.ends_with("] Rolled back - 0.5.7 restored and running")));
}
#[test]
fn idle_without_a_drop_is_a_restart_the_version_decides() {
let started = Instant::now();
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.7")]),
(
"/admin/update/status",
vec![applying("download"), ok(json!({"status": "idle"}))],
),
("/health", vec![health("0.5.7")]),
]);
assert!(
matches!(outcome, Outcome::RolledBack { .. }),
"{outcome:?} {lines:#?}"
);
assert!(started.elapsed() < FAST.back_within, "not a timeout");
}
#[test]
fn completed_without_a_swap_is_up_to_date() {
let started = Instant::now();
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.7")]),
(
"/admin/update/status",
vec![ok(json!({"status": "completed", "stage": null}))],
),
("/health", vec![health("0.5.7")]),
]);
assert!(
matches!(&outcome, Outcome::UpToDate { version, stale: None } if version == "0.5.7"),
"{outcome:?} {lines:#?}"
);
assert!(started.elapsed() < FAST.back_within);
}
#[test]
fn a_dev_build_compares_the_daemon_against_its_own_identity() {
let dev = "0.5.9-dev.12+g73ec01c";
let port = serve(vec![("/health", vec![health(dev)])]);
let (outcome, _) = on_rail(up_to_date("0.5.8".into(), dev, port));
assert!(
matches!(&outcome, Outcome::UpToDate { stale: None, .. }),
"{outcome:?}"
);
let port = serve(vec![("/health", vec![health("0.5.9-dev.11+g0a1b2c3")])]);
let (outcome, _) = on_rail(up_to_date("0.5.8".into(), dev, port));
assert!(
matches!(&outcome, Outcome::UpToDate { stale: Some(v), .. } if v == "0.5.9-dev.11+g0a1b2c3"),
"{outcome:?}"
);
let r = release();
assert!(serves_from(dev, &r, dev));
assert!(serves_from("0.5.7", &r, dev));
assert!(!serves_from("0.5.9-dev.11+g0a1b2c3", &r, dev));
assert!(!serves_from("0.5.8", &r, dev));
}
#[test]
fn up_to_date_with_a_stale_daemon_needs_a_restart() {
let port = serve(vec![("/health", vec![health("0.5.7")])]);
let (outcome, lines) = on_rail(async {
let outcome = up_to_date("0.5.8".into(), "0.5.8", port).await;
let exit = finish(&outcome, &AutoUpdate::default(), chrono::Utc::now());
(outcome, exit)
});
let (outcome, exit) = outcome;
assert!(
matches!(&outcome, Outcome::UpToDate { stale: Some(v), .. } if v == "0.5.7"),
"{outcome:?}"
);
assert_eq!(exit, 7);
assert_eq!(
lines[lines.len() - 3..],
[
"Installed 0.5.8 - the background service is still on 0.5.7. Run: openlatch restart",
"! Automatic updates: off - turn on: [update] auto_update = true",
r#"result=up_to_date version=0.5.8 hint="openlatch restart" exit=7"#,
]
);
let (card, warnings) = outcome.card(&AutoUpdate::default(), chrono::Utc::now(), None);
assert_eq!(card.verdict, ui::Verdict::NeedsAttention);
assert_eq!(card.title, "OpenLatch 0.5.8 is installed");
assert_eq!(
warnings[0],
"The background service is still on 0.5.7 · restart it: `openlatch restart`"
);
let current = serve(vec![("/health", vec![health("0.5.8")])]);
let ((same, none), _) = on_rail(async {
let closed = TcpListener::bind("127.0.0.1:0")
.and_then(|l| l.local_addr())
.expect("a free port")
.port();
(
up_to_date("0.5.8".into(), "0.5.8", current).await,
up_to_date("0.5.8".into(), "0.5.8", closed).await,
)
});
for outcome in [same, none] {
assert!(
matches!(outcome, Outcome::UpToDate { stale: None, .. }),
"{outcome:?}"
);
assert_eq!(outcome.exit_code(), 0);
}
}
#[test]
fn the_restart_is_timed_from_the_drop() {
let mut health_script = vec![Down; 3 + 40];
health_script.push(health("0.5.8"));
let (outcome, lines) = apply_through(vec![
("/admin/update", vec![accepted("0.5.7")]),
(
"/admin/update/status",
vec![
applying("download"),
Down,
Down,
Down,
ok(json!({"status": "completed"})),
],
),
("/health", health_script),
]);
assert!(
matches!(outcome, Outcome::Updated { .. }),
"{outcome:?} {lines:#?}"
);
let back = lines
.iter()
.find_map(|l| {
l.split("Restarted the background service - back in ")
.nth(1)
})
.unwrap_or_else(|| panic!("{lines:#?}"));
let secs: f64 = back.trim_end_matches('s').parse().expect("seconds");
assert!(secs >= 0.1, "timed from the drop: {back} {lines:#?}");
}
#[test]
fn no_daemon_is_found_fast() {
let closed = TcpListener::bind("127.0.0.1:0")
.and_then(|l| l.local_addr())
.expect("a free port")
.port();
let started = Instant::now();
let (outcome, _) = on_rail(up_to_date("0.5.8".into(), "0.5.8", closed));
assert!(matches!(outcome, Outcome::UpToDate { stale: None, .. }));
assert!(
started.elapsed() < Duration::from_millis(1500),
"{:?}",
started.elapsed()
);
}
#[test]
fn a_refusal_before_the_check_writes_its_result_line() {
let (exit, lines) = on_rail(async {
ui::intro("Updating OpenLatch");
stop(&Stop::new(cargo_refusal(CARGO_SUGGESTION), 5))
});
assert_eq!(exit, 5);
assert_eq!(lines[0], "Updating OpenLatch");
assert!(lines[1].ends_with("(OL-1505)"), "{lines:#?}");
assert_eq!(
lines.last().map(String::as_str),
Some("result=failed code=OL-1505 exit=5")
);
}
}