diff --git a/core/archipelago/src/api/rpc/package/async_lifecycle.rs b/core/archipelago/src/api/rpc/package/async_lifecycle.rs index 80678e92..63ad037b 100644 --- a/core/archipelago/src/api/rpc/package/async_lifecycle.rs +++ b/core/archipelago/src/api/rpc/package/async_lifecycle.rs @@ -362,7 +362,11 @@ impl RpcHandler { set_package_state( &handler.state_manager, &package_id_spawn, - if result.get("status").and_then(|v| v.as_str()) == Some("up-to-date") { + if result.get("status").and_then(|v| v.as_str()) == Some("staged") { + PackageState::Stopped + } else if result.get("status").and_then(|v| v.as_str()) + == Some("up-to-date") + { pre_state.clone().unwrap_or(PackageState::Running) } else { PackageState::Running diff --git a/core/archipelago/src/api/rpc/package/config.rs b/core/archipelago/src/api/rpc/package/config.rs index 20cf592d..28586bfd 100644 --- a/core/archipelago/src/api/rpc/package/config.rs +++ b/core/archipelago/src/api/rpc/package/config.rs @@ -526,7 +526,18 @@ pub(in crate::api::rpc) async fn get_containers_for_app(package_id: &str) -> Res .await .context("podman ps timed out while listing containers")? .context("Failed to list containers")?; - let stdout = String::from_utf8_lossy(&output.stdout); + containers_from_list_output(package_id, &output) +} + +fn containers_from_list_output( + package_id: &str, + output: &std::process::Output, +) -> Result> { + anyhow::ensure!( + output.status.success(), + "podman ps failed while listing containers" + ); + let stdout = std::str::from_utf8(&output.stdout).context("Invalid container list response")?; let all: Vec<&str> = stdout.lines().filter(|s| !s.is_empty()).collect(); let patterns = all_container_names(package_id); @@ -543,6 +554,32 @@ pub(in crate::api::rpc) async fn get_containers_for_app(package_id: &str) -> Res mod tests { use super::{all_container_names, get_data_dirs_for_app, get_health_check_args}; + #[test] + fn failed_container_listing_is_not_an_absent_app() { + use std::os::unix::process::ExitStatusExt; + let mut output = std::process::Output { + status: std::process::ExitStatus::from_raw(1 << 8), + stdout: vec![], + stderr: b"store unavailable".to_vec(), + }; + assert!(super::containers_from_list_output("node-demo-music", &output).is_err()); + output.stdout = b"node-demo-music\n".to_vec(); + assert!(super::containers_from_list_output("node-demo-music", &output).is_err()); + output.status = std::process::ExitStatus::from_raw(0); + assert_eq!( + super::containers_from_list_output("node-demo-music", &output).unwrap(), + vec!["node-demo-music"] + ); + output.stdout.clear(); + assert!( + super::containers_from_list_output("node-demo-music", &output) + .unwrap() + .is_empty() + ); + output.stdout = vec![0xff]; + assert!(super::containers_from_list_output("node-demo-music", &output).is_err()); + } + #[test] fn bitcoin_variant_container_names_are_precise() { let core = all_container_names("bitcoin-core"); diff --git a/core/archipelago/src/api/rpc/package/runtime.rs b/core/archipelago/src/api/rpc/package/runtime.rs index ca081db9..9a46c547 100644 --- a/core/archipelago/src/api/rpc/package/runtime.rs +++ b/core/archipelago/src/api/rpc/package/runtime.rs @@ -126,7 +126,19 @@ impl RpcHandler { Err(e) => { tracing::error!("package.start {} failed: {:#}", package_id_owned, e); install_log(&format!("START FAIL: {} — {:#}", package_id_owned, e)).await; - if let Some(prev) = pre_state { + if e.downcast_ref::() + .is_some() + { + // Installed is neutral and scanner-owned. Restoring the + // prior Stopped state would conceal a failed cleanup; + // keeping Starting would prevent scanner convergence. + set_package_state( + &state_manager, + &package_id_owned, + PackageState::Installed, + ) + .await; + } else if let Some(prev) = pre_state { set_package_state(&state_manager, &package_id_owned, prev).await; } else { warn!( diff --git a/core/archipelago/src/api/rpc/package/update.rs b/core/archipelago/src/api/rpc/package/update.rs index f6daf30e..19957959 100644 --- a/core/archipelago/src/api/rpc/package/update.rs +++ b/core/archipelago/src/api/rpc/package/update.rs @@ -4,7 +4,7 @@ //! remove old container(s) → recreate (orchestrator-first, legacy fallback) → verify running. //! Data volumes are preserved (bind mounts, not stored in container). -use super::config::get_containers_for_app; +use super::config::{all_container_names, get_containers_for_app}; use super::install::install_log; use super::progress::parse_pull_progress; use super::runtime::stop_timeout_secs; @@ -13,6 +13,7 @@ use crate::api::rpc::RpcHandler; use crate::container::image_versions; use crate::data_model::{InstallPhase, PackageState}; use anyhow::{Context, Result}; +use std::{collections::HashSet, path::Path}; use tokio::io::{AsyncBufReadExt, BufReader}; use tracing::{error, info, warn}; @@ -60,9 +61,24 @@ impl RpcHandler { let targets = pinned .as_ref() .map(|target| self.resolve_images_to_pull(package_id, target)); + // A stopped Quadlet normally removes its --rm container. Absence is + // not an install decision: only a known single managed app with durable + // installed AND stopped evidence may enter the recreate path. + let known_managed = + if should_try_orchestrator_update(package_id, self.orchestrator.is_some()) { + self.orchestrator + .as_ref() + .expect("orchestrator presence checked") + .knows_app(orchestrator_update_app_id(package_id)) + .await + } else { + false + }; + let markers = UpdateMarkers::load(&self.config.data_dir, package_id).await?; + let installed = inspect_update_images(package_id).await?; + validate_update_presence(package_id, known_managed, !installed.is_empty(), &markers)?; if let Some(targets) = &targets { - let installed = inspect_update_images(package_id).await?; - if !update_targets_need_change(targets, &installed)? { + if !update_targets_need_change(targets, &installed)? && !markers.stopped { install_log(&format!( "UPDATE SKIP: {} — target versions already installed", package_id @@ -111,6 +127,19 @@ impl RpcHandler { if let Some(orchestrator) = self.orchestrator.as_ref() { match orchestrator.upgrade(orchestrator_app_id).await { Ok(()) => { + if let Some(image) = orchestrator + .staged_upgrade_image(orchestrator_app_id) + .await? + { + // The orchestrator proved the pinned image exists and + // persisted its reviewed manifest. No running container + // or healthy service is claimed for a stopped update. + self.clear_install_progress(package_id).await; + return Ok(serde_json::json!({ + "status": "staged", "state": "stopped", + "package_id": package_id, "image": image, + })); + } if let Some(targets) = &targets { verify_update_targets( targets, @@ -574,13 +603,69 @@ impl RpcHandler { } } -async fn inspect_update_images(package_id: &str) -> Result> { - let containers = get_containers_for_app(package_id).await?; +#[derive(Default)] +struct UpdateMarkers { + installed: bool, + stopped: bool, + uninstalled: bool, +} + +impl UpdateMarkers { + async fn load(data_dir: &Path, package_id: &str) -> Result { + let mut names = all_container_names(package_id); + names.push(package_id.to_string()); + names.push(orchestrator_update_app_id(package_id).to_string()); + let installed = read_update_markers(data_dir, "installed-apps.json").await?; + let stopped = read_update_markers(data_dir, "user-stopped.json").await?; + let uninstalled = read_update_markers(data_dir, "user-uninstalled.json").await?; + Ok(Self { + installed: names.iter().any(|name| installed.contains(name)), + stopped: names.iter().any(|name| stopped.contains(name)), + uninstalled: names.iter().any(|name| uninstalled.contains(name)), + }) + } +} + +// The recovery loaders intentionally return empty on malformed state. An update +// must instead distinguish unavailable evidence from a proven absence of an +// uninstall decision before it can recreate a removed Quadlet container. +async fn read_update_markers(data_dir: &Path, filename: &str) -> Result> { + match tokio::fs::read(data_dir.join(filename)).await { + Ok(bytes) => serde_json::from_slice(&bytes) + .with_context(|| format!("Cannot read {filename}; update cancelled")), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(HashSet::new()), + Err(error) => { + Err(error).with_context(|| format!("Cannot read {filename}; update cancelled")) + } + } +} + +fn validate_update_presence( + package_id: &str, + known_managed: bool, + has_containers: bool, + markers: &UpdateMarkers, +) -> Result<()> { anyhow::ensure!( - !containers.is_empty(), - "No containers found for {}", + !markers.uninstalled, + "{} was uninstalled; use an explicit install instead of update", package_id ); + anyhow::ensure!( + has_containers || (known_managed && markers.installed && markers.stopped), + "No containers or durable stopped installation found for {}", + package_id + ); + Ok(()) +} + +async fn inspect_update_images(package_id: &str) -> Result> { + let containers = get_containers_for_app(package_id).await?; + if containers.is_empty() { + // The caller decides whether durable installation evidence authorizes + // a missing runtime. Never invoke bare `podman inspect` here. + return Ok(Vec::new()); + } let mut command = tokio::process::Command::new("podman"); command.arg("inspect").args(&containers).kill_on_drop(true); let output = tokio::time::timeout(std::time::Duration::from_secs(30), command.output()) @@ -626,6 +711,7 @@ fn verify_update_targets( targets: &[(String, String)], installed: &[(String, String)], ) -> Result<()> { + anyhow::ensure!(!targets.is_empty(), "No update targets resolved"); for (app_id, target) in targets { let running = installed_image_for_target(app_id, installed).ok_or_else(|| { anyhow::anyhow!("Update {}: target container missing after recreate", app_id) @@ -651,6 +737,7 @@ fn update_targets_need_change( installed: &[(String, String)], ) -> Result { use std::cmp::Ordering; + anyhow::ensure!(!targets.is_empty(), "No update targets resolved"); let mut changed = false; for (app_id, target) in targets { let running = installed_image_for_target(app_id, installed); @@ -771,9 +858,83 @@ mod tests { use super::{ candidate_app_ids_for_container, immutable_update_image, orchestrator_update_app_id, should_try_orchestrator_update, update_targets_need_change, uses_legacy_update_flow, - verify_update_targets, + validate_update_presence, verify_update_targets, UpdateMarkers, }; + #[tokio::test] + async fn stopped_installed_managed_app_updates_after_quadlet_container_disappears() { + let root = tempfile::tempdir().unwrap(); + crate::crash_recovery::mark_installed(root.path(), "node-demo-music").await; + crate::crash_recovery::mark_user_stopped(root.path(), "node-demo-music").await; + let markers = UpdateMarkers::load(root.path(), "node-demo-music") + .await + .unwrap(); + validate_update_presence("node-demo-music", true, false, &markers).unwrap(); + let target = vec![("node-demo-music".into(), "localhost/music:2".into())]; + assert!(super::update_targets_need_change(&target, &[]).unwrap()); + // Success must still prove that upgrade actually created the target. + assert!(verify_update_targets(&target, &[]).is_err()); + assert!(verify_update_targets(&target, &target).is_ok()); + assert!(crate::crash_recovery::load_user_stopped(root.path()) + .await + .contains("node-demo-music")); + } + + #[test] + fn absent_catalog_app_or_unmanaged_runtime_is_not_an_update_installation() { + for (managed, installed, stopped) in [ + (true, false, false), + (true, false, true), + (true, true, false), + (false, true, true), + ] { + let markers = UpdateMarkers { + installed, + stopped, + uninstalled: false, + }; + assert!(validate_update_presence("optional", managed, false, &markers).is_err()); + } + // Existing legacy containers remain updateable without modern markers. + assert!(validate_update_presence("legacy", false, true, &UpdateMarkers::default()).is_ok()); + assert!(super::update_targets_need_change(&[], &[]).is_err()); + assert!(verify_update_targets(&[], &[]).is_err()); + } + + #[tokio::test] + async fn uninstall_tombstone_beats_stale_installed_and_stopped_markers() { + let root = tempfile::tempdir().unwrap(); + crate::crash_recovery::mark_installed(root.path(), "node-demo-music").await; + crate::crash_recovery::mark_user_stopped(root.path(), "node-demo-music").await; + crate::crash_recovery::mark_user_uninstalled(root.path(), "archy-node-demo-music").await; + let markers = UpdateMarkers::load(root.path(), "node-demo-music") + .await + .unwrap(); + for present in [false, true] { + assert!(validate_update_presence("node-demo-music", true, present, &markers).is_err()); + } + } + + #[tokio::test] + async fn damaged_lifecycle_markers_cannot_authorize_recreation() { + let root = tempfile::tempdir().unwrap(); + let absent = UpdateMarkers::load(root.path(), "optional").await.unwrap(); + assert!(validate_update_presence("optional", true, false, &absent).is_err()); + for name in [ + "installed-apps.json", + "user-stopped.json", + "user-uninstalled.json", + ] { + tokio::fs::write(root.path().join(name), b"not-json") + .await + .unwrap(); + assert!(UpdateMarkers::load(root.path(), "optional").await.is_err()); + tokio::fs::remove_file(root.path().join(name)) + .await + .unwrap(); + } + } + #[tokio::test] async fn stack_image_failure_precedes_every_lifecycle_action() { use std::sync::{Arc, Mutex}; diff --git a/core/archipelago/src/container/docker_packages.rs b/core/archipelago/src/container/docker_packages.rs index 2b341b80..1fe2f1d0 100644 --- a/core/archipelago/src/container/docker_packages.rs +++ b/core/archipelago/src/container/docker_packages.rs @@ -298,7 +298,22 @@ impl DockerPackageScanner { uninstall_stage: None, }; - apply_manifest_presentation(&app_id, &mut package); + match super::staged_update::installed_manifest(data_dir, &app_id, &container.image) + .await + { + Ok(Some(manifest)) => { + apply_manifest_value(&serde_json::to_value(manifest)?, &mut package) + } + Ok(None) => apply_manifest_presentation(&app_id, &mut package), + Err(error) => { + tracing::warn!(app_id, error = %error, "Cannot verify retained installed manifest entry point"); + // Do not invent a launch path from a later catalog when its + // binding to the running version cannot be established. + if let Some(installed) = package.installed.as_mut() { + installed.interface_addresses.clear(); + } + } + } packages.insert(app_id.clone(), package); info!( "Detected container: {} ({})", diff --git a/core/archipelago/src/container/hooks.rs b/core/archipelago/src/container/hooks.rs index 4541792c..0ca32aa1 100644 --- a/core/archipelago/src/container/hooks.rs +++ b/core/archipelago/src/container/hooks.rs @@ -85,6 +85,18 @@ pub async fn run_post_install(manifest: &AppManifest, container_name: &str, data } } +/// Strict completion for an explicitly staged update. Keep the durable pin on failure. +pub(super) async fn run_post_install_strict( + manifest: &AppManifest, + container: &str, + data_dir: &Path, +) -> Result<()> { + for step in &manifest.app.hooks.post_install { + run_step(step, container, &manifest.app.id, data_dir).await?; + } + Ok(()) +} + async fn run_step(step: &HookStep, container: &str, app_id: &str, data_dir: &Path) -> Result<()> { match step { HookStep::Exec { exec } => { diff --git a/core/archipelago/src/container/mod.rs b/core/archipelago/src/container/mod.rs index 60a71a1a..0153fd0d 100644 --- a/core/archipelago/src/container/mod.rs +++ b/core/archipelago/src/container/mod.rs @@ -1,5 +1,4 @@ pub mod app_catalog; -pub mod node_catalog; pub mod app_gate_config; pub mod bitcoin_ui; pub mod boot_reconciler; @@ -14,11 +13,12 @@ pub mod image_policy; pub mod image_versions; pub mod lnd; pub mod migration_backup; +pub mod node_catalog; pub mod npm; pub mod prod_orchestrator; pub mod quadlet; -pub mod registry; pub mod registration_pin; +pub mod registry; pub mod secrets; pub mod traits; pub mod ui_detection; @@ -29,3 +29,5 @@ pub use dev_orchestrator::DevContainerOrchestrator; pub use docker_packages::DockerPackageScanner; pub use prod_orchestrator::ProdContainerOrchestrator; pub use traits::ContainerOrchestrator; + +mod staged_update; diff --git a/core/archipelago/src/container/prod_orchestrator.rs b/core/archipelago/src/container/prod_orchestrator.rs index 2dfcbbbc..2c547f16 100644 --- a/core/archipelago/src/container/prod_orchestrator.rs +++ b/core/archipelago/src/container/prod_orchestrator.rs @@ -1422,6 +1422,17 @@ fn host_port_collisions<'m>( out } +/// A staged start failed and its runtime could not be proven stopped. Callers +/// must not restore a cached Stopped state; the scanner/reconciler owns recovery. +#[derive(Debug)] +pub(crate) struct StagedCleanupFailure(String); +impl std::fmt::Display for StagedCleanupFailure { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + std::fmt::Display::fmt(&self.0, f) + } +} +impl std::error::Error for StagedCleanupFailure {} + /// Internal: track a manifest together with the absolute directory it was loaded /// from, so Build sources can resolve relative `context:` paths. #[derive(Debug, Clone)] @@ -2043,7 +2054,9 @@ impl ProdContainerOrchestrator { } } } - let mut report = ReconcileReport::default(); + // Stopped/disabled apps are excluded from ordinary start reconciliation. + // Pending cleanup must still run for them, including after a restart. + let mut report = self.reconcile_staged_updates().await; let disk_gb = self.disk_gb().await; let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB || crate::settings::bitcoin_storage::load(&self.data_dir) @@ -2336,6 +2349,12 @@ impl ProdContainerOrchestrator { let lock = self.app_lock(&app_id).await; let _guard = lock.lock().await; + if let Some(record) = super::staged_update::load(&self.data_dir, &app_id).await? { + let name = compute_container_name(&record.manifest); + self.stop_staged_runtime(&name).await?; + return Ok(ReconcileAction::Left("staged-update-awaiting-start".into())); + } + self.ensure_app_secrets(&app_id).await?; // Don't fight the Bitcoin-implementation switch: bitcoin-core and @@ -2849,11 +2868,171 @@ impl ProdContainerOrchestrator { } /// Build-or-pull, create, start. Assumes the per-app mutex is already held. + async fn reconcile_staged_updates(&self) -> ReconcileReport { + let mut report = ReconcileReport::default(); + let ids = match super::staged_update::pending_ids(&self.data_dir).await { + Ok(ids) => ids, + Err(error) => { + report + .failures + .push(("staged-updates".into(), error.to_string())); + return report; + } + }; + for app_id in ids { + if crate::app_ops::lifecycle_op_in_flight(&app_id) { + continue; + } + let lock = self.app_lock(&app_id).await; + let _guard = lock.lock().await; + // Re-read after acquiring the lock: a successful explicit start + // may have completed and cleared the stage while this pass waited. + let result = async { + if let Some(record) = super::staged_update::load(&self.data_dir, &app_id).await? { + self.stop_staged_runtime(&compute_container_name(&record.manifest)) + .await?; + report.record( + &app_id, + ReconcileAction::Left("staged-update-awaiting-start".into()), + ); + } + Ok::<_, anyhow::Error>(()) + } + .await; + if let Err(error) = result { + report.failures.push((app_id, error.to_string())); + } + } + report + } + + async fn stop_staged_runtime(&self, name: &str) -> Result<()> { + let mut errors = Vec::new(); + if let Err(e) = self.remove_quadlet_unit_if_present(name).await { + errors.push(format!("could not disable staged unit: {e:#}")); + } + let inventory = self + .runtime + .list_containers() + .await + .context("Cannot verify staged runtime inventory")?; + if inventory.iter().any(|c| { + c.name.trim_start_matches('/') == name + && !matches!( + c.state, + ContainerState::Stopped | ContainerState::Exited | ContainerState::Created + ) + }) { + if let Err(e) = self.runtime.stop_container(name).await { + errors.push(format!("could not stop staged container: {e:#}")); + } + } + let inventory = self + .runtime + .list_containers() + .await + .context("Cannot verify staged runtime stopped")?; + if inventory.iter().any(|c| { + c.name.trim_start_matches('/') == name + && !matches!( + c.state, + ContainerState::Stopped | ContainerState::Exited | ContainerState::Created + ) + }) { + errors.push("staged container is still active".into()); + } + anyhow::ensure!(errors.is_empty(), "{}", errors.join("; ")); + Ok(()) + } + + async fn start_staged_update(&self, app_id: &str) -> Result { + use super::staged_update; + if staged_update::load(&self.data_dir, app_id).await?.is_none() { + return Ok(false); + } + let lock = self.app_lock(app_id).await; + let _guard = lock.lock().await; + let record = staged_update::load(&self.data_dir, app_id) + .await? + .context("Staged update changed; retry start")?; + let name = compute_container_name(&record.manifest); + anyhow::ensure!( + !staged_update::marked(&self.data_dir, "user-uninstalled.json", app_id, &name).await?, + "App was uninstalled; staged update cannot reinstall it" + ); + anyhow::ensure!( + record.ready, + "Update preparation was interrupted; retry update before start" + ); + anyhow::ensure!( + self.runtime.image_exists(record.image()?).await?, + "Staged image is unavailable; refusing to start another version" + ); + // Clear only this explicit start's stopped intent. The pending manifest + // remains durable until all creation, hooks and readiness checks succeed. + crate::crash_recovery::clear_user_stopped(&self.data_dir, app_id).await; + crate::crash_recovery::clear_user_stopped(&self.data_dir, &name).await; + crate::crash_recovery::clear_user_stopped(&self.data_dir, &format!("archy-{app_id}")).await; + self.state.write().await.disabled.remove(app_id); + let lm = LoadedManifest { + manifest: record.manifest.clone(), + manifest_dir: record.manifest_dir.clone(), + }; + let result = async { + self.remove_quadlet_unit_if_present(&name).await?; + if self + .runtime + .list_containers() + .await? + .iter() + .any(|c| c.name.trim_start_matches('/') == name) + { + self.runtime.stop_container(&name).await?; + self.runtime.remove_container(&name).await?; + } + self.install_fresh_with_pin(&lm, true).await?; + let observed = self.runtime.get_container_status(&name).await?; + anyhow::ensure!( + observed.state == ContainerState::Running, + "Pinned start did not become running" + ); + anyhow::ensure!( + staged_update::matches_image(record.image()?, &observed.image), + "Pinned start produced a different image" + ); + // Installed presentation is durable before pending intent disappears. + // A crash on either side of clear cannot fall back to a later catalog. + staged_update::save_installed(&self.data_dir, &record).await?; + staged_update::clear(&self.data_dir, app_id).await + } + .await; + if let Err(error) = result { + crate::crash_recovery::mark_user_stopped(&self.data_dir, app_id).await; + self.state.write().await.disabled.insert(app_id.into()); + let retained = staged_update::save(&self.data_dir, &record).await; + let cleanup = self.stop_staged_runtime(&name).await; + if let Err(cleanup) = cleanup { + return Err(StagedCleanupFailure(format!( + "Staged start failed: {error:#}; cleanup failed: {cleanup:#}; pin persistence: {:?}", retained.err() + )).into()); + } + retained.context("Could not retain failed staged start pin")?; + return Err(error); + } + Ok(true) + } + async fn install_fresh(&self, lm: &LoadedManifest) -> Result<()> { + self.install_fresh_with_pin(lm, false).await + } + + async fn install_fresh_with_pin(&self, lm: &LoadedManifest, pinned: bool) -> Result<()> { self.ensure_app_secrets(&lm.manifest.app.id).await?; let mut resolved_manifest = lm.manifest.clone(); self.resolve_dynamic_env(&mut resolved_manifest).await?; - resolve_catalog_image(&mut resolved_manifest); + if !pinned { + resolve_catalog_image(&mut resolved_manifest); + } let resolved = resolved_manifest.app.container.resolve().ok_or_else(|| { anyhow::anyhow!( @@ -3008,7 +3187,17 @@ impl ProdContainerOrchestrator { // freshly created container — exactly when container mutations (e.g. // indeedhub's nginx X-Frame-Options strip + nostr-provider injection) must // be re-applied. Best-effort + idempotent: never fails the install. - crate::container::hooks::run_post_install(&resolved_manifest, &name, &self.data_dir).await; + if pinned { + crate::container::hooks::run_post_install_strict( + &resolved_manifest, + &name, + &self.data_dir, + ) + .await?; + } else { + crate::container::hooks::run_post_install(&resolved_manifest, &name, &self.data_dir) + .await; + } if uses_pasta_network(&resolved_manifest) { if let Err(err) = wait_for_manifest_host_ports( &resolved_manifest, @@ -4825,6 +5014,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { } async fn install(&self, app_id: &str) -> Result { + anyhow::ensure!(super::staged_update::load(&self.data_dir, app_id).await?.is_none(), + "An update is staged; use explicit start, or uninstall before selecting another installation"); let lm = self.loaded(app_id).await?; // Optional shared-service preconditions are checked before recording // installation or creating anything. A headless adapter must not claim @@ -4896,7 +5087,7 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { // point). Just delegate. self.prepare_for_start(&lm.manifest).await?; let action = self.ensure_running(&lm).await?; - match action { + let result = match action { ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed => { Ok(name) } @@ -4915,10 +5106,17 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { self.install_fresh(&lm).await?; Ok(name) } + }; + if result.is_ok() { + super::staged_update::clear_installed(&self.data_dir, app_id).await?; } + result } async fn start(&self, app_id: &str) -> Result<()> { + if self.start_staged_update(app_id).await? { + return Ok(()); + } if let Some(members) = self.mempool_umbrella_members(app_id).await { tracing::info!( app_id, @@ -5136,6 +5334,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { self.remove_quadlet_unit_if_present(&name).await?; } self.state.write().await.disabled.insert(app_id.to_string()); + super::staged_update::clear(&self.data_dir, app_id).await?; + super::staged_update::clear_installed(&self.data_dir, app_id).await?; return Ok(()); } let lm = self.loaded(app_id).await?; @@ -5189,15 +5389,83 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { // stale claim behind would let desired-state recovery recreate the very // app that was just uninstalled. crate::crash_recovery::clear_installed(&self.data_dir, app_id).await; + super::staged_update::clear(&self.data_dir, app_id).await?; + super::staged_update::clear_installed(&self.data_dir, app_id).await?; Ok(()) } /// Upgrade: stop-remove-reinstall (re-pulls or rebuilds as required). async fn upgrade(&self, app_id: &str) -> Result<()> { - let lm = self.loaded(app_id).await?; + use super::staged_update; let lock = self.app_lock(app_id).await; let _guard = lock.lock().await; + let existing = staged_update::load(&self.data_dir, app_id).await?; + let lm = if let Some(record) = &existing { + LoadedManifest { + manifest: record.manifest.clone(), + manifest_dir: record.manifest_dir.clone(), + } + } else { + self.loaded(app_id).await? + }; let name = compute_container_name(&lm.manifest); + anyhow::ensure!( + !staged_update::marked(&self.data_dir, "user-uninstalled.json", app_id, &name).await?, + "App was uninstalled; update cannot reinstall it" + ); + let stopped = + staged_update::marked(&self.data_dir, "user-stopped.json", app_id, &name).await?; + if stopped || existing.is_some() { + anyhow::ensure!( + staged_update::marked(&self.data_dir, "installed-apps.json", app_id, &name).await?, + "Stopped update needs durable installation evidence" + ); + let mut record = match existing { + Some(record) => record, + None => { + let mut manifest = lm.manifest.clone(); + resolve_catalog_image(&mut manifest); + staged_update::StagedUpdate { + manifest, + manifest_dir: lm.manifest_dir.clone(), + ready: false, + } + } + }; + record.image()?; + crate::crash_recovery::mark_user_stopped(&self.data_dir, app_id).await; + anyhow::ensure!( + staged_update::marked(&self.data_dir, "user-stopped.json", app_id, &name).await?, + "Could not persist stopped update intent" + ); + // Freeze the approved manifest before any pull. Retries cannot + // silently switch to another catalog revision after an interruption. + record.ready = false; + staged_update::save(&self.data_dir, &record).await?; + self.remove_quadlet_unit_if_present(&name).await?; + if self + .runtime + .list_containers() + .await? + .iter() + .any(|c| c.name.trim_start_matches('/') == name) + { + self.runtime.stop_container(&name).await?; + self.runtime.remove_container(&name).await?; + } + let pinned = LoadedManifest { + manifest: record.manifest.clone(), + manifest_dir: record.manifest_dir.clone(), + }; + self.ensure_resolved_source_available(&pinned).await?; + anyhow::ensure!( + self.runtime.image_exists(record.image()?).await?, + "Staged image unavailable after preparation" + ); + record.ready = true; + staged_update::save(&self.data_dir, &record).await?; + return Ok(()); + } let mut resolved = lm.manifest.clone(); resolve_catalog_image(&mut resolved); if resolved.app.container.build.is_none() { @@ -5213,7 +5481,10 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { running.image, target ), - Some(std::cmp::Ordering::Equal) => return Ok(()), + Some(std::cmp::Ordering::Equal) => { + staged_update::clear_installed(&self.data_dir, app_id).await?; + return Ok(()); + } _ => {} } } @@ -5221,7 +5492,22 @@ impl ContainerOrchestrator for ProdContainerOrchestrator { } let _ = self.runtime.stop_container(&name).await; let _ = self.runtime.remove_container(&name).await; - self.install_fresh(&lm).await + self.install_fresh(&lm).await?; + staged_update::clear_installed(&self.data_dir, app_id).await + } + + async fn staged_upgrade_image(&self, app_id: &str) -> Result> { + match super::staged_update::load(&self.data_dir, app_id).await? { + Some(record) => { + anyhow::ensure!(record.ready, "Staged update is not ready"); + anyhow::ensure!( + self.runtime.image_exists(record.image()?).await?, + "Staged image unavailable" + ); + Ok(Some(record.image()?.to_string())) + } + None => Ok(None), + } } async fn status(&self, app_id: &str) -> Result { @@ -5778,6 +6064,7 @@ mod tests { fail_image_exists: StdMutex>, /// If set, `start_container` for this container fails with this message. fail_start: StdMutex>, + fail_stop: StdMutex>, } impl MockRuntime { @@ -5841,6 +6128,12 @@ mod tests { ) -> Result { self.record(format!("create_container:{name}:offset={port_offset}")); self.set_state(name, ContainerState::Created); + if let Some(image) = manifest.app.container.image_ref() { + self.running_images + .lock() + .unwrap() + .insert(name.into(), image); + } self.created_env .lock() .unwrap() @@ -5862,6 +6155,9 @@ mod tests { } async fn stop_container(&self, name: &str) -> Result<()> { self.record(format!("stop_container:{name}")); + if let Some(error) = self.fail_stop.lock().unwrap().remove(name) { + return Err(anyhow::anyhow!(error)); + } self.set_state(name, ContainerState::Stopped); Ok(()) } @@ -6310,6 +6606,333 @@ app: } } + async fn stopped_update_fixture() -> (Arc, ProdContainerOrchestrator, String) { + let rt = Arc::new(MockRuntime::default()); + let orch = orch_with(rt.clone()).await; + let image = format!("docker.io/library/alpine@sha256:{}", "a".repeat(64)); + let mut manifest = pull_manifest("staged-music", &image); + manifest + .app + .environment + .push("STAGED_VERSION=original".into()); + orch.insert_manifest_for_test(manifest, PathBuf::from("/tmp")) + .await; + crate::crash_recovery::mark_installed(&orch.data_dir, "staged-music").await; + crate::crash_recovery::mark_user_stopped(&orch.data_dir, "staged-music").await; + (rt, orch, image) + } + + #[tokio::test] + async fn stopped_upgrade_stages_without_start_then_starts_pinned_manifest() { + let (rt, orch, image) = stopped_update_fixture().await; + let sentinel = orch.data_dir.join("persistent-song-and-secret-fixture"); + tokio::fs::write(&sentinel, b"preserved").await.unwrap(); + orch.upgrade("staged-music").await.unwrap(); + assert_eq!( + orch.staged_upgrade_image("staged-music").await.unwrap(), + Some(image) + ); + assert!(!rt + .calls() + .iter() + .any(|c| c.starts_with("create_container:") || c.starts_with("start_container:"))); + let replacement = pull_manifest( + "staged-music", + &format!("docker.io/library/alpine@sha256:{}", "b".repeat(64)), + ); + orch.insert_manifest_for_test(replacement, PathBuf::from("/tmp/new-catalog")) + .await; + let lm = orch.loaded("staged-music").await.unwrap(); + assert_eq!( + orch.ensure_running(&lm).await.unwrap(), + ReconcileAction::Left("staged-update-awaiting-start".into()) + ); + orch.start("staged-music").await.unwrap(); + assert!(rt + .created_env_for("staged-music") + .contains(&"STAGED_VERSION=original".into())); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Running + ); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_none()); + assert!(!crate::crash_recovery::load_user_stopped(&orch.data_dir) + .await + .contains("staged-music")); + assert_eq!(tokio::fs::read(sentinel).await.unwrap(), b"preserved"); + } + + #[tokio::test] + async fn interrupted_stopped_pull_retains_original_pin_and_requires_retry() { + let (rt, orch, image) = stopped_update_fixture().await; + *rt.fail_pull.lock().unwrap() = Some("interrupted pull".into()); + assert!(orch.upgrade("staged-music").await.is_err()); + assert!(orch + .start("staged-music") + .await + .unwrap_err() + .to_string() + .contains("interrupted")); + let record = super::super::staged_update::load(&orch.data_dir, "staged-music") + .await + .unwrap() + .unwrap(); + assert!(!record.ready); + assert_eq!(record.image().unwrap(), image); + *rt.fail_pull.lock().unwrap() = None; + orch.insert_manifest_for_test( + pull_manifest("staged-music", "docker.io/library/alpine:latest"), + PathBuf::from("/tmp"), + ) + .await; + orch.upgrade("staged-music").await.unwrap(); + assert_eq!( + orch.staged_upgrade_image("staged-music").await.unwrap(), + Some(image) + ); + assert!(!rt.calls().iter().any(|c| c.starts_with("start_container:"))); + } + + #[tokio::test] + async fn failed_staged_start_preserves_pin_and_stop_until_retry_succeeds() { + let (rt, orch, _) = stopped_update_fixture().await; + orch.upgrade("staged-music").await.unwrap(); + rt.fail_start + .lock() + .unwrap() + .insert("staged-music".into(), "injected start failure".into()); + assert!(orch.start("staged-music").await.is_err()); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_some()); + assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir) + .await + .contains("staged-music")); + assert!(matches!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Stopped | ContainerState::Created + )); + orch.start("staged-music").await.unwrap(); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_none()); + } + + #[tokio::test] + async fn staged_missing_image_uninstall_and_damaged_record_refuse_start() { + let (rt, orch, image) = stopped_update_fixture().await; + orch.upgrade("staged-music").await.unwrap(); + rt.images.lock().unwrap().remove(&image); + assert!(orch + .start("staged-music") + .await + .unwrap_err() + .to_string() + .contains("unavailable")); + rt.mark_image_present(&image); + crate::crash_recovery::mark_user_uninstalled(&orch.data_dir, "staged-music").await; + assert!(orch + .start("staged-music") + .await + .unwrap_err() + .to_string() + .contains("uninstalled")); + assert!(orch.upgrade("staged-music").await.is_err()); + tokio::fs::write( + orch.data_dir.join("staged-updates/staged-music.json"), + b"damaged", + ) + .await + .unwrap(); + assert!(orch.start("staged-music").await.is_err()); + assert!(!rt.calls().iter().any(|c| c.starts_with("start_container:"))); + } + + #[tokio::test] + async fn staged_failed_post_install_hook_keeps_pin_and_stops_container() { + let (rt, orch, _) = stopped_update_fixture().await; + let mut lm = orch.loaded("staged-music").await.unwrap(); + lm.manifest.app.hooks.post_install = serde_yaml::from_str( + "- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n", + ) + .unwrap(); + orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir) + .await; + orch.upgrade("staged-music").await.unwrap(); + assert!(orch.start("staged-music").await.is_err()); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_some()); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Stopped + ); + assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir) + .await + .contains("staged-music")); + } + + #[tokio::test] + async fn failed_staged_cleanup_reports_active_runtime_and_reconcile_retries_stop() { + let (rt, orch, _) = stopped_update_fixture().await; + let mut lm = orch.loaded("staged-music").await.unwrap(); + lm.manifest.app.hooks.post_install = serde_yaml::from_str( + "- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n", + ) + .unwrap(); + orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir) + .await; + orch.upgrade("staged-music").await.unwrap(); + rt.fail_stop + .lock() + .unwrap() + .insert("staged-music".into(), "stop refused".into()); + let error = orch.start("staged-music").await.unwrap_err(); + assert!(error.downcast_ref::().is_some()); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Running + ); + let report = orch.reconcile_all().await; + assert!(report.failures.is_empty()); + assert!(report + .actions + .iter() + .any(|(app, action)| app == "staged-music" + && action == &ReconcileAction::Left("staged-update-awaiting-start".into()))); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Stopped + ); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_some()); + } + + #[tokio::test] + async fn pending_stage_enforces_stop_after_restart_when_stop_marker_cannot_be_saved() { + let (rt, orch, _) = stopped_update_fixture().await; + let mut lm = orch.loaded("staged-music").await.unwrap(); + lm.manifest.app.hooks.post_install = serde_yaml::from_str( + "- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n", + ) + .unwrap(); + orch.insert_manifest_for_test(lm.manifest.clone(), lm.manifest_dir.clone()) + .await; + orch.upgrade("staged-music").await.unwrap(); + let marker = orch.data_dir.join("user-stopped.json"); + tokio::fs::remove_file(&marker).await.unwrap(); + tokio::fs::create_dir(&marker).await.unwrap(); + rt.fail_stop + .lock() + .unwrap() + .insert("staged-music".into(), "stop interrupted".into()); + let failure = orch.start("staged-music").await.unwrap_err(); + assert!(failure.downcast_ref::().is_some()); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Running + ); + let starts = rt + .calls() + .iter() + .filter(|c| c.starts_with("start_container:")) + .count(); + let mut resumed = orch_with(rt.clone()).await; + resumed.set_data_dir(orch.data_dir.clone()); + resumed + .insert_manifest_for_test(lm.manifest, lm.manifest_dir) + .await; + let report = resumed.reconcile_all().await; + assert!(report.failures.is_empty()); + assert_eq!( + rt.get_container_status("staged-music").await.unwrap().state, + ContainerState::Stopped + ); + assert_eq!( + rt.calls() + .iter() + .filter(|c| c.starts_with("start_container:")) + .count(), + starts + ); + assert!(resumed + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_some()); + } + + #[tokio::test] + async fn staged_configuration_failure_keeps_stop_and_pin_before_create() { + let (rt, mut orch, _) = stopped_update_fixture().await; + orch.secrets_dir = orch.data_dir.join("fixture-secrets"); + let mut lm = orch.loaded("staged-music").await.unwrap(); + lm.manifest.app.container.secret_env = + serde_yaml::from_str("- key: REQUIRED_SECRET\n secret_file: missing-staged-secret\n") + .unwrap(); + orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir) + .await; + orch.upgrade("staged-music").await.unwrap(); + assert!(orch.start("staged-music").await.is_err()); + assert!(orch + .staged_upgrade_image("staged-music") + .await + .unwrap() + .is_some()); + assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir) + .await + .contains("staged-music")); + assert!(!rt + .calls() + .iter() + .any(|c| c.starts_with("create_container:"))); + } + + #[tokio::test] + async fn staged_record_persistence_failure_precedes_runtime_changes() { + let (rt, orch, _) = stopped_update_fixture().await; + tokio::fs::write(orch.data_dir.join("staged-updates"), b"blocked directory") + .await + .unwrap(); + assert!(orch.upgrade("staged-music").await.is_err()); + assert!(!rt.calls().iter().any(|c| c.starts_with("pull_image:") + || c.starts_with("stop_container:") + || c.starts_with("remove_container:") + || c.starts_with("create_container:"))); + } + + #[tokio::test] + async fn stopped_unpinned_update_refuses_without_starting() { + let (rt, orch, _) = stopped_update_fixture().await; + orch.insert_manifest_for_test( + pull_manifest("staged-music", "docker.io/library/alpine:3.20"), + PathBuf::from("/tmp"), + ) + .await; + assert!(orch + .upgrade("staged-music") + .await + .unwrap_err() + .to_string() + .contains("digest-pinned")); + assert!(!rt.calls().iter().any(|c| c.starts_with("pull_image:") + || c.starts_with("create_container:") + || c.starts_with("start_container:"))); + } + #[tokio::test] async fn install_fresh_pull() { let rt = Arc::new(MockRuntime::default()); diff --git a/core/archipelago/src/container/staged_update.rs b/core/archipelago/src/container/staged_update.rs new file mode 100644 index 00000000..f0e022ff --- /dev/null +++ b/core/archipelago/src/container/staged_update.rs @@ -0,0 +1,314 @@ +//! Durable, immutable preparation for updates that must remain stopped. +use anyhow::{Context, Result}; +use archipelago_container::AppManifest; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::{ + collections::HashSet, + path::{Path, PathBuf}, +}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub(super) struct StagedUpdate { + pub manifest: AppManifest, + pub manifest_dir: PathBuf, + pub ready: bool, +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Envelope { + schema: u32, + payload: String, + checksum: String, +} + +impl StagedUpdate { + pub fn image(&self) -> Result<&str> { + anyhow::ensure!( + self.manifest.app.container.build.is_none(), + "Stopped updates of build-based apps require immutable image resolution" + ); + let image = self + .manifest + .app + .container + .image + .as_deref() + .context("Missing staged image")?; + let digest = image.rsplit_once("@sha256:").map(|(_, d)| d); + anyhow::ensure!( + digest.is_some_and(|d| d.len() == 64 && d.bytes().all(|b| b.is_ascii_hexdigit())), + "Stopped update requires a digest-pinned image; app remains stopped" + ); + Ok(image) + } +} + +fn record_path(root: &Path, directory: &str, app: &str) -> Result { + anyhow::ensure!( + !app.is_empty() + && app.len() <= 128 + && app + .bytes() + .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_')), + "Invalid staged app id" + ); + Ok(root.join(directory).join(format!("{app}.json"))) +} + +fn path(root: &Path, app: &str) -> Result { + record_path(root, "staged-updates", app) +} + +pub(super) async fn load(root: &Path, app: &str) -> Result> { + load_at(root, "staged-updates", app).await +} + +async fn load_at(root: &Path, directory: &str, app: &str) -> Result> { + let bytes = match tokio::fs::read(record_path(root, directory, app)?).await { + Ok(bytes) => bytes, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None), + Err(e) => return Err(e.into()), + }; + let envelope: Envelope = serde_json::from_slice(&bytes) + .context("Damaged staged update; refusing another version")?; + anyhow::ensure!(envelope.schema == 1, "Unsupported staged update schema"); + anyhow::ensure!( + envelope.checksum == hex::encode(Sha256::digest(envelope.payload.as_bytes())), + "Staged update checksum mismatch" + ); + let record: StagedUpdate = + serde_json::from_str(&envelope.payload).context("Damaged staged manifest")?; + record + .manifest + .validate() + .context("Invalid staged manifest")?; + anyhow::ensure!( + record.manifest.app.id == app, + "Staged update belongs to another app" + ); + record.image()?; + Ok(Some(record)) +} + +pub(super) async fn save(root: &Path, record: &StagedUpdate) -> Result<()> { + save_at(root, "staged-updates", record).await +} + +pub(super) async fn save_installed(root: &Path, record: &StagedUpdate) -> Result<()> { + anyhow::ensure!(record.ready, "Cannot publish incomplete installed manifest"); + save_at(root, "installed-manifests", record).await +} + +pub(super) async fn installed_manifest( + root: &Path, + app: &str, + observed_image: &str, +) -> Result> { + let Some(record) = load_at(root, "installed-manifests", app).await? else { + return Ok(None); + }; + let expected = record.image()?; + let matches = matches_image(expected, observed_image); + Ok((record.ready && matches).then_some(record.manifest)) +} + +pub(super) fn matches_image(expected: &str, observed: &str) -> bool { + expected == observed + || match ( + expected.rsplit_once("@sha256:"), + observed.rsplit_once("@sha256:"), + ) { + (Some((_, a)), Some((_, b))) => a.len() == 64 && a.eq_ignore_ascii_case(b), + _ => false, + } +} + +async fn save_at(root: &Path, directory: &str, record: &StagedUpdate) -> Result<()> { + use std::{ + io::Write, + os::unix::fs::{DirBuilderExt, OpenOptionsExt, PermissionsExt}, + }; + record.image()?; + let target = record_path(root, directory, &record.manifest.app.id)?; + record + .manifest + .validate() + .context("Invalid staged manifest")?; + let payload = serde_json::to_string(record)?; + let checksum = hex::encode(Sha256::digest(payload.as_bytes())); + let bytes = serde_json::to_vec(&Envelope { + schema: 1, + payload, + checksum, + })?; + // Bounded metadata commit stays in this poll while the caller owns its + // lifecycle lock; cancellation cannot leave a late detached writer. + (|| -> Result<()> { + let parent = target.parent().context("Missing staged update directory")?; + match std::fs::DirBuilder::new().mode(0o700).create(parent) { + Ok(()) => {} + Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {} + Err(e) => return Err(e.into()), + } + std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700))?; + // Persist the directory entry itself before an acknowledged record can + // depend on it surviving a crash (syncing the child alone is insufficient). + std::fs::File::open(parent.parent().context("Missing staged update parent")?)? + .sync_all()?; + let tmp = parent.join(format!(".{}.tmp", uuid::Uuid::new_v4())); + let result = (|| -> Result<()> { + let mut file = std::fs::OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&tmp)?; + file.write_all(&bytes)?; + file.sync_all()?; + std::fs::rename(&tmp, &target)?; + std::fs::File::open(parent)?.sync_all()?; + Ok(()) + })(); + if result.is_err() { + let _ = std::fs::remove_file(tmp); + } + result + })() +} + +pub(super) async fn clear(root: &Path, app: &str) -> Result<()> { + clear_at(root, "staged-updates", app).await +} + +pub(super) async fn clear_installed(root: &Path, app: &str) -> Result<()> { + clear_at(root, "installed-manifests", app).await +} + +async fn clear_at(root: &Path, directory: &str, app: &str) -> Result<()> { + let target = record_path(root, directory, app)?; + // Bounded metadata commit stays in this poll while the caller owns its + // lifecycle lock; cancellation cannot leave a late detached writer. + (|| -> Result<()> { + match std::fs::remove_file(&target) { + Ok(()) => std::fs::File::open(target.parent().unwrap())?.sync_all()?, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => {} + Err(e) => return Err(e.into()), + } + Ok(()) + })() +} + +pub(super) async fn marked(root: &Path, file: &str, app: &str, container: &str) -> Result { + let bytes = match tokio::fs::read(root.join(file)).await { + Ok(bytes) => bytes, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(false), + Err(e) => return Err(e).with_context(|| format!("Cannot read {file}")), + }; + let set: HashSet = + serde_json::from_slice(&bytes).with_context(|| format!("Damaged {file}"))?; + Ok(set.contains(app) || set.contains(container) || set.contains(&format!("archy-{app}"))) +} + +#[cfg(test)] +mod tests { + use super::*; + use futures_util::FutureExt; + use std::os::unix::fs::PermissionsExt; + + #[tokio::test] + async fn envelope_rejects_corruption_schema_and_invalid_manifest_and_is_private() { + let root = tempfile::tempdir().unwrap(); + let manifest = AppManifest::parse(&format!("app:\n id: music\n name: Music\n version: 1.0.0\n container:\n image: docker.io/library/alpine@sha256:{}\n", "a".repeat(64))).unwrap(); + let record = StagedUpdate { + manifest, + manifest_dir: root.path().into(), + ready: true, + }; + save(root.path(), &record) + .now_or_never() + .expect("Journal commit must finish before lifecycle lock cancellation is possible") + .unwrap(); + let target = path(root.path(), "music").unwrap(); + assert_eq!( + std::fs::metadata(target.parent().unwrap()) + .unwrap() + .permissions() + .mode() + & 0o777, + 0o700 + ); + assert_eq!( + std::fs::metadata(&target).unwrap().permissions().mode() & 0o777, + 0o600 + ); + assert!(load(root.path(), "music").await.unwrap().unwrap().ready); + save_installed(root.path(), &record).await.unwrap(); + clear(root.path(), "music") + .now_or_never() + .expect("Journal clear must not leave a detached late deletion") + .unwrap(); + assert!(load(root.path(), "music").await.unwrap().is_none()); + assert!( + installed_manifest(root.path(), "music", record.image().unwrap()) + .await + .unwrap() + .is_some() + ); + assert!(installed_manifest( + root.path(), + "music", + &format!("docker.io/library/alpine@sha256:{}", "b".repeat(64)) + ) + .await + .unwrap() + .is_none()); + save(root.path(), &record).await.unwrap(); + let original = tokio::fs::read(&target).await.unwrap(); + let mut envelope: Envelope = serde_json::from_slice(&original).unwrap(); + envelope.payload = envelope.payload.replace("Music", "Changed Music"); + tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap()) + .await + .unwrap(); + assert!(load(root.path(), "music") + .await + .unwrap_err() + .to_string() + .contains("checksum")); + envelope = serde_json::from_slice(&original).unwrap(); + envelope.schema = 2; + tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap()) + .await + .unwrap(); + assert!(load(root.path(), "music").await.is_err()); + envelope = serde_json::from_slice(&original).unwrap(); + let mut value: serde_json::Value = serde_json::from_str(&envelope.payload).unwrap(); + value["manifest"]["app"]["container"]["image"] = serde_json::Value::Null; + envelope.payload = serde_json::to_string(&value).unwrap(); + envelope.checksum = hex::encode(Sha256::digest(envelope.payload.as_bytes())); + tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap()) + .await + .unwrap(); + assert!(load(root.path(), "music").await.is_err()); + } +} + +pub(super) async fn pending_ids(root: &Path) -> Result> { + let mut dir = match tokio::fs::read_dir(root.join("staged-updates")).await { + Ok(dir) => dir, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(vec![]), + Err(e) => return Err(e.into()), + }; + let mut ids = Vec::new(); + while let Some(entry) = dir.next_entry().await? { + if let Some(name) = entry + .file_name() + .to_str() + .and_then(|name| name.strip_suffix(".json")) + { + ids.push(name.to_string()); + } + } + Ok(ids) +} diff --git a/core/archipelago/src/container/traits.rs b/core/archipelago/src/container/traits.rs index 1a71a0a0..a681f0c1 100644 --- a/core/archipelago/src/container/traits.rs +++ b/core/archipelago/src/container/traits.rs @@ -62,6 +62,11 @@ pub trait ContainerOrchestrator: Send + Sync { /// Pull/rebuild the image and recreate the container from scratch. async fn upgrade(&self, app_id: &str) -> Result<()>; + /// Exact image prepared by a stopped upgrade; no running container is claimed. + async fn staged_upgrade_image(&self, _app_id: &str) -> Result> { + Ok(None) + } + /// Current state of a single container. async fn status(&self, app_id: &str) -> Result;