1use std::io::{Read, Seek, SeekFrom, Write};
27use std::path::{Path, PathBuf};
28use std::sync::OnceLock;
29use std::time::{Duration, Instant};
30
31use serde_json::{Value, json};
32
33use super::bootstrap::{AGENT_HOME, AGENT_UNIT, AGENT_USER, elf_arch, sha256_hex};
34use super::{AgentClient, Servers, health, store};
35use crate::error::{Error, Result};
36use crate::exec::{ExecEvent, ExecOptions, Stdin};
37use crate::sandbox::Sandbox;
38use crate::stack::now_secs;
39
40pub const PROTOCOL: u64 = 2;
44pub const MIN_PROTOCOL: u64 = 1;
46pub const SSH_PROTOCOL: u64 = 2;
48
49pub const CHUNK: usize = 3 << 20;
51pub const MAX_BINARY: u64 = 512 << 20;
53pub const CONFIRM_WINDOW_S: u64 = 120;
56const COME_BACK: Duration = Duration::from_secs(100);
59const ROLLBACK_SEEN: Duration = Duration::from_secs(90);
61
62pub const HELPER: &str = "/usr/local/lib/isb/agent-upgrade";
63pub const HELPER_UNIT: &str = "isb-agent-upgrade";
64
65pub fn dir(state_dir: &Path) -> PathBuf {
67 state_dir.join("upgrade")
68}
69
70fn box_dir() -> String {
72 format!("{AGENT_HOME}/state/upgrade")
73}
74
75pub fn build_id() -> &'static str {
78 static ID: OnceLock<String> = OnceLock::new();
79 ID.get_or_init(|| file_sha256(Path::new("/proc/self/exe")).unwrap_or_default())
80}
81
82fn file_sha256(p: &Path) -> std::io::Result<String> {
83 let mut f = std::fs::File::open(p)?;
84 let mut ctx = ring::digest::Context::new(&ring::digest::SHA256);
85 let mut buf = vec![0u8; 1 << 20];
86 loop {
87 let n = f.read(&mut buf)?;
88 if n == 0 {
89 break;
90 }
91 ctx.update(&buf[..n]);
92 }
93 Ok(crate::machine::hex(ctx.finish().as_ref()))
94}
95
96pub fn protocol_of(hb: &Value) -> u64 {
98 hb["protocol"].as_u64().unwrap_or(1)
99}
100
101pub fn compatible(name: &str, hb: &Value, need: u64) -> Result<()> {
104 if hb.is_null() {
105 return Ok(());
106 }
107 let p = protocol_of(hb);
108 let ver = hb["isb"].as_str().unwrap_or("unknown");
109 if p > PROTOCOL {
110 return Err(Error::invalid(format!(
111 "server {name} runs isb {ver} (agent protocol {p}), newer than this control plane understands (protocol {PROTOCOL}): upgrade the control plane"
112 )));
113 }
114 if p < need {
115 return Err(Error::invalid(format!(
116 "server {name} runs isb {ver} (agent protocol {p}); this needs protocol {need}: upgrade it with `isb server upgrade {name}`"
117 )));
118 }
119 Ok(())
120}
121
122pub fn heartbeat_fields(state_dir: &Path) -> Value {
124 json!({
125 "protocol": PROTOCOL,
126 "build": build_id(),
127 "arch": std::env::consts::ARCH,
128 "upgrade": status(&dir(state_dir)),
129 })
130}
131
132pub fn status(dir: &Path) -> Value {
134 let helper = Path::new(&format!("/etc/systemd/system/{HELPER_UNIT}.path")).exists();
135 let last = std::fs::read(dir.join("result"))
136 .ok()
137 .and_then(|b| serde_json::from_slice::<Value>(&b).ok())
138 .unwrap_or(Value::Null);
139 json!({"helper": helper, "last": last})
140}
141
142pub fn helper_script() -> String {
148 let dir = box_dir();
149 let window = CONFIRM_WINDOW_S;
150 format!(
151 r#"#!/bin/sh
152# Installs the isb binary the control plane staged for this agent, restarts
153# the agent, and restores the previous binary unless the control plane
154# confirms within {window} s that the new agent answers. Written by isb;
155# run by {HELPER_UNIT}.service.
156set -u
157d={dir}
158bin=/usr/local/bin/isb
159[ -f "$d/request" ] || exit 0
160rm -f "$d/request" "$d/confirmed"
161want=$(tr -dc 0-9a-f < "$d/sha256" 2>/dev/null | head -c 64)
162result() {{
163 printf '{{"state":"%s","sha256":"%s","at":%s,"message":"%s"}}\n' "$1" "$want" "$(date +%s)" "$2" > "$d/result.new"
164 chown {user}:{user} "$d/result.new" 2>/dev/null
165 mv -f "$d/result.new" "$d/result"
166 echo "isb-agent-upgrade: $1: $2"
167}}
168[ "${{#want}}" -eq 64 ] || {{ result failed "no sha256 was staged"; exit 1; }}
169# Check our own copy, so what is checked is what is installed.
170tmp=$(mktemp /usr/local/bin/.isb.XXXXXX) || {{ result failed "mktemp failed"; exit 1; }}
171if ! cp "$d/isb.new" "$tmp" || ! chmod 0755 "$tmp"; then
172 rm -f "$tmp"; result failed "the staged binary could not be copied"; exit 1
173fi
174if ! echo "$want $tmp" | sha256sum -c --quiet - >/dev/null 2>&1; then
175 rm -f "$tmp"; result failed "the staged binary does not match its sha256"; exit 1
176fi
177if ! runuser -u {user} -- "$tmp" --version >/dev/null 2>&1; then
178 rm -f "$tmp"; result failed "the staged binary does not run"; exit 1
179fi
180cp -p "$bin" "$bin.prev" || {{ rm -f "$tmp"; result failed "could not keep the previous binary"; exit 1; }}
181mv -f "$tmp" "$bin"
182result restarting "installed; waiting for the control plane to confirm"
183systemctl reset-failed {unit} 2>/dev/null
184systemctl restart {unit}
185i=0
186while [ "$i" -lt {window} ]; do
187 if [ -f "$d/confirmed" ] && [ "$(tr -dc 0-9a-f < "$d/confirmed")" = "$want" ]; then
188 rm -f "$d/isb.new" "$d/isb.new.part" "$d/confirmed"
189 result done "upgraded"
190 exit 0
191 fi
192 sleep 1
193 i=$((i + 1))
194done
195cp -p "$bin.prev" "$bin.rollback" && mv -f "$bin.rollback" "$bin"
196# A binary that keeps exiting may have hit the unit's start limit.
197systemctl reset-failed {unit} 2>/dev/null
198systemctl restart {unit}
199result rolled_back "the new agent was not confirmed within {window} s: the previous binary is back"
200exit 1
201"#,
202 user = AGENT_USER,
203 unit = AGENT_UNIT,
204 )
205}
206
207pub fn helper_install() -> String {
211 let dir = box_dir();
212 let mut s = String::new();
213 s.push_str("install -d -m 0755 /usr/local/lib/isb\n");
214 s.push_str(&format!("cat > {HELPER}.new <<'ISB_HELPER_EOF'\n"));
215 s.push_str(&helper_script());
216 s.push_str("ISB_HELPER_EOF\n");
217 s.push_str(&format!(
218 "chmod 0755 {HELPER}.new && mv -f {HELPER}.new {HELPER}\n"
219 ));
220 s.push_str(&format!(
221 "cat > /etc/systemd/system/{HELPER_UNIT}.service <<'ISB_HELPER_EOF'
222[Unit]
223Description=Install the isb agent binary its control plane staged
224
225[Service]
226Type=oneshot
227ExecStart={HELPER}
228TimeoutStartSec={timeout}
229ISB_HELPER_EOF
230cat > /etc/systemd/system/{HELPER_UNIT}.path <<'ISB_HELPER_EOF'
231[Unit]
232Description=Watch for an isb agent upgrade its control plane staged
233
234[Path]
235PathExists={dir}/request
236
237[Install]
238WantedBy=multi-user.target
239ISB_HELPER_EOF
240[ -d {home}/state ] || install -d -o {user} -g {user} -m 0750 {home}/state
241install -d -o {user} -g {user} -m 0700 {dir}
242systemctl daemon-reload
243systemctl enable --now {HELPER_UNIT}.path >/dev/null 2>&1
244",
245 timeout = CONFIRM_WINDOW_S + 180,
246 home = AGENT_HOME,
247 user = AGENT_USER,
248 ));
249 s
250}
251
252pub fn stage_chunk(dir: &Path, offset: u64, body: &[u8]) -> Result<u64> {
257 std::fs::create_dir_all(dir)?;
258 let part = dir.join("isb.new.part");
259 let mut f = std::fs::OpenOptions::new()
260 .create(true)
261 .write(true)
262 .truncate(offset == 0)
263 .open(&part)?;
264 let len = f.metadata()?.len();
265 if offset != len {
266 return Err(Error::invalid(format!(
267 "upload: chunk at {offset}, but {len} bytes are staged"
268 )));
269 }
270 if offset + body.len() as u64 > MAX_BINARY {
271 return Err(Error::invalid("upload: larger than any isb binary"));
272 }
273 f.seek(SeekFrom::Start(offset))?;
274 f.write_all(body)?;
275 Ok(offset + body.len() as u64)
276}
277
278pub fn stage_apply(dir: &Path, sha256: &str, size: u64) -> Result<()> {
281 let part = dir.join("isb.new.part");
282 let len = std::fs::metadata(&part)
283 .map_err(|_| Error::invalid("upload: nothing staged"))?
284 .len();
285 if len != size {
286 return Err(Error::invalid(format!(
287 "upload: {len} bytes staged, {size} expected"
288 )));
289 }
290 let got = file_sha256(&part)?;
291 if !got.eq_ignore_ascii_case(sha256) {
292 return Err(Error::invalid(format!(
293 "upload: sha256 {got} does not match {sha256}"
294 )));
295 }
296 let mut head = [0u8; 20];
297 std::fs::File::open(&part)?.read_exact(&mut head)?;
298 match elf_arch(&head) {
299 Some(a) if a == std::env::consts::ARCH => {}
300 Some(a) => {
301 return Err(Error::invalid(format!(
302 "the binary is for {a}; this server is {}",
303 std::env::consts::ARCH
304 )));
305 }
306 None => return Err(Error::invalid("the binary is not a Linux executable")),
307 }
308 if status(dir)["helper"] != true {
309 return Err(Error::invalid(format!(
310 "this server has no upgrade helper ({HELPER_UNIT}.path): it was bootstrapped by an older isb; upgrade it by hand once (docs/operations/upgrades.md)"
311 )));
312 }
313 std::fs::rename(&part, dir.join("isb.new"))?;
314 for f in ["confirmed", "result"] {
315 let _ = std::fs::remove_file(dir.join(f));
316 }
317 std::fs::write(
318 dir.join("sha256"),
319 format!("{}\n", sha256.to_ascii_lowercase()),
320 )?;
321 std::fs::write(dir.join("request"), format!("{}\n", now_secs()))?;
322 Ok(())
323}
324
325pub fn confirm(dir: &Path, sha256: &str) -> Result<()> {
328 if !build_id().eq_ignore_ascii_case(sha256) {
329 return Err(Error::invalid(format!(
330 "this agent runs build {}, not {sha256}",
331 build_id()
332 )));
333 }
334 std::fs::write(
335 dir.join("confirmed"),
336 format!("{}\n", sha256.to_ascii_lowercase()),
337 )?;
338 Ok(())
339}
340
341#[derive(Debug, Clone)]
345pub enum Source {
346 Own,
348 Release(String),
350 File(PathBuf),
352}
353
354impl Servers {
355 pub fn upgrade(&self, client: &crate::Client, name: &str, source: &Source) -> Result<Value> {
358 let _busy = self.begin_upgrade(name)?;
359 let rec = self.record(name)?;
360 let c = self.client(name)?;
361 let hb = c
362 .internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT)
363 .map_err(|e| {
364 Error::invalid(format!(
365 "server {name} does not answer ({e}): an upgrade needs its agent running"
366 ))
367 })?;
368 let arch = match (hb["arch"].as_str(), &rec.vm) {
369 (Some(a), _) => a.to_string(),
370 (None, Some(_)) => std::env::consts::ARCH.to_string(),
372 (None, None) => {
373 return Err(Error::invalid(format!(
374 "server {name} runs isb {} without the upgrade routes: upgrade it by hand once (docs/operations/upgrades.md)",
375 hb["isb"].as_str().unwrap_or("?")
376 )));
377 }
378 };
379 let bin = self.binary(source, &arch)?;
380 let sha = sha256_hex(&bin);
381 let from = json!({"isb": hb["isb"], "build": hb["build"]});
382 if hb["build"].as_str() == Some(sha.as_str()) {
383 return Ok(
384 json!({"name": name, "upgraded": false, "from": from, "to": from,
385 "note": "already runs this build"}),
386 );
387 }
388 match &rec.vm {
389 Some(v) => stage_vm(client, v, &bin, &sha)?,
390 None => stage_mtls(&c, &bin, &sha)?,
391 }
392 let now = match come_back(&c, &sha) {
393 Ok(hb) => hb,
394 Err(why) => {
395 let (e, seen) = rollback_error(&c, name, &sha, why);
396 if let Some(hb) = seen {
397 self.observe(name, hb);
398 }
399 return Err(e);
400 }
401 };
402 c.internal(
403 "POST",
404 "/internal/v1/upgrade/confirm",
405 Some(&json!({"sha256": sha})),
406 health::TIMEOUT,
407 )?;
408 let to = json!({"isb": now["isb"], "build": now["build"]});
409 self.upgraded(name, &now);
410 Ok(json!({"name": name, "upgraded": true, "from": from, "to": to}))
411 }
412
413 fn begin_upgrade(&self, name: &str) -> Result<UpgradeGuard<'_>> {
415 if !self.upgrading.lock().unwrap().insert(name.to_string()) {
416 return Err(Error::invalid(format!(
417 "server {name} is being upgraded already"
418 )));
419 }
420 Ok(UpgradeGuard {
421 set: &self.upgrading,
422 name: name.to_string(),
423 })
424 }
425
426 fn binary(&self, source: &Source, arch: &str) -> Result<Vec<u8>> {
427 let b = match source {
428 Source::Own => return super::bootstrap::own_binary(arch),
429 Source::File(f) => {
430 std::fs::read(f).map_err(|e| Error::invalid(format!("{}: {e}", f.display())))?
431 }
432 Source::Release(v) => {
433 let scratch = self.dir.join(format!("tmp-upgrade-{}", now_secs()));
434 std::fs::create_dir_all(&scratch)?;
435 let dst = scratch.join("isb");
436 let r = crate::machine::download_release(v, arch, &scratch, &dst)
437 .and_then(|()| Ok(std::fs::read(&dst)?));
438 let _ = std::fs::remove_dir_all(&scratch);
439 r?
440 }
441 };
442 match elf_arch(&b) {
443 Some(a) if a == arch => Ok(b),
444 Some(a) => Err(Error::invalid(format!(
445 "the isb binary is for {a}; the server is {arch}"
446 ))),
447 None => Err(Error::invalid("the isb binary is not a Linux executable")),
448 }
449 }
450
451 fn upgraded(&self, name: &str, hb: &Value) {
453 {
454 let mut st = self.store.lock().unwrap();
455 if let Some(r) = st.servers.get_mut(name) {
456 r.isb_version = hb["isb"].as_str().unwrap_or("").to_string();
457 let _ = st.save();
458 }
459 }
460 self.observe(name, hb.clone());
461 }
462
463 fn observe(&self, name: &str, hb: Value) {
465 self.health
466 .lock()
467 .unwrap()
468 .entry(name.to_string())
469 .or_default()
470 .observe(Ok(hb), now_secs());
471 }
472}
473
474struct UpgradeGuard<'a> {
475 set: &'a std::sync::Mutex<std::collections::BTreeSet<String>>,
476 name: String,
477}
478
479impl Drop for UpgradeGuard<'_> {
480 fn drop(&mut self) {
481 self.set.lock().unwrap().remove(&self.name);
482 }
483}
484
485fn stage_mtls(c: &AgentClient, bin: &[u8], sha: &str) -> Result<()> {
487 let mut offset = 0usize;
488 for chunk in bin.chunks(CHUNK) {
489 c.internal_bytes(
490 &format!("/internal/v1/upgrade/chunk?offset={offset}"),
491 chunk,
492 Duration::from_secs(120),
493 )?;
494 offset += chunk.len();
495 }
496 c.internal(
497 "POST",
498 "/internal/v1/upgrade/apply",
499 Some(&json!({"sha256": sha, "size": bin.len()})),
500 Duration::from_secs(60),
501 )?;
502 Ok(())
503}
504
505fn stage_vm(client: &crate::Client, v: &store::VmRecord, bin: &[u8], sha: &str) -> Result<()> {
508 let sys = client.clone().project(&v.project);
509 let upload = "/root/isb-agent.upgrade";
510 sys.push_file(&v.instance, upload, bin, 0, 0, 0o600)?;
511 let dir = box_dir();
512 let script = format!(
513 "set -eu\n{helper}echo \"{sha} {upload}\" | sha256sum -c --quiet -\n\
514 install -o {user} -g {user} -m 0600 {upload} {dir}/isb.new\nrm -f {upload} {dir}/confirmed {dir}/result\n\
515 echo {sha} > {dir}/sha256\nchown {user}:{user} {dir}/sha256\ndate +%s > {dir}/request\n",
516 helper = helper_install(),
517 user = AGENT_USER,
518 );
519 let sb = Sandbox::get(&sys, &v.instance)?;
520 let mut s = sb.exec_stream(
521 ["sh", "-s"],
522 ExecOptions::default()
523 .env(
524 "PATH",
525 "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin",
526 )
527 .timeout(Duration::from_secs(120))
528 .stdin(Stdin::Bytes(script.into_bytes())),
529 )?;
530 let mut out = Vec::new();
531 while let Some(ev) = s.next_event() {
532 match ev {
533 ExecEvent::Stdout(b) | ExecEvent::Stderr(b) => out.extend(b),
534 }
535 }
536 let code = s.wait()?;
537 if code != 0 {
538 let text = String::from_utf8_lossy(&out);
539 return Err(Error::OperationFailed {
540 step: format!("stage the upgrade in VM {}", v.instance),
541 message: format!("exit {code}: {}", text.trim()),
542 });
543 }
544 Ok(())
545}
546
547fn come_back(c: &AgentClient, sha: &str) -> std::result::Result<Value, String> {
549 let started = Instant::now();
550 let mut last = String::from("no answer yet");
551 while started.elapsed() < COME_BACK {
552 match c.internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT) {
553 Ok(hb) if hb["build"].as_str() == Some(sha) => return Ok(hb),
554 Ok(hb) => {
555 last = match hb["upgrade"]["last"]["state"].as_str() {
556 Some("failed") if hb["upgrade"]["last"]["sha256"] == sha => {
557 return Err(format!(
559 "the server's upgrade helper refused it: {}",
560 hb["upgrade"]["last"]["message"].as_str().unwrap_or("")
561 ));
562 }
563 _ => format!(
564 "still answering with build {}",
565 hb["build"].as_str().unwrap_or("?")
566 ),
567 }
568 }
569 Err(e) => last = e.to_string(),
570 }
571 std::thread::sleep(Duration::from_secs(2));
572 }
573 Err(format!(
574 "the agent did not answer with the new build within {}s ({last})",
575 COME_BACK.as_secs()
576 ))
577}
578
579fn rollback_error(c: &AgentClient, name: &str, sha: &str, why: String) -> (Error, Option<Value>) {
583 let started = Instant::now();
584 let mut rolled = None;
585 if !why.contains("refused it") {
586 while started.elapsed() < ROLLBACK_SEEN {
587 if let Ok(hb) = c.internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT) {
588 let last = &hb["upgrade"]["last"];
589 if last["sha256"] == sha && last["state"] == "rolled_back" {
590 rolled = Some(hb);
591 break;
592 }
593 }
594 std::thread::sleep(Duration::from_secs(3));
595 }
596 }
597 let message = match &rolled {
598 Some(hb) => format!(
599 "{why}; the server restored its previous binary and answers again (isb {}, build {})",
600 hb["isb"].as_str().unwrap_or("?"),
601 short(hb["build"].as_str().unwrap_or("?"))
602 ),
603 None if why.contains("refused it") => why,
604 None => format!(
605 "{why}; its upgrade helper restores the previous binary {CONFIRM_WINDOW_S}s after the restart: check `isb server show {name}`"
606 ),
607 };
608 let e = Error::OperationFailed {
609 step: format!("upgrade server {name}"),
610 message,
611 };
612 (e, rolled)
613}
614
615fn short(b: &str) -> &str {
616 &b[..b.len().min(12)]
617}
618
619#[cfg(test)]
620mod tests {
621 use super::*;
622
623 #[test]
624 fn protocols_are_checked_both_ways() {
625 compatible("s", &Value::Null, SSH_PROTOCOL).unwrap();
626 compatible("s", &json!({"isb": "0.7.0"}), MIN_PROTOCOL).unwrap();
627 let old = compatible("s", &json!({"isb": "0.7.0"}), SSH_PROTOCOL).unwrap_err();
628 assert!(old.to_string().contains("isb server upgrade s"), "{old}");
629 let new =
630 compatible("s", &json!({"isb": "9.0.0", "protocol": PROTOCOL + 1}), 1).unwrap_err();
631 assert!(
632 new.to_string().contains("upgrade the control plane"),
633 "{new}"
634 );
635 compatible("s", &json!({"protocol": PROTOCOL}), SSH_PROTOCOL).unwrap();
636 }
637
638 #[test]
639 fn chunks_stage_in_order_and_apply_checks_the_hash() {
640 let d = tempfile::tempdir().unwrap();
641 let mut bin = b"\x7fELF\x02\x01\x01\0\0\0\0\0\0\0\0\0\x02\0".to_vec();
642 bin.extend(match std::env::consts::ARCH {
643 "aarch64" => [0xb7, 0],
644 _ => [0x3e, 0],
645 });
646 bin.extend(vec![7u8; 100]);
647 let sha = sha256_hex(&bin);
648 assert_eq!(stage_chunk(d.path(), 0, &bin[..50]).unwrap(), 50);
649 assert!(
650 stage_chunk(d.path(), 10, &bin[50..]).is_err(),
651 "out of order"
652 );
653 assert_eq!(
654 stage_chunk(d.path(), 50, &bin[50..]).unwrap(),
655 bin.len() as u64
656 );
657 assert!(stage_apply(d.path(), &sha, 5).is_err(), "wrong size");
658 let bad = stage_apply(d.path(), &"0".repeat(64), bin.len() as u64).unwrap_err();
659 assert!(bad.to_string().contains("does not match"), "{bad}");
660 if status(d.path())["helper"] != true {
662 let e = stage_apply(d.path(), &sha, bin.len() as u64).unwrap_err();
663 assert!(e.to_string().contains("upgrade helper"), "{e}");
664 }
665 assert_eq!(stage_chunk(d.path(), 0, b"x").unwrap(), 1);
667 assert!(confirm(d.path(), &"f".repeat(64)).is_err());
668 }
669
670 #[test]
671 fn the_helper_checks_restarts_and_rolls_back() {
672 let s = helper_script();
673 assert!(s.contains("sha256sum -c"));
674 assert!(s.contains("runuser -u isb --"));
675 assert!(s.contains("systemctl restart isb-agent.service"));
676 assert!(s.contains("mv -f \"$tmp\" \"$bin\""));
677 assert!(s.contains("rolled_back"));
678 assert!(s.contains(&format!("-lt {CONFIRM_WINDOW_S}")));
679 let i = helper_install();
680 assert!(i.contains("PathExists=/var/lib/isb/state/upgrade/request"));
681 assert!(i.contains("systemctl enable --now isb-agent-upgrade.path"));
682 assert_eq!(i.matches("\nISB_HELPER_EOF\n").count(), 3);
684 }
685}