From 6a342f665ce032f8553e92e72d3d52155cbc3a9e Mon Sep 17 00:00:00 2001 From: archipelago Date: Thu, 8 Oct 2026 05:20:40 -0400 Subject: [PATCH] Keep automatic recovery outside managed update ownership --- core/archipelago/src/crash_recovery.rs | 297 +++++++++++++++++- core/archipelago/src/health_monitor.rs | 14 +- core/archipelago/src/main.rs | 2 +- .../managed-update-recovery-implementation.md | 39 +++ 4 files changed, 338 insertions(+), 14 deletions(-) diff --git a/core/archipelago/src/crash_recovery.rs b/core/archipelago/src/crash_recovery.rs index 2765d0ac..a2cc1d30 100644 --- a/core/archipelago/src/crash_recovery.rs +++ b/core/archipelago/src/crash_recovery.rs @@ -483,9 +483,48 @@ pub async fn save_container_snapshot(data_dir: &Path) -> Result<()> { Ok(()) } +// Every automatic mutation shares lifecycle admission with managed updates. +// A stack is indivisible here: alias repair can affect a held sibling even if +// the particular stopped member has no saved unit of its own. +async fn automatic_recovery_intent_allowed(data_dir: &Path, names: &[&str]) -> bool { + if names + .iter() + .any(|name| !crate::health_monitor::automatic_recovery_allowed(data_dir, name)) + { + return false; + } + for file in [USER_STOPPED_FILE, USER_UNINSTALLED_FILE] { + let values: std::collections::HashSet = match fs::read(data_dir.join(file)).await { + Ok(bytes) => match serde_json::from_slice(&bytes) { + Ok(values) => values, + Err(_) => return false, + }, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Default::default(), + Err(_) => return false, + }; + if names.iter().any(|name| values.contains(*name)) { + return false; + } + } + true +} + +async fn automatic_recovery_admission( + data_dir: &Path, + names: &[&str], +) -> Option { + let guard = crate::container::update_transaction::Guard::acquire(data_dir).ok()?; + automatic_recovery_intent_allowed(data_dir, names) + .await + .then_some(guard) +} + /// Recover containers that were running before a crash. /// Attempts to start each container, logging success/failure. -pub async fn recover_containers(containers: &[RunningContainerRecord]) -> RecoveryReport { +pub async fn recover_containers( + data_dir: &Path, + containers: &[RunningContainerRecord], +) -> RecoveryReport { // Snapshot entries can outlive their containers (removed while we were // down, or podman storage partially reset by an unclean poweroff). // `podman start` on those fails permanently, and recovery runs BEFORE the @@ -517,12 +556,22 @@ pub async fn recover_containers(containers: &[RunningContainerRecord]) -> Recove pending_boot_starts_add(containers.iter().map(|r| r.name.clone())); for (i, record) in containers.iter().enumerate() { + let scope = stack_recovery_specs() + .iter() + .find(|stack| stack.containers.contains(&record.name.as_str())) + .map(|stack| stack.containers.to_vec()) + .unwrap_or_else(|| vec![record.name.as_str()]); + let Some(_admission) = automatic_recovery_admission(data_dir, &scope).await else { + pending_boot_start_done(&record.name); + continue; + }; // Skip containers that are already up — `podman start` on a running // container produces the noisy benign conmon "Failed to create // container" + cgroup Permission-denied journal pair (fleet log // sweep 2026-07-22). if container_state(&record.name).await.as_deref() == Some("running") { report.recovered += 1; + pending_boot_start_done(&record.name); continue; } info!( @@ -556,14 +605,11 @@ pub async fn recover_containers(containers: &[RunningContainerRecord]) -> Recove ); tokio::time::sleep(std::time::Duration::from_secs(10)).await; } - let mut cmd = tokio::process::Command::new("podman"); - cmd.args(["start", &record.name]); - let result = command_with_timeout( - cmd, - Duration::from_secs(timeout_secs), - &format!("podman start {}", record.name), - ) - .await; + if !automatic_recovery_intent_allowed(data_dir, &scope).await { + break; + } + let result = + podman_output(&["start", &record.name], Duration::from_secs(timeout_secs)).await; match result { Ok(output) if output.status.success() => { @@ -677,6 +723,13 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport { }; for stack in stack_recovery_specs() { + let Some(_admission) = automatic_recovery_admission(data_dir, stack.containers).await + else { + for name in stack.containers { + pending_boot_start_done(name); + } + continue; + }; if !stack_anchor_container_exists(stack).await { continue; } @@ -694,6 +747,9 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport { if pending.is_empty() { continue; } + if !automatic_recovery_intent_allowed(data_dir, stack.containers).await { + continue; + } info!("Recovering stopped {} stack containers", stack.name); repair_stack_network_aliases(stack).await; pending_boot_starts_add(pending.iter().cloned()); @@ -720,9 +776,21 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport { } } + if !automatic_recovery_intent_allowed(data_dir, stack.containers).await { + for name in &pending { + pending_boot_start_done(name); + } + break; + } repair_stack_network_aliases(stack).await; - wait_before_stack_container_recovery(stack, container).await; + wait_before_stack_container_recovery(data_dir, stack, container).await; + if !automatic_recovery_intent_allowed(data_dir, stack.containers).await { + for name in &pending { + pending_boot_start_done(name); + } + break; + } report.total += 1; if start_existing_container(container).await { report.recovered += 1; @@ -736,13 +804,20 @@ async fn start_stopped_app_stacks(data_dir: &Path) -> RecoveryReport { report } -async fn wait_before_stack_container_recovery(stack: &StackRecoverySpec, container: &str) { +async fn wait_before_stack_container_recovery( + data_dir: &Path, + stack: &StackRecoverySpec, + container: &str, +) { if stack.name != "indeedhub" || container != "indeedhub" { return; } for _ in 0..60 { if indeedhub_recovery_dependencies_running().await { + if !automatic_recovery_intent_allowed(data_dir, stack.containers).await { + return; + } repair_stack_network_aliases(stack).await; break; } @@ -869,7 +944,7 @@ async fn start_stopped_containers_for( records.len(), user_stopped.len() ); - recover_containers(&records).await + recover_containers(data_dir, &records).await } fn should_auto_start_stopped_container(name: &str, include_stack_members: bool) -> bool { @@ -1139,7 +1214,18 @@ async fn podman_status(args: &[&str], timeout: Duration) -> Option Output>>> = std::cell::RefCell::new(None); +} + async fn podman_output(args: &[&str], timeout: Duration) -> Result { + #[cfg(test)] + if let Some(output) = + RECOVERY_PODMAN_TEST.with(|hook| hook.borrow_mut().as_mut().map(|f| f(args))) + { + return Ok(output); + } let mut cmd = tokio::process::Command::new("podman"); cmd.args(args); command_with_timeout(cmd, timeout, &format!("podman {}", args.join(" "))).await @@ -1227,6 +1313,193 @@ mod tests { use super::*; use tempfile::TempDir; + struct RecoveryPodmanFixture; + impl Drop for RecoveryPodmanFixture { + fn drop(&mut self) { + RECOVERY_PODMAN_TEST.with(|hook| *hook.borrow_mut() = None); + } + } + fn recovery_podman_fixture( + root: &Path, + ) -> ( + RecoveryPodmanFixture, + std::rc::Rc>>>, + ) { + use std::os::unix::process::ExitStatusExt; + let calls = std::rc::Rc::new(std::cell::RefCell::new(Vec::new())); + let captured = calls.clone(); + let data = root.to_path_buf(); + let mut started = false; + RECOVERY_PODMAN_TEST.with(|hook| { + *hook.borrow_mut() = Some(Box::new(move |args| { + captured + .borrow_mut() + .push(args.iter().map(|v| v.to_string()).collect()); + let mut status = 0; + let mut stdout = String::new(); + match args.first().copied() { + Some("inspect") if args.get(1).is_some_and(|v| v.starts_with("immich_")) => { + if args.last() == Some(&"{{json .NetworkSettings.Networks}}") { + stdout = "{}".into(); + } else { + stdout = if args[1] == "immich_redis" && !started { + "exited" + } else { + "running" + } + .into(); + } + } + Some("inspect") => status = 1, + Some("ps") => stdout = "immich_redis\n".into(), + Some("start") | Some("network") => { + assert!( + crate::container::update_transaction::Guard::acquire(&data).is_err(), + "mutation escaped lifecycle guard" + ); + if args[0] == "start" { + started = true; + } + } + _ => panic!("unexpected recovery command: {args:?}"), + } + Output { + status: std::process::ExitStatus::from_raw(status << 8), + stdout: stdout.into_bytes(), + stderr: vec![], + } + })) + }); + (RecoveryPodmanFixture, calls) + } + + #[tokio::test] + async fn stack_recovery_never_mutates_saved_held_or_operator_disabled_siblings() { + for reason in [ + "saved", + "broken-saved", + "held", + "broken-hold", + "stopped", + "uninstalled", + "broken-intent", + "locked", + ] { + let root = TempDir::new().unwrap(); + let transaction = root.path().join("update-transactions"); + let mut lock = None; + match reason { + "saved" | "broken-saved" => { + let dir = transaction.join("installed-units"); + std::fs::create_dir_all(&dir).unwrap(); + let body = if reason == "saved" { + serde_json::json!({"schema":1,"operation":uuid::Uuid::new_v4().to_string(),"name":"immich_server","body":"[Container]\nImage=retained\n","mode":384}).to_string() + } else { + "broken".into() + }; + std::fs::write(dir.join("immich_server.json"), body).unwrap(); + } + "held" | "broken-hold" => { + let dir = transaction.join("holds"); + std::fs::create_dir_all(&dir).unwrap(); + std::fs::write( + dir.join("immich_server"), + if reason == "held" { + uuid::Uuid::new_v4().to_string() + } else { + "broken".into() + }, + ) + .unwrap(); + } + "locked" => { + lock = Some( + crate::container::update_transaction::Guard::acquire(root.path()).unwrap(), + ) + } + "stopped" => { + std::fs::write(root.path().join(USER_STOPPED_FILE), r#"["immich_server"]"#) + .unwrap() + } + "uninstalled" => std::fs::write( + root.path().join(USER_UNINSTALLED_FILE), + r#"["immich_server"]"#, + ) + .unwrap(), + _ => std::fs::write(root.path().join(USER_UNINSTALLED_FILE), "broken").unwrap(), + } + let (_fixture, calls) = recovery_podman_fixture(root.path()); + let report = start_stopped_app_stacks(root.path()).await; + assert_eq!(report.recovered, 0, "{reason}"); + assert!( + !calls + .borrow() + .iter() + .any(|args| matches!(args[0].as_str(), "network" | "start")), + "{reason}" + ); + let snapshot = recover_containers( + root.path(), + &[RunningContainerRecord { + name: "immich_redis".into(), + image: String::new(), + }], + ) + .await; + assert_eq!( + snapshot.recovered, 0, + "snapshot sibling admission: {reason}" + ); + assert!( + !calls.borrow().iter().any(|args| args[0] == "start"), + "snapshot {reason}" + ); + drop(lock); + } + } + + #[tokio::test] + async fn stack_and_snapshot_recovery_hold_lifecycle_guard_for_real_mutation_path() { + let root = TempDir::new().unwrap(); + let (_fixture, calls) = recovery_podman_fixture(root.path()); + let report = start_stopped_app_stacks(root.path()).await; + assert_eq!(report.recovered, 1); + assert!(calls + .borrow() + .iter() + .any(|args| args == &["start", "immich_redis"])); + assert!(crate::container::update_transaction::Guard::acquire(root.path()).is_ok()); + let (_positive_snapshot_fixture, positive_calls) = recovery_podman_fixture(root.path()); + let positive = recover_containers( + root.path(), + &[RunningContainerRecord { + name: "immich_redis".into(), + image: String::new(), + }], + ) + .await; + assert_eq!(positive.recovered, 1); + assert!(positive_calls + .borrow() + .iter() + .any(|args| args == &["start", "immich_redis"])); + let (_snapshot_fixture, snapshot_calls) = recovery_podman_fixture(root.path()); + std::fs::write(root.path().join(USER_STOPPED_FILE), r#"["immich_redis"]"#).unwrap(); + let report = recover_containers( + root.path(), + &[RunningContainerRecord { + name: "immich_redis".into(), + image: String::new(), + }], + ) + .await; + assert_eq!(report.recovered, 0); + assert!(!snapshot_calls + .borrow() + .iter() + .any(|args| args[0] == "start")); + } + #[tokio::test] async fn if_recorded_distinguishes_no_record_from_empty_record() { let tmp = TempDir::new().unwrap(); diff --git a/core/archipelago/src/health_monitor.rs b/core/archipelago/src/health_monitor.rs index 0423d319..76c9d981 100644 --- a/core/archipelago/src/health_monitor.rs +++ b/core/archipelago/src/health_monitor.rs @@ -728,7 +728,7 @@ fn parse_health_from_status(status: &str) -> Option { /// Try to recover a container. Running containers need a real restart so /// rootless network helpers such as pasta are recreated; `podman start` is a /// no-op for a running container with a missing host listener. -fn automatic_recovery_allowed(data_dir: &Path, name: &str) -> bool { +pub(crate) fn automatic_recovery_allowed(data_dir: &Path, name: &str) -> bool { match ( crate::container::supervised_update::installed_unit(data_dir, name), crate::container::update_transaction::is_held(data_dir, name), @@ -739,6 +739,11 @@ fn automatic_recovery_allowed(data_dir: &Path, name: &str) -> bool { } async fn restart_container(name: &str, state: &str, data_dir: &Path) -> bool { + // Keep admission held across the command, so an update cannot acquire + // ownership after the policy check and race this automatic restart. + let Ok(_admission) = crate::container::update_transaction::Guard::acquire(data_dir) else { + return false; + }; if !automatic_recovery_allowed(data_dir, name) { warn!(container = %name, "Automatic restart refused: reviewed managed runtime needs explicit recovery"); return false; @@ -1229,6 +1234,13 @@ pub fn spawn_health_monitor(state: Arc, data_dir: PathBuf) { mod tests { use super::*; + #[tokio::test] + async fn automatic_restart_refuses_competing_lifecycle_before_command() { + let root = tempfile::TempDir::new().unwrap(); + let _guard = crate::container::update_transaction::Guard::acquire(root.path()).unwrap(); + assert!(!restart_container("indeedhub-ffmpeg", "running", root.path()).await); + } + #[tokio::test] async fn automatic_recovery_never_restarts_saved_held_or_damaged_managed_runtime() { let root = tempfile::tempdir().unwrap(); diff --git a/core/archipelago/src/main.rs b/core/archipelago/src/main.rs index 8067eb5d..2ae869b5 100644 --- a/core/archipelago/src/main.rs +++ b/core/archipelago/src/main.rs @@ -284,7 +284,7 @@ async fn main() -> Result<()> { "🔧 Recovering {} containers from previous crash...", containers.len() ); - let report = crash_recovery::recover_containers(&containers).await; + let report = crash_recovery::recover_containers(&config.data_dir, &containers).await; info!( "🔧 Recovery complete: {}/{} containers restarted (failed: {:?})", report.recovered, report.total, report.failed diff --git a/docs/managed-update-recovery-implementation.md b/docs/managed-update-recovery-implementation.md index 94a04192..dd2f7a29 100644 --- a/docs/managed-update-recovery-implementation.md +++ b/docs/managed-update-recovery-implementation.md @@ -586,3 +586,42 @@ executed exactly once. This source change still requires matching embedded-helpe build and actual transaction acceptance; it does not close cutover or rollback. The current guest was QMP-paused without reboot to serialize worker runtime and backend process-fixture qualification. + +### Actual automatic-recovery race found and contained (2026-10-08) + +The 180-second helper matching fixture executable built with stable inputs: +full SHA256 `9acf7970f1e2409733e22789a6264b347b3748ba45cee4b79c9e51f6c662acd9`, +stripped VM `41a5fca8551e1d213c0d69ec24aa60534e4610282fbe3433beaa507933108761`, +helper `6fc3f978cb88dbf022dc5bc07eaf0337c6b6b79ff42b20cee7b33ed7c100b879`. +Actual manager startup and read-only API/Redis/PostgreSQL readiness passed with +all seven identities retained. A first request correctly refused stale reviewed +original-unit hashes after prior recovery, without creating a transaction. The +old synthetic plan was archived; independently verified replacement plan +`2ae10528245c5304d8107236ec91f1dda7609b85affab5aea385baf2f87cfbca` retained the +intentional API post-install exit77 hook and exact current original-unit pins. + +Operation `23550e9e-cc0a-427b-9813-7aa1bfc0df5d` stopped the frontend cleanly, +then refused the legacy worker's unknown process-exit result before backup or +target startup. Podman could not observe PID death after SIGKILL; the 45-second +systemd stop budget killed conmon, producing died exit code -1, not a proven +137/143 termination. The stop safety check correctly refused this evidence. + +Separately, actual manager logs prove the periodic `crash_recovery` stack path +restarted that same stopped original worker while transaction holds were active. +At 07:24:17 it logged Recovering stack container, and at 07:24:18 started original +ID `981d58a1f5613ccbf92d91c1dd404f247718ba7146448db9cdfc5b0acffd8527`. +Supervised recovery then refused the unexpected live replacement. This is a +confirmed competing recovery path, not merely a timeout inference. + +The fixture manager was stopped and barred, and the guest QMP-paused on its same +boot. **Operation23550e9e remains unresolved Restoring/Recovering with holds and +journal intact.** Never rewrite it to Restored or adopt the restarted identity as +proof of successful rollback. No Yaya service or app data was changed. + +Source correction applies the existing saved-unit/hold fail-closed policy to +whole stacks, holds the lifecycle lock across every alias/start mutation, reads +stopped/uninstalled intent strictly and rechecks it before later mutations. +Snapshot recovery uses the same stack-wide admission. Actual-path fake Podman +regressions exercise managed/held/corrupt/operator-disabled siblings and ensure +accepted mutation runs under the lifecycle lock. Compilation, isolated execution, +matching executable and fresh actual transaction acceptance remain required.