diff --git a/core/archipelago/src/api/rpc/package/async_lifecycle.rs b/core/archipelago/src/api/rpc/package/async_lifecycle.rs index 63ad037b..fe67b145 100644 --- a/core/archipelago/src/api/rpc/package/async_lifecycle.rs +++ b/core/archipelago/src/api/rpc/package/async_lifecycle.rs @@ -377,12 +377,34 @@ impl RpcHandler { Err(e) => { error!("package.update {} failed: {:#}", package_id_spawn, e); install_log(&format!("UPDATE FAIL: {} — {:#}", package_id_spawn, e)).await; - // Inner handler already ran rollback_update + cleared - // update state, but be defensive: revert to pre-state - // in case the inner flow died before its cleanup. - if let Some(prev) = pre_state { - set_package_state(&handler.state_manager, &package_id_spawn, prev).await; - } + // Release the transitional overlay before asking the scanner + // for real state. Prior Running is not proof of successful + // rollback, and a failed preflight is not proof of Stopped. + handler + .state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(&package_id_spawn) { + finish_failed_update(entry); + } + data.notifications.retain(|item| { + item.id != format!("update-failed-{package_id_spawn}") + }); + data.notifications.push(crate::data_model::Notification { + id: format!("update-failed-{package_id_spawn}"), + level: crate::data_model::NotificationLevel::Error, + title: format!("Could not update {package_id_spawn}"), + message: format!( + "{e}. Runtime recovery does not roll back database changes." + ), + timestamp: chrono::Utc::now().to_rfc3339(), + app_id: Some(package_id_spawn.clone()), + }); + while data.notifications.len() > 20 { + data.notifications.remove(0); + } + }) + .await; + kick_scanner_and_wait(&handler).await; } } }); @@ -585,3 +607,34 @@ async fn kick_scanner_and_wait(handler: &RpcHandler) { }) .await; } + +fn finish_failed_update(entry: &mut crate::data_model::PackageDataEntry) { + if entry.state == PackageState::Updating { + entry.state = PackageState::Installed; + } + entry.install_progress = None; +} +#[cfg(test)] +mod update_completion_tests { + use super::*; + #[test] + fn failure_releases_spinner_without_inventing_stopped_or_restored_runtime() { + let mut entry = super::super::progress::create_installing_entry("movie"); + entry.state = PackageState::Updating; + finish_failed_update(&mut entry); + assert_eq!(entry.state, PackageState::Installed); + assert!(entry.install_progress.is_none()); + for actual in [ + PackageState::Running, + PackageState::Stopped, + PackageState::Exited, + ] { + entry.state = actual.clone(); + finish_failed_update(&mut entry); + assert_eq!( + entry.state, actual, + "Fresh scanner evidence must win over old pre-update intent" + ); + } + } +} diff --git a/core/archipelago/src/api/rpc/package/install.rs b/core/archipelago/src/api/rpc/package/install.rs index 05d8c92b..1609c037 100644 --- a/core/archipelago/src/api/rpc/package/install.rs +++ b/core/archipelago/src/api/rpc/package/install.rs @@ -283,6 +283,7 @@ impl RpcHandler { let lifecycle_guard = crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?; lifecycle_guard.require_clear()?; + lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?; let docker_image = params .get("dockerImage") diff --git a/core/archipelago/src/api/rpc/package/runtime.rs b/core/archipelago/src/api/rpc/package/runtime.rs index 88b16795..edba9f64 100644 --- a/core/archipelago/src/api/rpc/package/runtime.rs +++ b/core/archipelago/src/api/rpc/package/runtime.rs @@ -63,6 +63,7 @@ impl RpcHandler { let lifecycle_guard = crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?; lifecycle_guard.require_clear()?; + lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?; // A cuprate node that starts on a too-small disk fills it and takes // Archipelago down with it (no upstream pruning — see // dependencies::check_cuprate_disk_compatibility). Fail the start @@ -174,6 +175,7 @@ impl RpcHandler { let lifecycle_guard = crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?; lifecycle_guard.require_clear()?; + lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?; let single_orchestrator_app = self.orchestrator.is_some() && uses_single_orchestrator_app(package_id); @@ -280,6 +282,7 @@ impl RpcHandler { let lifecycle_guard = crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?; lifecycle_guard.require_clear()?; + lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?; // Restart is stop + recreate, so on a disk that shrank below the cuprate // minimum after install it resumes the doomed unprunable sync just like // start would — same gate, same "fail before clearing user-stopped / @@ -389,6 +392,7 @@ impl RpcHandler { let lifecycle_guard = crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?; lifecycle_guard.require_clear()?; + lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?; let preserve_data = params .get("preserve_data") .and_then(|v| v.as_bool()) diff --git a/core/archipelago/src/api/rpc/package/update.rs b/core/archipelago/src/api/rpc/package/update.rs index 4d73e58a..037f0f93 100644 --- a/core/archipelago/src/api/rpc/package/update.rs +++ b/core/archipelago/src/api/rpc/package/update.rs @@ -481,7 +481,9 @@ impl RpcHandler { if let Some(entry) = data.package_data.get_mut(package_id) { // Don't overwrite state from scanner — just clear if still Updating if entry.state == PackageState::Updating { - entry.state = PackageState::Stopped; + // Unknown is not stopped: the authoritative scanner will + // refresh actual retained runtime immediately in the wrapper. + entry.state = PackageState::Installed; } } self.state_manager.update_data(data).await; diff --git a/core/archipelago/src/container/mod.rs b/core/archipelago/src/container/mod.rs index 056e21da..98db75e1 100644 --- a/core/archipelago/src/container/mod.rs +++ b/core/archipelago/src/container/mod.rs @@ -33,3 +33,5 @@ pub use traits::ContainerOrchestrator; mod staged_update; pub(crate) mod update_transaction; + +pub(crate) mod supervised_update; diff --git a/core/archipelago/src/container/prod_orchestrator.rs b/core/archipelago/src/container/prod_orchestrator.rs index a5c0d696..fbb9478d 100644 --- a/core/archipelago/src/container/prod_orchestrator.rs +++ b/core/archipelago/src/container/prod_orchestrator.rs @@ -1994,6 +1994,16 @@ impl ProdContainerOrchestrator { return report; } }; + let held_names = match _update_guard.held_names() { + Ok(names) => names, + Err(error) => { + let mut report = ReconcileReport::default(); + report + .failures + .push(("update-recovery".into(), format!("{error:#}"))); + return report; + } + }; let user_stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await; // Durable desired-state signal: the container names that were running at // the last periodic snapshot. Used below to recreate a previously-running @@ -2015,6 +2025,7 @@ impl ProdContainerOrchestrator { let filtered = state .manifests .iter() + .filter(|(_, lm)| !held_names.contains(&compute_container_name(&lm.manifest))) .filter(|(app_id, _)| !state.disabled.contains(*app_id)) .filter(|(app_id, lm)| { dependency_required.contains(*app_id) @@ -3442,6 +3453,9 @@ impl ProdContainerOrchestrator { /// changes restart the service and retain a durable pending marker until /// that succeeds, including across daemon restarts and failed reloads. async fn sync_quadlet_unit(&self, lm: &LoadedManifest, name: &str) -> Result<()> { + if super::update_transaction::is_held(&self.data_dir, name)? { + return Ok(()); // Preserve the recovered unit instead of current catalog drift. + } // Companions: same reasoning as migrate_to_quadlet_if_needed — // companion.rs renders these units with a different shape, syncing // here would clobber them. diff --git a/core/archipelago/src/container/supervised_update.rs b/core/archipelago/src/container/supervised_update.rs new file mode 100644 index 00000000..37fe6a94 --- /dev/null +++ b/core/archipelago/src/container/supervised_update.rs @@ -0,0 +1,809 @@ +//! 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 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, + Committed, + Restored, +} +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Journal { + schema: u8, + id: String, + package: String, + phase: Phase, + members: Vec, +} + +pub(crate) trait Supervisor: Sync { + /// 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 read(&self, name: &str) -> impl Future> + Send; + fn write( + &self, + name: &str, + 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; +} +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) +} +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) + && 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 +} +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| 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, + }; + 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 { + 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]; + 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)?; + for member in record.members.iter().rev() { + supervisor.stop(&member.original.name).await?; + } + for member in &record.members { + supervisor + .write( + &member.original.name, + &[member.original.body.clone()], + &member.target_body, + ) + .await?; + } + supervisor.reload().await?; + record.phase = Phase::Starting; + save(guard, record)?; + for member in &record.members { + 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)?; + for member in &record.members { + guard.release_hold(&member.original.name, &record.id)?; + } + Ok(()) +} +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)?; + for member in &record.members { + guard.release_hold(&member.original.name, &record.id)?; + } + return Ok(()); + } + // Refuse to overwrite a foreign edit before stopping any surviving member. + for member in &record.members { + 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.name, + &[ + 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 { + 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) +} +pub(crate) async fn recover(guard: &Guard, supervisor: &impl Supervisor) -> Result<()> { + for mut record in records(guard)? { + match record.phase { + Phase::Committed | Phase::Aborted => { + for member in &record.members { + guard.release_hold(&member.original.name, &record.id)?; + } + } + Phase::Restored => {} // A later retry may own the current hold. + _ => restore(guard, &mut record, supervisor).await?, + } + } + Ok(()) +} + +#[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, + 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), + 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), + 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 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 read(&self, _name: &str) -> Result { + Ok(self.body.lock().unwrap().clone()) + } + async fn write(&self, _name: &str, 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) + } + } + #[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_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 record = records(&guard).unwrap().pop().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; + 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; + 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()); + } +} diff --git a/core/archipelago/src/container/update_transaction.rs b/core/archipelago/src/container/update_transaction.rs index c632e9a2..429b5f86 100644 --- a/core/archipelago/src/container/update_transaction.rs +++ b/core/archipelago/src/container/update_transaction.rs @@ -74,6 +74,9 @@ pub(crate) struct Guard { root: PathBuf, } impl Guard { + pub(crate) fn directory(&self) -> &Path { + &self.root + } pub(crate) fn acquire(data: &Path) -> Result { let root = data.join("update-transactions"); match std::fs::DirBuilder::new().mode(0o700).create(&root) { @@ -153,7 +156,76 @@ impl Guard { } Ok(records) } + pub(crate) fn hold(&self, name: &str, operation: &str) -> Result<()> { + anyhow::ensure!(name_ok(name), "Invalid held container name"); + uuid::Uuid::parse_str(operation)?; + let directory = self.root.join("holds"); + match std::fs::DirBuilder::new().mode(0o700).create(&directory) { + Ok(()) => {} + Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {} + Err(e) => return Err(e.into()), + } + let temporary = directory.join(format!(".{}.tmp", uuid::Uuid::new_v4())); + let mut file = std::fs::OpenOptions::new() + .create_new(true) + .write(true) + .mode(0o600) + .open(&temporary)?; + file.write_all(operation.as_bytes())?; + file.sync_all()?; + std::fs::rename(temporary, directory.join(name))?; + std::fs::File::open(&directory)?.sync_all()?; + std::fs::File::open(&self.root)?.sync_all()?; + Ok(()) + } + pub(crate) fn release_hold(&self, name: &str, operation: &str) -> Result<()> { + anyhow::ensure!(name_ok(name), "Invalid held container name"); + let path = self.root.join("holds").join(name); + match std::fs::read_to_string(&path) { + Ok(owner) if owner == operation => { + std::fs::remove_file(&path)?; + std::fs::File::open(path.parent().unwrap())?.sync_all()?; + } + Ok(_) => {} + Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} + Err(e) => return Err(e.into()), + } + Ok(()) + } + pub(crate) fn held_names(&self) -> Result> { + let directory = self.root.join("holds"); + let mut held = HashSet::new(); + let entries = match std::fs::read_dir(directory) { + Ok(entries) => entries, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(held), + Err(e) => return Err(e.into()), + }; + for entry in entries { + let entry = entry?; + let name = entry + .file_name() + .into_string() + .map_err(|_| anyhow::anyhow!("Invalid recovery hold name"))?; + if name.starts_with('.') { + continue; + } + anyhow::ensure!( + name_ok(&name) && entry.file_type()?.is_file() && entry.metadata()?.len() == 36, + "Damaged recovery hold" + ); + uuid::Uuid::parse_str(&std::fs::read_to_string(entry.path())?)?; + held.insert(name); + } + Ok(held) + } + pub(crate) fn require_unheld(&self, names: &[String]) -> Result<()> { + let held = self.held_names()?; + anyhow::ensure!(!names.iter().any(|name|held.contains(name)), + "Original app runtime is retained after rollback; retry its update before changing lifecycle"); + Ok(()) + } pub(crate) fn require_clear(&self) -> Result<()> { + super::supervised_update::require_clear(self)?; anyhow::ensure!( self.records()? .iter() @@ -260,6 +332,9 @@ pub(crate) async fn execute( members, }; guard.save(&record)?; + for member in &record.members { + guard.hold(&member.original.name, &record.operation)?; + } match replace(&guard,&mut record,runtime).await { Ok(())=>Ok(()), Err(error)=>match restore(&guard,&mut record,runtime).await { @@ -342,7 +417,11 @@ async fn replace(guard: &Guard, record: &mut Record, runtime: &impl Runtime) -> // Original backups are retained. Committing does not remove volumes or // silently discard the only recoverable original after a schema migration. record.phase = Phase::Committed; - guard.save(record) + guard.save(record)?; + for member in &record.members { + guard.release_hold(&member.original.name, &record.operation)?; + } + Ok(()) } async fn restore(guard: &Guard, record: &mut Record, runtime: &impl Runtime) -> Result<()> { // Preflight every original before mutating any replacement. Never infer @@ -401,8 +480,12 @@ async fn restore(guard: &Guard, record: &mut Record, runtime: &impl Runtime) -> "Original state was not restored" ); } + for member in &record.members { + guard.hold(&member.original.name, &record.operation)?; + } record.phase = Phase::Restored; - guard.save(record) + guard.save(record)?; + Ok(()) } /// Called before ordinary reconciliation. An unresolved recovery must prevent /// reconciliation from guessing a replacement or starting stopped originals. @@ -411,11 +494,32 @@ pub(crate) async fn recover(data: &Path, runtime: &impl Runtime) -> Result Result { + anyhow::ensure!(name_ok(name), "Invalid recovery name"); + let path = data.join("update-transactions/holds").join(name); + match std::fs::symlink_metadata(&path) { + Ok(meta) => { + anyhow::ensure!( + meta.is_file() && meta.len() == 36, + "Damaged update recovery hold" + ); + uuid::Uuid::parse_str(&std::fs::read_to_string(path)?)?; + Ok(true) + } + Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false), + Err(e) => Err(e.into()), + } +} + pub(crate) struct Podman; impl Podman { async fn output(args: &[&str]) -> Result { @@ -868,6 +972,27 @@ mod tests { assert!(recover(root.path(), &runtime).await.is_err()); assert!(runtime.calls.lock().unwrap().is_empty()); } + #[tokio::test] + async fn restored_original_remains_held_until_explicit_successful_update() { + let root = tempfile::tempdir().unwrap(); + let runtime = Mock::new(); + runtime.fail_health.store(true, Ordering::SeqCst); + let guard = Guard::acquire(root.path()).unwrap(); + assert!(execute(&guard, "stack", &Mock::targets(), &runtime) + .await + .is_err()); + assert!(is_held(root.path(), "db").unwrap()); + assert!(guard.require_unheld(&["db".into()]).is_err()); + drop(guard); + let guard = recover(root.path(), &runtime).await.unwrap(); + assert!(is_held(root.path(), "db").unwrap()); + runtime.fail_health.store(false, Ordering::SeqCst); + execute(&guard, "stack", &Mock::targets(), &runtime) + .await + .unwrap(); + assert!(!is_held(root.path(), "db").unwrap()); + guard.require_unheld(&["db".into(), "web".into()]).unwrap(); + } #[test] fn lifecycle_lock_excludes_competing_commands() { let root = tempfile::tempdir().unwrap();