From 6c030a8109fa4fe79e27287d9992002b4037912a Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 01:52:36 -0400 Subject: [PATCH] Plan reviewed Quadlet migrations and preserve private writable-layer snapshots --- .../src/container/supervised_update.rs | 260 ++++++++++++++++++ 1 file changed, 260 insertions(+) diff --git a/core/archipelago/src/container/supervised_update.rs b/core/archipelago/src/container/supervised_update.rs index 37fe6a94..0f4d736b 100644 --- a/core/archipelago/src/container/supervised_update.rs +++ b/core/archipelago/src/container/supervised_update.rs @@ -97,6 +97,114 @@ pub(crate) trait Supervisor: Sync { fn observed(&self, name: &str) -> impl Future>> + Send; fn healthy(&self, name: &str) -> impl Future> + Send; } +/// Local recovery images are never pushed or exported. A matching tag may be +/// reused after a lost commit response only when image labels bind the exact +/// original container and transaction. Mounted data is deliberately excluded. +pub(crate) async fn capture_local_recovery_image( + original: &Unit, + operation: &str, + tag: &str, +) -> Result { + use super::update_transaction::{Podman, Runtime}; + anyhow::ensure!( + digest(&original.container_id) + && uuid::Uuid::parse_str(operation)?.to_string() == operation + && tag.starts_with(&format!("localhost/archy-update-recovery:{operation}-")) + && tag + .rsplit('-') + .next() + .is_some_and(|v| v.parse::().is_ok()), + "Invalid local recovery image ownership" + ); + async fn command(args: &[&str]) -> Result { + tokio::time::timeout( + std::time::Duration::from_secs(300), + tokio::process::Command::new("podman") + .args(args) + .kill_on_drop(true) + .output(), + ) + .await + .context("Recovery image operation timed out; original runtime retained")? + .context("Recovery image runtime unavailable") + } + async fn inspect(tag: &str, original: &Unit, operation: &str) -> Result { + let result = command(&["image", "inspect", tag]).await?; + anyhow::ensure!( + result.status.success(), + "Local recovery image inspection failed" + ); + let rows: Vec = serde_json::from_slice(&result.stdout)?; + anyhow::ensure!(rows.len() == 1, "Ambiguous recovery image"); + let row = &rows[0]; + let labels = row + .pointer("/Config/Labels") + .context("Recovery image has no ownership labels")?; + anyhow::ensure!( + labels + .get("io.archipelago.recovery.operation") + .and_then(|v| v.as_str()) + == Some(operation) + && labels + .get("io.archipelago.recovery.container") + .and_then(|v| v.as_str()) + == Some(original.container_id.as_str()), + "Existing recovery image belongs to another operation" + ); + let image = row + .get("Id") + .or_else(|| row.get("ID")) + .and_then(|v| v.as_str()) + .context("Recovery image identity missing")? + .trim_start_matches("sha256:") + .to_string(); + anyhow::ensure!(digest(&image), "Invalid recovery image digest"); + Ok(RecoveryImage { + image, + source_container_id: original.container_id.clone(), + operation_id: operation.into(), + }) + } + let current = Podman + .inspect(&original.name) + .await? + .context("Original container missing before snapshot")?; + anyhow::ensure!( + current.id == original.container_id + && current.image == original.image + && current.running == original.running + && current.config_sha256 == original.config_sha256, + "Original runtime changed before writable-layer snapshot" + ); + let exists = command(&["image", "exists", tag]).await?; + match exists.status.code() { + Some(0) => return inspect(tag, original, operation).await, + Some(1) => {} + _ => anyhow::bail!("Recovery image inventory unavailable"), + } + let operation_label = format!("LABEL io.archipelago.recovery.operation={operation}"); + let container_label = format!( + "LABEL io.archipelago.recovery.container={}", + original.container_id + ); + let result = command(&[ + "commit", + "--pause=true", + "--include-volumes=false", + "--change", + &operation_label, + "--change", + &container_label, + &original.container_id, + tag, + ]) + .await?; + anyhow::ensure!( + result.status.success(), + "Writable-layer snapshot failed; original runtime retained" + ); + inspect(tag, original, operation).await +} fn simple(value: &str) -> bool { !value.is_empty() && value.len() <= 128 @@ -151,6 +259,122 @@ pub(crate) fn pin_body(body: &str, name: &str, image: &str) -> Result { ); Ok(output) } +/// Three-way configuration migration: the previous renderer output must come +/// from the installed immutable manifest, never today's mutable catalog. Keep +/// unrelated operator overrides; a conflicting required change is a preflight +/// error rather than an overwrite. Values are never included in errors. +pub(crate) fn merge_reviewed_unit(original: &str, previous: &str, next: &str) -> Result { + type Key = (String, String); + fn parse(body: &str) -> Result<(Vec, std::collections::BTreeMap>)> { + anyhow::ensure!(body.len() <= 1024 * 1024, "Unit exceeds migration limit"); + let mut section = String::new(); + let mut sections = std::collections::HashSet::new(); + let mut order = Vec::new(); + let mut values = std::collections::BTreeMap::>::new(); + for raw in body.lines() { + let line = raw.trim(); + if line.is_empty() || line.starts_with('#') || line.starts_with(';') { + continue; + } + anyhow::ensure!( + !line.ends_with('\\') && !line.starts_with(".include"), + "Continued or included unit cannot be migrated automatically" + ); + if line.starts_with('[') { + anyhow::ensure!( + line.ends_with(']') && sections.insert(line.to_string()), + "Repeated or malformed unit section" + ); + section = line.into(); + continue; + } + let (directive, value) = line.split_once('=').context("Malformed unit directive")?; + anyhow::ensure!( + !section.is_empty() + && !value.is_empty() + && directive.bytes().all(|b| b.is_ascii_alphanumeric()) + && !matches!(directive, "EnvironmentFile" | "EnvFile"), + "Unsupported unit reset or external configuration" + ); + let key = if directive == "Environment" { + let env = value.strip_prefix('"').unwrap_or(value); + let (name, _) = env + .split_once('=') + .context("Unsupported environment directive")?; + anyhow::ensure!( + !name.is_empty() + && name.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'_') + && (value.starts_with('"') && value.ends_with('"') + || !value.contains(char::is_whitespace)), + "Ambiguous environment directive" + ); + format!("Environment:{name}") + } else { + directive.into() + }; + let key = (section.clone(), key); + if !values.contains_key(&key) { + order.push(key.clone()); + } + values.entry(key).or_default().push(raw.to_string()); + } + Ok((order, values)) + } + let (mut order, original) = parse(original)?; + let (_, previous) = parse(previous)?; + let (next_order, next) = parse(next)?; + let mut merged = original.clone(); + let keys: std::collections::BTreeSet<_> = previous.keys().chain(next.keys()).cloned().collect(); + for key in keys { + let old = previous.get(&key); + let wanted = next.get(&key); + if old == wanted { + continue; + } + let actual = original.get(&key); + anyhow::ensure!( + actual == old || actual == wanted, + "Required unit configuration conflicts with an operator override in {} {}", + key.0, + key.1 + ); + match wanted { + Some(lines) => { + merged.insert(key, lines.clone()); + } + None => { + merged.remove(&key); + } + } + } + for key in next_order { + if !order.contains(&key) { + order.push(key); + } + } + // Group sections once; preserve directive and repeated-value order within + // each section, including operator-only directives absent from both plans. + let mut sections = Vec::::new(); + for (section, _) in &order { + if !sections.contains(section) { + sections.push(section.clone()); + } + } + let mut output = String::new(); + for section in sections { + output.push_str(§ion); + output.push('\n'); + for key in order.iter().filter(|key| key.0 == section) { + if let Some(lines) = merged.get(key) { + for line in lines { + output.push_str(line); + output.push('\n'); + } + } + } + } + Ok(output) +} fn root(guard: &Guard) -> Result { let path = guard.directory().join("supervised"); match std::fs::DirBuilder::new().mode(0o700).create(&path) { @@ -672,6 +896,42 @@ mod tests { Ok(true) } } + #[test] + fn reviewed_migration_keeps_unrelated_operator_values_and_applies_required_new_fields() { + let old = "[Container]\nImage=old\nEnvironment=FEATURE=old\nEnvironment=PORT=1\n[Service]\nRestart=always\n"; + let actual = old + .replace("PORT=1", "PORT=42") + .replace("[Service]", "Environment=OPERATOR=mine\n[Service]"); + let next = old + .replace("Image=old", "Image=new") + .replace("FEATURE=old", "FEATURE=new") + .replace("[Service]", "Volume=/identity:/run/identity:ro\n[Service]"); + let merged = merge_reviewed_unit(&actual, old, &next).unwrap(); + assert!(merged.contains("Image=new")); + assert!(merged.contains("FEATURE=new")); + assert!(merged.contains("PORT=42")); + assert!(merged.contains("OPERATOR=mine")); + assert!(merged.contains("Volume=/identity:/run/identity:ro")); + assert_eq!(merge_reviewed_unit(&merged, old, &next).unwrap(), merged); + } + #[test] + fn conflicting_required_environment_or_mount_migration_is_rejected_before_lifecycle() { + let old = "[Container]\nEnvironment=FEATURE=old\nVolume=/old:/data\n"; + assert!(merge_reviewed_unit( + &old.replace("FEATURE=old", "FEATURE=operator"), + old, + &old.replace("FEATURE=old", "FEATURE=new") + ) + .is_err()); + assert!(merge_reviewed_unit( + &old.replace("/old:/data", "/operator:/data"), + old, + &old.replace("/old:/data", "/required:/data") + ) + .is_err()); + assert!(merge_reviewed_unit("[Container]\nEnvironment=\n", old, old).is_err()); + assert!(merge_reviewed_unit("[Container]\nEnvironment=A=1 B=2\n", old, old).is_err()); + } #[tokio::test] async fn snapshot_failure_never_stops_or_recreates_original_runtime() { let root = tempfile::tempdir().unwrap();