//! Quadlet recovery preserves exact original launch configuration and image, not //! ephemeral --rm container IDs. Persistent application data is never rolled back. use super::update_transaction::{Guard, Observed, Target}; use anyhow::{Context, Result}; use serde::{Deserialize, Serialize}; use std::{ future::Future, io::Write, os::unix::fs::{DirBuilderExt, OpenOptionsExt}, path::{Path, PathBuf}, }; #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] pub(crate) struct Unit { pub name: String, pub body: String, pub image: String, pub container_id: String, pub file_mode: u32, pub running: bool, pub config_sha256: String, } #[derive(Clone, Debug, Serialize, Deserialize)] pub(crate) struct PreparedTarget { pub body: String, pub manifest: archipelago_container::AppManifest, } #[derive(Clone, Debug, Serialize, Deserialize)] pub(crate) struct RecoveryImage { pub image: String, pub source_container_id: String, pub operation_id: String, } #[derive(Clone, Debug, Serialize, Deserialize)] struct Member { original: Unit, target: Target, target_body: String, target_manifest: archipelago_container::AppManifest, pinned_original_body: String, original_tag: String, recovery_image: Option, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] enum Phase { Prepared, Aborted, Editing, Starting, Restoring, Committed, Restored, } #[derive(Clone, Debug, Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct Journal { schema: u8, id: String, package: String, phase: Phase, members: Vec, #[serde(default)] cleanup_done: bool, #[serde(default)] target_startup_began: bool, } #[derive(Clone, Copy, Debug, serde::Serialize)] #[serde(rename_all = "snake_case")] pub(crate) enum Completion { Committed, Restored, Aborted, } pub(crate) trait Supervisor: Sync { /// Called only after every original writable recovery image is durable. /// Legacy acquisition may gracefully stop AutoRemove writers, so its /// destructive obligation is journaled before entering the controller. /// It must fence/drain writers and finish coherent mounted-data backup; /// image snapshots alone never establish application consistency. /// Acquisition/release are idempotent and operation-owned across restart. fn begin_barrier( &self, operation: &str, originals: &[Unit], recovery: bool, ) -> impl Future> + Send; fn verify_barrier(&self, operation: &str) -> impl Future> + Send; fn release_barrier( &self, operation: &str, outcome: Completion, ) -> impl Future> + Send; /// Only an internal reviewed signed-manifest planner may supply this value; /// browser parameters must never become a unit body or hook recipe. fn prepare_target( &self, target: &Target, original: &Unit, ) -> impl Future> + Send; fn target_hooks( &self, name: &str, manifest: &archipelago_container::AppManifest, ) -> impl Future> + Send; /// Commit the exact original writable layer to a local-only operation-owned /// image, explicitly pausing and excluding mounted volumes. A retry must /// recover its matching image rather than overwrite an unrelated tag. /// This does not establish application-level write quiescence or DB backup. fn snapshot( &self, original: &Unit, operation_id: &str, tag: &str, ) -> impl Future> + Send; fn capture(&self, name: &str) -> impl Future> + Send; fn validate_original_file(&self, original: &Unit) -> impl Future> + Send; fn read(&self, name: &str) -> impl Future> + Send; fn write( &self, original: &Unit, expected: &[String], body: &str, ) -> impl Future> + Send; fn pin(&self, image: &str, tag: &str) -> impl Future> + Send; fn stop(&self, name: &str) -> impl Future> + Send; fn reload(&self) -> impl Future> + Send; fn start(&self, name: &str) -> impl Future> + Send; fn observed(&self, name: &str) -> impl Future>> + Send; fn healthy(&self, name: &str) -> impl Future> + Send; } /// Local recovery images are never pushed or exported. A matching tag may be /// reused after a lost commit response only when image labels bind the exact /// original container and transaction. Mounted data is deliberately excluded. pub(crate) async fn capture_local_recovery_image( original: &Unit, operation: &str, tag: &str, ) -> Result { use super::update_transaction::{Podman, Runtime}; anyhow::ensure!( digest(&original.container_id) && uuid::Uuid::parse_str(operation)?.to_string() == operation && tag.starts_with(&format!("localhost/archy-update-recovery:{operation}-")) && tag .rsplit('-') .next() .is_some_and(|v| v.parse::().is_ok()), "Invalid local recovery image ownership" ); async fn command(args: &[&str]) -> Result { tokio::time::timeout( std::time::Duration::from_secs(300), tokio::process::Command::new("podman") .args(args) .kill_on_drop(true) .output(), ) .await .context("Recovery image operation timed out; original runtime retained")? .context("Recovery image runtime unavailable") } async fn inspect(tag: &str, original: &Unit, operation: &str) -> Result { let result = command(&["image", "inspect", tag]).await?; anyhow::ensure!( result.status.success(), "Local recovery image inspection failed" ); let rows: Vec = serde_json::from_slice(&result.stdout)?; anyhow::ensure!(rows.len() == 1, "Ambiguous recovery image"); let row = &rows[0]; let labels = row .pointer("/Config/Labels") .context("Recovery image has no ownership labels")?; anyhow::ensure!( labels .get("io.archipelago.recovery.operation") .and_then(|v| v.as_str()) == Some(operation) && labels .get("io.archipelago.recovery.container") .and_then(|v| v.as_str()) == Some(original.container_id.as_str()), "Existing recovery image belongs to another operation" ); let image = row .get("Id") .or_else(|| row.get("ID")) .and_then(|v| v.as_str()) .context("Recovery image identity missing")? .trim_start_matches("sha256:") .to_string(); anyhow::ensure!(digest(&image), "Invalid recovery image digest"); Ok(RecoveryImage { image, source_container_id: original.container_id.clone(), operation_id: operation.into(), }) } let current = Podman .inspect(&original.name) .await? .context("Original container missing before snapshot")?; anyhow::ensure!( current.id == original.container_id && current.image == original.image && current.running == original.running && current.config_sha256 == original.config_sha256, "Original runtime changed before writable-layer snapshot" ); let exists = command(&["image", "exists", tag]).await?; match exists.status.code() { Some(0) => return inspect(tag, original, operation).await, Some(1) => {} _ => anyhow::bail!("Recovery image inventory unavailable"), } let operation_label = format!("LABEL io.archipelago.recovery.operation={operation}"); let container_label = format!( "LABEL io.archipelago.recovery.container={}", original.container_id ); let result = command(&[ "commit", "--pause=true", "--include-volumes=false", "--change", &operation_label, "--change", &container_label, &original.container_id, tag, ]) .await?; anyhow::ensure!( result.status.success(), "Writable-layer snapshot failed; original runtime retained" ); inspect(tag, original, operation).await } fn simple(value: &str) -> bool { !value.is_empty() && value.len() <= 128 && value .bytes() .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_')) } fn digest(value: &str) -> bool { value.len() == 64 && value.bytes().all(|b| b.is_ascii_hexdigit()) } /// Parse only the renderer-owned unit shape. Ambiguous images, includes and /// external environment files cannot become a guessed recovery recipe. pub(crate) fn pin_body(body: &str, name: &str, image: &str) -> Result { anyhow::ensure!( body.len() <= 1024 * 1024 && simple(name), "Invalid original unit" ); let mut section = ""; let mut images = 0; let mut names = 0; let mut output = String::new(); for line in body.lines() { let trimmed = line.trim(); anyhow::ensure!( !trimmed.ends_with('\\') && !trimmed.starts_with(".include") && !trimmed.starts_with("EnvironmentFile=") && !trimmed.starts_with("EnvFile="), "Unit has external/continued configuration; exact recovery is not supported yet" ); if trimmed.starts_with('[') { section = trimmed; } if section == "[Container]" && trimmed.starts_with("Image=") { images += 1; output.push_str(&format!("Image={image}\n")); } else { if section == "[Container]" && trimmed.starts_with("ContainerName=") { names += 1; anyhow::ensure!( trimmed == format!("ContainerName={name}"), "Unit container ownership mismatch" ); } output.push_str(line); output.push('\n'); } } anyhow::ensure!( images == 1 && names == 1, "Original Quadlet must bind one container and image" ); Ok(output) } /// Three-way configuration migration: the previous renderer output must come /// from the installed immutable manifest, never today's mutable catalog. Keep /// unrelated operator overrides; a conflicting required change is a preflight /// error rather than an overwrite. Values are never included in errors. pub(crate) fn merge_reviewed_unit(original: &str, previous: &str, next: &str) -> Result { type Key = (String, String); fn parse(body: &str) -> Result<(Vec, std::collections::BTreeMap>)> { anyhow::ensure!(body.len() <= 1024 * 1024, "Unit exceeds migration limit"); let mut section = String::new(); let mut sections = std::collections::HashSet::new(); let mut order = Vec::new(); let mut values = std::collections::BTreeMap::>::new(); for raw in body.lines() { let line = raw.trim(); if line.is_empty() || line.starts_with('#') || line.starts_with(';') { continue; } anyhow::ensure!( !line.ends_with('\\') && !line.starts_with(".include"), "Continued or included unit cannot be migrated automatically" ); if line.starts_with('[') { anyhow::ensure!( line.ends_with(']') && sections.insert(line.to_string()), "Repeated or malformed unit section" ); section = line.into(); continue; } let (directive, value) = line.split_once('=').context("Malformed unit directive")?; anyhow::ensure!( !section.is_empty() && !value.is_empty() && directive.bytes().all(|b| b.is_ascii_alphanumeric()) && !matches!(directive, "EnvironmentFile" | "EnvFile"), "Unsupported unit reset or external configuration" ); let key = if directive == "Environment" { let env = value.strip_prefix('"').unwrap_or(value); let (name, _) = env .split_once('=') .context("Unsupported environment directive")?; anyhow::ensure!( !name.is_empty() && name.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'_') && (value.starts_with('"') && value.ends_with('"') || !value.contains(char::is_whitespace)), "Ambiguous environment directive" ); format!("Environment:{name}") } else { directive.into() }; let key = (section.clone(), key); if !values.contains_key(&key) { order.push(key.clone()); } values.entry(key).or_default().push(raw.to_string()); } Ok((order, values)) } let (mut order, original) = parse(original)?; let (_, previous) = parse(previous)?; let (next_order, next) = parse(next)?; let mut merged = original.clone(); let keys: std::collections::BTreeSet<_> = previous.keys().chain(next.keys()).cloned().collect(); for key in keys { let old = previous.get(&key); let wanted = next.get(&key); if old == wanted { continue; } let actual = original.get(&key); anyhow::ensure!( actual == old || actual == wanted, "Required unit configuration conflicts with an operator override in {} {}", key.0, key.1 ); match wanted { Some(lines) => { merged.insert(key, lines.clone()); } None => { merged.remove(&key); } } } for key in next_order { if !order.contains(&key) { order.push(key); } } // Group sections once; preserve directive and repeated-value order within // each section, including operator-only directives absent from both plans. let mut sections = Vec::::new(); for (section, _) in &order { if !sections.contains(section) { sections.push(section.clone()); } } let mut output = String::new(); for section in sections { output.push_str(§ion); output.push('\n'); for key in order.iter().filter(|key| key.0 == section) { if let Some(lines) = merged.get(key) { for line in lines { output.push_str(line); output.push('\n'); } } } } Ok(output) } fn root(guard: &Guard) -> Result { let path = guard.directory().join("supervised"); match std::fs::DirBuilder::new().mode(0o700).create(&path) { Ok(()) => {} Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {} Err(e) => return Err(e.into()), } anyhow::ensure!( !std::fs::symlink_metadata(&path)?.file_type().is_symlink(), "Invalid supervised journal directory" ); Ok(path) } fn validate(record: &Journal) -> Result<()> { anyhow::ensure!( record.schema == 1 && simple(&record.package) && uuid::Uuid::parse_str(&record.id)?.to_string() == record.id && !record.members.is_empty() && record.members.len() <= 32, "Invalid supervised update journal" ); let mut names = std::collections::HashSet::new(); for (index, member) in record.members.iter().enumerate() { anyhow::ensure!( simple(&member.original.name) && names.insert(&member.original.name) && member.target.name == member.original.name && digest(&member.target.image) && digest(&member.original.image) && digest(&member.original.container_id) && member.original.file_mode & !0o777 == 0 && member.original.file_mode & 0o022 == 0 && digest(&member.original.config_sha256), "Invalid original supervised identity" ); anyhow::ensure!( member.original_tag == format!("localhost/archy-update-recovery:{}-{index}", record.id), "Recovery image pin changed" ); if let Some(image) = &member.recovery_image { anyhow::ensure!( digest(&image.image) && image.source_container_id == member.original.container_id && image.operation_id == record.id, "Recovery image ownership changed" ); } anyhow::ensure!( matches!(record.phase, Phase::Prepared | Phase::Aborted) || member.recovery_image.is_some(), "Destructive update lacks a durable writable-layer recovery image" ); let restore_image = member .recovery_image .as_ref() .map(|value| value.image.as_str()) .unwrap_or(&member.original.image); anyhow::ensure!( member.pinned_original_body == pin_body( &member.original.body, &member.original.name, &format!("sha256:{restore_image}") )? && member.target_body == pin_body( &member.target_body, &member.original.name, &member.target.reference )? && member.target_manifest.app.container.image.as_deref() == Some(member.target.reference.as_str()), "Saved unit recipe changed" ); } Ok(()) } fn save(guard: &Guard, record: &Journal) -> Result<()> { validate(record)?; let dir = root(guard)?; let bytes = serde_json::to_vec(record)?; anyhow::ensure!( bytes.len() <= 4 * 1024 * 1024, "Supervised recovery journal too large" ); let temporary = dir.join(format!(".{}.tmp", uuid::Uuid::new_v4())); let result = (|| -> Result<()> { let mut file = std::fs::OpenOptions::new() .create_new(true) .write(true) .mode(0o600) .open(&temporary)?; file.write_all(&bytes)?; file.sync_all()?; std::fs::rename(&temporary, dir.join(format!("{}.json", record.id)))?; std::fs::File::open(&dir)?.sync_all()?; std::fs::File::open(guard.directory())?.sync_all()?; Ok(()) })(); if result.is_err() { let _ = std::fs::remove_file(temporary); } result } // The committed unit is installation evidence, not a cache of today's catalog. // Keep it until explicit uninstall or the next reviewed managed transaction. #[derive(Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct InstalledUnit { schema: u8, operation: String, name: String, body: String, mode: u32, } fn installed_path(data: &Path, name: &str) -> Result { anyhow::ensure!(simple(name), "Invalid managed member name"); Ok(data .join("update-transactions/installed-units") .join(format!("{name}.json"))) } fn publish_installed(guard: &Guard, record: &Journal) -> Result<()> { anyhow::ensure!( record.phase == Phase::Committed, "Only committed units may be published" ); let dir = guard.directory().join("installed-units"); match std::fs::DirBuilder::new().mode(0o700).create(&dir) { Ok(()) => {} Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {} Err(e) => return Err(e.into()), } anyhow::ensure!( std::fs::symlink_metadata(&dir)?.is_dir(), "Invalid installed unit directory" ); for member in &record.members { let saved = InstalledUnit { schema: 1, operation: record.id.clone(), name: member.original.name.clone(), body: member.target_body.clone(), mode: member.original.file_mode, }; let temporary = dir.join(format!(".{}.tmp", uuid::Uuid::new_v4())); let mut file = std::fs::OpenOptions::new() .write(true) .create_new(true) .mode(0o600) .open(&temporary)?; file.write_all(&serde_json::to_vec(&saved)?)?; file.sync_all()?; std::fs::rename(&temporary, dir.join(format!("{}.json", saved.name)))?; } std::fs::File::open(&dir)?.sync_all()?; std::fs::File::open(guard.directory())?.sync_all()?; Ok(()) } pub(crate) fn installed_unit(data: &Path, name: &str) -> Result> { let path = installed_path(data, name)?; let meta = match std::fs::symlink_metadata(&path) { Ok(meta) => meta, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None), Err(e) => return Err(e.into()), }; anyhow::ensure!( meta.is_file() && meta.len() <= 4 * 1024 * 1024, "Invalid installed managed recipe" ); let saved: InstalledUnit = serde_json::from_slice(&std::fs::read(path)?)?; anyhow::ensure!( saved.schema == 1 && saved.name == name && uuid::Uuid::parse_str(&saved.operation)?.to_string() == saved.operation && saved.mode & !0o777 == 0 && saved.mode & 0o022 == 0, "Invalid installed managed recipe binding" ); Ok(Some((saved.body, saved.mode))) } pub(crate) fn forget_installed(data: &Path, name: &str) -> Result<()> { let path = installed_path(data, name)?; match std::fs::remove_file(&path) { Ok(()) => std::fs::File::open(path.parent().unwrap())?.sync_all()?, Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} Err(e) => return Err(e.into()), }; Ok(()) } fn records(guard: &Guard) -> Result> { let dir = root(guard)?; let mut records = Vec::new(); for entry in std::fs::read_dir(dir)? { let entry = entry?; if entry.path().extension().and_then(|v| v.to_str()) != Some("json") { continue; } anyhow::ensure!( records.len() < 128 && entry.file_type()?.is_file() && entry.metadata()?.len() <= 4 * 1024 * 1024, "Invalid supervised recovery inventory" ); let record: Journal = serde_json::from_slice(&std::fs::read(entry.path())?)?; validate(&record)?; anyhow::ensure!( entry.file_name() == format!("{}.json", record.id).as_str(), "Supervised journal name changed" ); records.push(record); } Ok(records) } pub(crate) fn require_clear(guard: &Guard) -> Result<()> { anyhow::ensure!( records(guard)?.iter().all(|r| r.cleanup_done && matches!(r.phase, Phase::Committed | Phase::Restored | Phase::Aborted)), "A supervised update needs recovery first" ); Ok(()) } pub(crate) async fn execute( guard: &Guard, package: &str, targets: &[Target], supervisor: &impl Supervisor, ) -> Result<()> { guard.require_clear()?; anyhow::ensure!( !targets.is_empty() && targets.len() <= 32, "Invalid supervised stack" ); let id = uuid::Uuid::new_v4().to_string(); let mut members = Vec::new(); for (index, target) in targets.iter().enumerate() { let original = supervisor.capture(&target.name).await?; // Inactive units require a durable explicit-start staging path; never // implement this by starting and then stopping a user's stopped member. anyhow::ensure!(original.running,"Stopped supervised member requires staged update; all original services remain unchanged"); let prepared = supervisor.prepare_target(target, &original).await?; anyhow::ensure!( prepared.manifest.app.container.image.as_deref() == Some(target.reference.as_str()), "Reviewed target manifest image changed" ); let target_body = pin_body(&prepared.body, &original.name, &target.reference)?; anyhow::ensure!( target_body == prepared.body, "Reviewed target unit image changed" ); let pinned_original_body = pin_body( &original.body, &original.name, &format!("sha256:{}", original.image), )?; members.push(Member { original, target: target.clone(), target_body, target_manifest: prepared.manifest, pinned_original_body, original_tag: format!("localhost/archy-update-recovery:{id}-{index}"), recovery_image: None, }); } let mut record = Journal { schema: 1, id, package: package.into(), phase: Phase::Prepared, members, cleanup_done: false, target_startup_began: false, }; save(guard, &record)?; for member in &record.members { guard.hold(&member.original.name, &record.id)?; } let result = apply(guard, &mut record, supervisor).await; if let Err(error) = result { if record.phase == Phase::Committed { return Err(error .context("Update committed; completion cleanup must resume without rolling back")); } return match restore(guard,&mut record,supervisor).await { Ok(())=>Err(error.context("Original supervised image/configuration and running intent restored; container IDs may change and data was not rolled back")), Err(recovery)=>Err(error.context(format!("Supervised runtime recovery remains unresolved: {recovery:#}"))), }; } Ok(()) } async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> { for index in 0..record.members.len() { let member = &record.members[index]; supervisor.validate_original_file(&member.original).await?; anyhow::ensure!( supervisor.read(&member.original.name).await? == member.original.body, "Unit edited before update; originals retained" ); let image = supervisor .snapshot(&member.original, &record.id, &member.original_tag) .await?; anyhow::ensure!( digest(&image.image) && image.source_container_id == member.original.container_id && image.operation_id == record.id, "Writable-layer recovery ownership mismatch" ); let member = &mut record.members[index]; member.pinned_original_body = pin_body( &member.original.body, &member.original.name, &format!("sha256:{}", image.image), )?; member.recovery_image = Some(image); // The image acknowledgement becomes durable before any original stop. save(guard, record)?; } record.phase = Phase::Editing; save(guard, record)?; // The legacy controller may stop original writers. Every writable layer // is already recoverable and this obligation survives cancellation. let originals: Vec<_> = record .members .iter() .map(|member| member.original.clone()) .collect(); supervisor .begin_barrier(&record.id, &originals, false) .await?; supervisor.verify_barrier(&record.id).await?; for member in record.members.iter().rev() { supervisor.verify_barrier(&record.id).await?; supervisor.stop(&member.original.name).await?; } for member in &record.members { supervisor .write( &member.original, &[member.original.body.clone()], &member.target_body, ) .await?; } supervisor.reload().await?; record.phase = Phase::Starting; record.target_startup_began = true; record.cleanup_done = false; save(guard, record)?; for member in &record.members { anyhow::ensure!( supervisor.read(&member.original.name).await? == member.target_body, "Reviewed unit changed before target start; recovery required" ); supervisor.start(&member.original.name).await?; supervisor .target_hooks(&member.original.name, &member.target_manifest) .await?; let observed = supervisor .observed(&member.original.name) .await? .context("Updated supervised member missing")?; anyhow::ensure!( observed.running && observed.image == member.target.image && supervisor.healthy(&member.original.name).await?, "Updated supervised member failed verification" ); } record.phase = Phase::Committed; save(guard, record)?; publish_installed(guard, record)?; supervisor .release_barrier(&record.id, Completion::Committed) .await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } record.cleanup_done = true; save(guard, record) } async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> { if record.phase == Phase::Prepared { // Only image snapshots may have happened. Never stop/recreate an intact // app just because preflight or snapshotting failed on another member. for member in &record.members { let original = supervisor .observed(&member.original.name) .await? .context("Original service disappeared during preparation")?; anyhow::ensure!( supervisor.read(&member.original.name).await? == member.original.body && original.id == member.original.container_id && original.image == member.original.image && original.running == member.original.running && original.config_sha256 == member.original.config_sha256, "Original service changed during preparation; recovery requires inspection" ); } record.phase = Phase::Aborted; save(guard, record)?; supervisor .release_barrier(&record.id, Completion::Aborted) .await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } record.cleanup_done = true; return save(guard, record); } let originals: Vec<_> = record .members .iter() .map(|member| member.original.clone()) .collect(); record.phase = Phase::Restoring; save(guard, record)?; supervisor .begin_barrier(&record.id, &originals, true) .await?; supervisor.verify_barrier(&record.id).await?; // Refuse to overwrite a foreign edit before stopping any surviving member. for member in &record.members { supervisor.validate_original_file(&member.original).await?; let body = supervisor.read(&member.original.name).await?; anyhow::ensure!( [ &member.original.body, &member.target_body, &member.pinned_original_body ] .contains(&&body), "Foreign unit edit requires explicit recovery" ); supervisor .pin( &member .recovery_image .as_ref() .context("Missing recovery image")? .image, &member.original_tag, ) .await?; } for member in record.members.iter().rev() { supervisor.stop(&member.original.name).await?; } for member in &record.members { supervisor .write( &member.original, &[ member.original.body.clone(), member.target_body.clone(), member.pinned_original_body.clone(), ], &member.pinned_original_body, ) .await?; } supervisor.reload().await?; for member in &record.members { if member.original.running { anyhow::ensure!( supervisor.read(&member.original.name).await? == member.pinned_original_body, "Original recovery unit changed before start" ); supervisor.start(&member.original.name).await?; } let observed = supervisor.observed(&member.original.name).await?; if member.original.running { let current = observed.context("Original supervised service did not return")?; anyhow::ensure!( current.running && current.image == member .recovery_image .as_ref() .context("Missing recovery image")? .image && current.config_sha256 == member.original.config_sha256, "Original launch configuration did not recover" ); } else { anyhow::ensure!( observed.is_none_or(|v| !v.running), "Originally stopped service unexpectedly running" ); } } for member in &record.members { guard.hold(&member.original.name, &record.id)?; } record.phase = Phase::Restored; save(guard, record)?; supervisor .release_barrier(&record.id, Completion::Restored) .await?; record.cleanup_done = true; save(guard, record) } pub(crate) async fn recover(guard: &Guard, supervisor: &impl Supervisor) -> Result<()> { for mut record in records(guard)? { if record.cleanup_done { continue; } match record.phase { Phase::Committed | Phase::Aborted => { if record.phase == Phase::Committed { publish_installed(guard, &record)?; } let outcome = if record.phase == Phase::Committed { Completion::Committed } else { Completion::Aborted }; supervisor.release_barrier(&record.id, outcome).await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } } Phase::Restored => { supervisor .release_barrier(&record.id, Completion::Restored) .await?; } // Never release a newer owner. _ => restore(guard, &mut record, supervisor).await?, } record.cleanup_done = true; save(guard, &record)?; } Ok(()) } pub(crate) fn needs_recovery(guard: &Guard) -> Result { Ok(records(guard)?.iter().any(|record| !record.cleanup_done)) } #[cfg(test)] mod tests { use super::*; use std::sync::{ atomic::{AtomicBool, AtomicUsize, Ordering}, Mutex, }; struct Mock { original: Unit, body: Mutex, running: AtomicBool, calls: Mutex>, fail_new_hooks: AtomicBool, fail_snapshot: AtomicBool, fail_barrier_release: AtomicBool, barrier: Mutex>, generation: AtomicUsize, } impl Mock { fn new() -> Self { let body=format!("[Container]\nContainerName=movie\nImage=sha256:{}\nEnvironment=OPERATOR_VALUE=retained\nPull=never\n[Service]\nRestart=always\n", "a".repeat(64)); Self { original: Unit { name: "movie".into(), body: body.clone(), image: "a".repeat(64), container_id: format!("{:064x}", 1), file_mode: 0o600, running: true, config_sha256: "c".repeat(64), }, body: Mutex::new(body), running: AtomicBool::new(true), calls: Default::default(), fail_new_hooks: AtomicBool::new(false), fail_snapshot: AtomicBool::new(false), fail_barrier_release: AtomicBool::new(false), barrier: Mutex::new(None), generation: AtomicUsize::new(1), } } fn target() -> Target { Target { name: "movie".into(), reference: format!("localhost/new@sha256:{}", "b".repeat(64)), image: "b".repeat(64), } } } impl Supervisor for Mock { async fn begin_barrier( &self, operation: &str, _originals: &[Unit], _recovery: bool, ) -> Result<()> { let mut held = self.barrier.lock().unwrap(); anyhow::ensure!( held.as_deref().is_none_or(|id| id == operation), "Another drain barrier owns admission" ); *held = Some(operation.into()); Ok(()) } async fn verify_barrier(&self, operation: &str) -> Result<()> { anyhow::ensure!( self.barrier.lock().unwrap().as_deref() == Some(operation), "Write barrier not held" ); Ok(()) } async fn release_barrier(&self, operation: &str, _outcome: Completion) -> Result<()> { anyhow::ensure!( !self.fail_barrier_release.load(Ordering::SeqCst), "Barrier release unavailable" ); let mut held = self.barrier.lock().unwrap(); if held.as_deref() == Some(operation) { *held = None; } Ok(()) } async fn prepare_target( &self, target: &Target, _original: &Unit, ) -> Result { let manifest = archipelago_container::AppManifest::parse(&format!( "app:\n id: movie\n name: Movie\n version: 2.0.0\n container:\n image: {}\n", target.reference ))?; let body = pin_body(&self.original.body, "movie", &target.reference)?.replace( "Pull=never", "Environment=NEW_FEATURE=enabled\nVolume=/identity:/run/identity:ro\nPull=never", ); Ok(PreparedTarget { body, manifest }) } async fn target_hooks( &self, _name: &str, _manifest: &archipelago_container::AppManifest, ) -> Result<()> { self.calls.lock().unwrap().push("new-hooks".into()); anyhow::ensure!( !self.fail_new_hooks.load(Ordering::SeqCst), "New provider hook failed" ); Ok(()) } async fn snapshot( &self, original: &Unit, operation_id: &str, _tag: &str, ) -> Result { self.calls.lock().unwrap().push("snapshot-original".into()); anyhow::ensure!( !self.fail_snapshot.load(Ordering::SeqCst), "Snapshot failed" ); Ok(RecoveryImage { image: "e".repeat(64), source_container_id: original.container_id.clone(), operation_id: operation_id.into(), }) } async fn capture(&self, _name: &str) -> Result { Ok(self.original.clone()) } async fn validate_original_file(&self, _original: &Unit) -> Result<()> { Ok(()) } async fn read(&self, _name: &str) -> Result { Ok(self.body.lock().unwrap().clone()) } async fn write(&self, _original: &Unit, expected: &[String], body: &str) -> Result<()> { let mut current = self.body.lock().unwrap(); anyhow::ensure!(expected.contains(&*current), "Foreign unit edit"); *current = body.into(); self.calls.lock().unwrap().push("write".into()); Ok(()) } async fn pin(&self, _image: &str, _tag: &str) -> Result<()> { self.calls.lock().unwrap().push("pin-original".into()); Ok(()) } async fn stop(&self, _name: &str) -> Result<()> { self.calls.lock().unwrap().push("stop-unit".into()); self.running.store(false, Ordering::SeqCst); Ok(()) } async fn reload(&self) -> Result<()> { self.calls.lock().unwrap().push("reload".into()); Ok(()) } async fn start(&self, _name: &str) -> Result<()> { self.calls.lock().unwrap().push("start-unit".into()); self.running.store(true, Ordering::SeqCst); self.generation.fetch_add(1, Ordering::SeqCst); Ok(()) } async fn observed(&self, _name: &str) -> Result> { if !self.running.load(Ordering::SeqCst) { return Ok(None); } let new = self.body.lock().unwrap().contains("NEW_FEATURE=enabled"); Ok(Some(Observed { id: format!("{:064x}", self.generation.load(Ordering::SeqCst)), name: "movie".into(), image: if new { "b".repeat(64) } else if self.body.lock().unwrap().contains(&"e".repeat(64)) { "e".repeat(64) } else { "a".repeat(64) }, running: true, retainable: false, config_sha256: if new { "d".repeat(64) } else { "c".repeat(64) }, })) } async fn healthy(&self, _name: &str) -> Result { Ok(true) } } #[test] fn reviewed_migration_keeps_unrelated_operator_values_and_applies_required_new_fields() { let old = "[Container]\nImage=old\nEnvironment=FEATURE=old\nEnvironment=PORT=1\n[Service]\nRestart=always\n"; let actual = old .replace("PORT=1", "PORT=42") .replace("[Service]", "Environment=OPERATOR=mine\n[Service]"); let next = old .replace("Image=old", "Image=new") .replace("FEATURE=old", "FEATURE=new") .replace("[Service]", "Volume=/identity:/run/identity:ro\n[Service]"); let merged = merge_reviewed_unit(&actual, old, &next).unwrap(); assert!(merged.contains("Image=new")); assert!(merged.contains("FEATURE=new")); assert!(merged.contains("PORT=42")); assert!(merged.contains("OPERATOR=mine")); assert!(merged.contains("Volume=/identity:/run/identity:ro")); assert_eq!(merge_reviewed_unit(&merged, old, &next).unwrap(), merged); } #[test] fn conflicting_required_environment_or_mount_migration_is_rejected_before_lifecycle() { let old = "[Container]\nEnvironment=FEATURE=old\nVolume=/old:/data\n"; assert!(merge_reviewed_unit( &old.replace("FEATURE=old", "FEATURE=operator"), old, &old.replace("FEATURE=old", "FEATURE=new") ) .is_err()); assert!(merge_reviewed_unit( &old.replace("/old:/data", "/operator:/data"), old, &old.replace("/old:/data", "/required:/data") ) .is_err()); assert!(merge_reviewed_unit("[Container]\nEnvironment=\n", old, old).is_err()); assert!(merge_reviewed_unit("[Container]\nEnvironment=A=1 B=2\n", old, old).is_err()); } #[tokio::test] async fn lost_barrier_release_after_commit_never_rolls_back_a_verified_update() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); runtime.fail_barrier_release.store(true, Ordering::SeqCst); let error = execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap_err(); assert!(error.to_string().contains("Update committed")); assert_eq!(records(&guard).unwrap()[0].phase, Phase::Committed); assert_eq!( runtime.observed("movie").await.unwrap().unwrap().image, "b".repeat(64) ); assert_eq!( runtime .calls .lock() .unwrap() .iter() .filter(|call| *call == "stop-unit") .count(), 1 ); runtime.calls.lock().unwrap().clear(); runtime.fail_barrier_release.store(false, Ordering::SeqCst); recover(&guard, &runtime).await.unwrap(); assert!(runtime.calls.lock().unwrap().is_empty()); assert!(runtime.barrier.lock().unwrap().is_none()); assert!(!super::super::update_transaction::is_held(root.path(), "movie").unwrap()); } #[tokio::test] async fn snapshot_failure_never_stops_or_recreates_original_runtime() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); runtime.fail_snapshot.store(true, Ordering::SeqCst); assert!(execute(&guard, "movie", &[Mock::target()], &runtime) .await .is_err()); assert_eq!(*runtime.calls.lock().unwrap(), ["snapshot-original"]); let original = runtime.observed("movie").await.unwrap().unwrap(); assert_eq!(original.id, runtime.original.container_id); assert_eq!(original.image, runtime.original.image); assert_eq!(*runtime.body.lock().unwrap(), runtime.original.body); assert_eq!(records(&guard).unwrap()[0].phase, Phase::Aborted); assert!(!super::super::update_transaction::is_held(root.path(), "movie").unwrap()); recover(&guard, &runtime).await.unwrap(); assert_eq!(*runtime.calls.lock().unwrap(), ["snapshot-original"]); } #[tokio::test] async fn committed_unit_survives_restart_until_explicit_uninstall() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap(); let saved = installed_unit(root.path(), "movie").unwrap().unwrap(); assert_eq!(saved.0, *runtime.body.lock().unwrap()); assert!(saved.0.contains("OPERATOR_VALUE=retained")); runtime.calls.lock().unwrap().clear(); recover(&guard, &runtime).await.unwrap(); assert_eq!( installed_unit(root.path(), "movie").unwrap().unwrap(), saved ); assert!(runtime.calls.lock().unwrap().is_empty()); forget_installed(root.path(), "movie").unwrap(); recover(&guard, &runtime).await.unwrap(); assert!(installed_unit(root.path(), "movie").unwrap().is_none()); } #[tokio::test] async fn committed_restart_releases_only_its_own_hold_without_runtime_mutation() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap(); let mut record = records(&guard).unwrap().pop().unwrap(); record.cleanup_done = false; save(&guard, &record).unwrap(); guard.hold("movie", &record.id).unwrap(); runtime.calls.lock().unwrap().clear(); recover(&guard, &runtime).await.unwrap(); assert!(runtime.calls.lock().unwrap().is_empty()); assert!(!super::super::update_transaction::is_held(root.path(), "movie").unwrap()); let next = uuid::Uuid::new_v4().to_string(); guard.hold("movie", &next).unwrap(); recover(&guard, &runtime).await.unwrap(); assert!(super::super::update_transaction::is_held(root.path(), "movie").unwrap()); } #[tokio::test] async fn forward_applies_reviewed_new_configuration_and_hooks_not_only_image() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap(); let body = runtime.body.lock().unwrap(); assert!(body.contains("OPERATOR_VALUE=retained")); assert!(body.contains("NEW_FEATURE=enabled")); assert!(body.contains("Volume=/identity:/run/identity:ro")); assert!(runtime.calls.lock().unwrap().contains(&"new-hooks".into())); assert_eq!(records(&guard).unwrap()[0].phase, Phase::Committed); assert!(!super::super::update_transaction::is_held(root.path(), "movie").unwrap()); } #[tokio::test] async fn auto_remove_recovery_restores_old_configuration_without_claiming_original_id() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); runtime.fail_new_hooks.store(true, Ordering::SeqCst); let error = execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap_err(); assert!(error.to_string().contains("container IDs may change")); assert_eq!( *runtime.body.lock().unwrap(), pin_body( &runtime.original.body, "movie", &format!("sha256:{}", "e".repeat(64)) ) .unwrap() ); let restored = runtime.observed("movie").await.unwrap().unwrap(); assert_eq!(restored.image, "e".repeat(64)); assert_eq!(restored.config_sha256, runtime.original.config_sha256); assert_ne!(restored.id, format!("{:064x}", 1)); assert!(super::super::update_transaction::is_held(root.path(), "movie").unwrap()); assert_eq!(records(&guard).unwrap()[0].phase, Phase::Restored); } #[tokio::test] async fn interrupted_restart_uses_saved_old_unit_even_after_catalog_plan_changes() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap(); let mut record = records(&guard).unwrap().pop().unwrap(); record.phase = Phase::Starting; record.cleanup_done = false; *runtime.barrier.lock().unwrap() = Some(record.id.clone()); save(&guard, &record).unwrap(); runtime.calls.lock().unwrap().clear(); recover(&guard, &runtime).await.unwrap(); assert_eq!( *runtime.body.lock().unwrap(), pin_body( &runtime.original.body, "movie", &format!("sha256:{}", "e".repeat(64)) ) .unwrap() ); assert!(!runtime.calls.lock().unwrap().contains(&"new-hooks".into())); } #[tokio::test] async fn foreign_unit_edit_blocks_recovery_before_stopping_other_services() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let runtime = Mock::new(); execute(&guard, "movie", &[Mock::target()], &runtime) .await .unwrap(); let mut record = records(&guard).unwrap().pop().unwrap(); record.phase = Phase::Starting; record.cleanup_done = false; *runtime.barrier.lock().unwrap() = Some(record.id.clone()); save(&guard, &record).unwrap(); *runtime.body.lock().unwrap() = "operator replaced unit".into(); runtime.calls.lock().unwrap().clear(); assert!(recover(&guard, &runtime).await.is_err()); assert!(runtime.calls.lock().unwrap().is_empty()); assert_eq!(*runtime.body.lock().unwrap(), "operator replaced unit"); } #[tokio::test] async fn unsupported_stopped_supervised_member_is_never_started_as_a_workaround() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); let mut runtime = Mock::new(); runtime.original.running = false; runtime.running.store(false, Ordering::SeqCst); assert!(execute(&guard, "movie", &[Mock::target()], &runtime) .await .is_err()); assert!(runtime.calls.lock().unwrap().is_empty()); assert!(records(&guard).unwrap().is_empty()); } }