Keep automatic recovery outside managed update ownership

This commit is contained in:
archipelago
2026-10-08 05:20:40 -04:00
parent 8ae0bdea6a
commit 6a342f665c
4 changed files with 338 additions and 14 deletions
+285 -12
View File
@@ -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<String> = 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<crate::container::update_transaction::Guard> {
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<std::process:
.map(|output| output.status)
}
#[cfg(test)]
thread_local! {
static RECOVERY_PODMAN_TEST: std::cell::RefCell<Option<Box<dyn FnMut(&[&str]) -> Output>>> = std::cell::RefCell::new(None);
}
async fn podman_output(args: &[&str], timeout: Duration) -> Result<Output> {
#[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<std::cell::RefCell<Vec<Vec<String>>>>,
) {
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();
+13 -1
View File
@@ -728,7 +728,7 @@ fn parse_health_from_status(status: &str) -> Option<String> {
/// 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<StateManager>, 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();
+1 -1
View File
@@ -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
@@ -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.