From 07c7eb0f1429c13eefedb904a6de4a41e1d3f432 Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 22:35:54 -0400 Subject: [PATCH] Preserve intact originals when pre-target Indee drain refuses --- .../src/container/supervised_update.rs | 502 +++++++++++++++++- 1 file changed, 497 insertions(+), 5 deletions(-) diff --git a/core/archipelago/src/container/supervised_update.rs b/core/archipelago/src/container/supervised_update.rs index a5084e6b..3ae3370a 100644 --- a/core/archipelago/src/container/supervised_update.rs +++ b/core/archipelago/src/container/supervised_update.rs @@ -39,6 +39,10 @@ struct Member { pinned_original_body: String, original_tag: String, recovery_image: Option, + /// A durable pre-target recovery decision. None is not yet decided; false + /// means recreate only a stopped/missing original, true preserves its ID. + #[serde(default)] + preserve_original: Option, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] enum Phase { @@ -423,7 +427,7 @@ fn root(guard: &Guard) -> Result { } fn validate(record: &Journal) -> Result<()> { anyhow::ensure!( - record.schema == 1 + matches!(record.schema, 1 | 2) && simple(&record.package) && uuid::Uuid::parse_str(&record.id)?.to_string() == record.id && !record.members.is_empty() @@ -461,6 +465,13 @@ fn validate(record: &Journal) -> Result<()> { || member.recovery_image.is_some(), "Destructive update lacks a durable writable-layer recovery image" ); + anyhow::ensure!( + member.preserve_original.is_none() + || (record.schema == 2 + && !record.target_startup_began + && matches!(record.phase, Phase::Restoring | Phase::Restored)), + "Preserved-original recovery cannot follow target startup" + ); let restore_image = member .recovery_image .as_ref() @@ -626,7 +637,11 @@ fn publish_installed(guard: &Guard, record: &Journal) -> Result<()> { operation: record.id.clone(), name: member.original.name.clone(), body: if record.phase == Phase::Restored { - member.pinned_original_body.clone() + if member.preserve_original == Some(true) { + member.original.body.clone() + } else { + member.pinned_original_body.clone() + } } else { member.target_body.clone() }, @@ -750,10 +765,11 @@ pub(crate) async fn execute( pinned_original_body, original_tag: format!("localhost/archy-update-recovery:{id}-{index}"), recovery_image: None, + preserve_original: None, }); } let mut record = Journal { - schema: 1, + schema: 2, id, package: package.into(), phase: Phase::Prepared, @@ -868,7 +884,139 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor record.cleanup_done = true; save(guard, record) } +/// Observe without mutation. A live original is never stopped because a drain +/// failed; a live same-operation recovery may be adopted after a lost start ACK. +async fn pre_target_state(member: &Member, supervisor: &impl Supervisor) -> Result<(bool, bool)> { + supervisor.validate_original_file(&member.original).await?; + let body = supervisor.read(&member.original.name).await?; + anyhow::ensure!( + [ + &member.original.body, + &member.target_body, + &member.pinned_original_body + ] + .contains(&&body), + "Foreign unit edit requires explicit recovery" + ); + let observed = supervisor.observed(&member.original.name).await?; + let intact = observed.as_ref().is_some_and(|current| { + body == member.original.body + && current.id == member.original.container_id + && current.image == member.original.image + && current.running == member.original.running + && current.config_sha256 == member.original.config_sha256 + }); + if member.preserve_original == Some(true) { + anyhow::ensure!( + intact, + "Preserved original changed; retain recovery hold for inspection" + ); + return Ok((true, false)); + } + if member.preserve_original.is_none() && intact { + return Ok((true, false)); + } + let recovery = &member + .recovery_image + .as_ref() + .context("Missing recovery image")? + .image; + if let Some(current) = observed { + let own_recovery = current.image == *recovery + && current.config_sha256 == member.original.config_sha256 + && body == member.pinned_original_body; + if current.running { + anyhow::ensure!( + member.original.running, + "Originally stopped member unexpectedly started; preserve stop intent and recovery hold" + ); + anyhow::ensure!( + own_recovery, + "Unexpected live replacement; original writers must not be stopped" + ); + return Ok((false, true)); + } + anyhow::ensure!( + own_recovery + || (current.id == member.original.container_id + && current.image == member.original.image + && current.config_sha256 == member.original.config_sha256), + "Unexpected stopped replacement; retain its data for inspection" + ); + } + Ok((false, false)) +} +async fn restore_before_target( + guard: &Guard, + record: &mut Journal, + supervisor: &impl Supervisor, +) -> Result<()> { + let mut choices = Vec::new(); + // Validate every member before changing even one service recipe. + for member in &record.members { + choices.push(pre_target_state(member, supervisor).await?.0); + } + for (member, preserve) in record.members.iter_mut().zip(choices) { + if let Some(saved) = member.preserve_original { + anyhow::ensure!(saved == preserve, "Original recovery decision changed"); + } + member.preserve_original = Some(preserve); + } + save(guard, record)?; + let mut changed = false; + for member in &record.members { + let (preserve, running_recovery) = pre_target_state(member, supervisor).await?; + if preserve || running_recovery { + continue; + } + supervisor + .pin( + &member + .recovery_image + .as_ref() + .context("Missing recovery image")? + .image, + &member.original_tag, + ) + .await?; + supervisor + .write( + &member.original, + &[ + member.original.body.clone(), + member.target_body.clone(), + member.pinned_original_body.clone(), + ], + &member.pinned_original_body, + ) + .await?; + changed = true; + } + if changed { + supervisor.reload().await?; + } + for member in &record.members { + let (preserve, running_recovery) = pre_target_state(member, supervisor).await?; + if !preserve && !running_recovery && member.original.running { + supervisor.start(&member.original.name).await?; + } + } + // Recheck preserved IDs and bodies as well as recreated runtime before + // publishing terminal installed recipes or reopening admission. + for member in &record.members { + let (preserve, running_recovery) = pre_target_state(member, supervisor).await?; + anyhow::ensure!( + preserve || running_recovery || !member.original.running, + "Original recovery did not return; retain admission hold" + ); + } + Ok(()) +} async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> { + // Older readers must refuse journals whose per-member preservation choices + // they cannot honor. Accept old journals for migration, then upgrade before + // persisting any recovery transition or performing recovery actions. + record.schema = 2; 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. @@ -908,6 +1056,10 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis .begin_barrier(&record.id, &originals, true) .await?; supervisor.verify_barrier(&record.id).await?; + if !record.target_startup_began { + restore_before_target(guard, record, supervisor).await?; + return complete_restore(guard, record, supervisor).await; + } // Refuse to overwrite a foreign edit before stopping any surviving member. for member in &record.members { supervisor.validate_original_file(&member.original).await?; @@ -978,6 +1130,13 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis ); } } + complete_restore(guard, record, supervisor).await +} +async fn complete_restore( + guard: &Guard, + record: &mut Journal, + supervisor: &impl Supervisor, +) -> Result<()> { for member in &record.members { guard.hold(&member.original.name, &record.id)?; } @@ -1014,6 +1173,16 @@ pub(crate) async fn recover(guard: &Guard, supervisor: &impl Supervisor) -> Resu } } Phase::Restored => { + if !record.target_startup_began { + for member in &record.members { + let (preserve, running_recovery) = + pre_target_state(member, supervisor).await?; + anyhow::ensure!( + preserve || running_recovery || !member.original.running, + "Restored original changed before admission release" + ); + } + } publish_installed(guard, &record)?; supervisor .release_barrier(&record.id, Completion::Restored) @@ -1124,7 +1293,8 @@ mod tests { "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( + let body = pin_body(&self.original.body, &self.original.name, &target.reference)? + .replace( "Pull=never", "Environment=NEW_FEATURE=enabled\nVolume=/identity:/run/identity:ro\nPull=never", ); @@ -1201,7 +1371,7 @@ mod tests { let new = self.body.lock().unwrap().contains("NEW_FEATURE=enabled"); Ok(Some(Observed { id: format!("{:064x}", self.generation.load(Ordering::SeqCst)), - name: "movie".into(), + name: self.original.name.clone(), image: if new { "b".repeat(64) } else if self.body.lock().unwrap().contains(&"e".repeat(64)) { @@ -1218,6 +1388,328 @@ mod tests { Ok(true) } } + struct DrainFailureStack { + members: std::collections::BTreeMap, + partial_frontend: bool, + lose_frontend_start: AtomicBool, + } + impl DrainFailureStack { + fn new(partial_frontend: bool) -> Self { + let members = ["frontend", "worker", "api", "storage"] + .into_iter() + .enumerate() + .map(|(index, name)| { + let mut member = Mock::new(); + member.original.name = name.into(); + member.original.body = member + .original + .body + .replace("ContainerName=movie", &format!("ContainerName={name}")); + *member.body.lock().unwrap() = member.original.body.clone(); + member.original.container_id = format!("{:064x}", index + 1); + member.generation.store(index + 1, Ordering::SeqCst); + (name.into(), member) + }) + .collect(); + Self { + members, + partial_frontend, + lose_frontend_start: AtomicBool::new(false), + } + } + fn targets(&self) -> Vec { + self.members + .keys() + .map(|name| { + let mut target = Mock::target(); + target.name = name.clone(); + target + }) + .collect() + } + } + impl Supervisor for DrainFailureStack { + async fn begin_barrier( + &self, + operation: &str, + originals: &[Unit], + recovery: bool, + ) -> Result<()> { + for member in self.members.values() { + member.begin_barrier(operation, originals, recovery).await?; + } + if !recovery { + if self.partial_frontend { + self.members["frontend"].stop("frontend").await?; + } + anyhow::bail!("Active worker refused drain; no work completion acknowledged"); + } + Ok(()) + } + async fn verify_barrier(&self, operation: &str) -> Result<()> { + for member in self.members.values() { + member.verify_barrier(operation).await?; + } + Ok(()) + } + async fn release_barrier(&self, operation: &str, outcome: Completion) -> Result<()> { + for member in self.members.values() { + member.release_barrier(operation, outcome).await?; + } + Ok(()) + } + async fn prepare_target(&self, target: &Target, original: &Unit) -> Result { + self.members[&target.name] + .prepare_target(target, original) + .await + } + async fn target_hooks( + &self, + name: &str, + manifest: &archipelago_container::AppManifest, + ) -> Result<()> { + self.members[name].target_hooks(name, manifest).await + } + async fn snapshot( + &self, + original: &Unit, + operation: &str, + tag: &str, + ) -> Result { + self.members[&original.name] + .snapshot(original, operation, tag) + .await + } + async fn capture(&self, name: &str) -> Result { + self.members[name].capture(name).await + } + async fn validate_original_file(&self, original: &Unit) -> Result<()> { + self.members[&original.name] + .validate_original_file(original) + .await + } + async fn read(&self, name: &str) -> Result { + self.members[name].read(name).await + } + async fn write(&self, original: &Unit, expected: &[String], body: &str) -> Result<()> { + self.members[&original.name] + .write(original, expected, body) + .await + } + async fn pin(&self, _image: &str, _tag: &str) -> Result<()> { + Ok(()) + } + async fn stop(&self, name: &str) -> Result<()> { + self.members[name].stop(name).await + } + async fn reload(&self) -> Result<()> { + Ok(()) + } + async fn start(&self, name: &str) -> Result<()> { + self.members[name].start(name).await?; + anyhow::ensure!( + name != "frontend" || !self.lose_frontend_start.swap(false, Ordering::SeqCst), + "Lost original start acknowledgement" + ); + Ok(()) + } + async fn observed(&self, name: &str) -> Result> { + self.members[name].observed(name).await + } + async fn healthy(&self, name: &str) -> Result { + self.members[name].healthy(name).await + } + } + #[tokio::test] + async fn failed_active_worker_drain_never_stops_or_recreates_intact_originals() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(false); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + let record = records(&guard).unwrap().pop().unwrap(); + assert_eq!(record.phase, Phase::Restored); + assert!(!record.target_startup_began); + assert_eq!(record.schema, 2); + let mut downgraded = record.clone(); + downgraded.schema = 1; + assert!(validate(&downgraded).is_err()); + for member in &record.members { + assert_eq!(member.preserve_original, Some(true)); + let live = &runtime.members[&member.original.name]; + assert_eq!(*live.calls.lock().unwrap(), ["snapshot-original"]); + assert_eq!( + live.observed(&member.original.name) + .await + .unwrap() + .unwrap() + .id, + member.original.container_id + ); + assert_eq!( + installed_unit(root.path(), &member.original.name) + .unwrap() + .unwrap() + .0, + member.original.body + ); + } + } + #[tokio::test] + async fn partial_frontend_drain_restores_only_missing_frontend_preserving_busy_writers() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(true); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + let record = records(&guard).unwrap().pop().unwrap(); + assert_eq!(record.phase, Phase::Restored); + for member in &record.members { + let live = &runtime.members[&member.original.name]; + if member.original.name == "frontend" { + assert_eq!(member.preserve_original, Some(false)); + assert_eq!( + *live.calls.lock().unwrap(), + ["snapshot-original", "stop-unit", "write", "start-unit"] + ); + assert_eq!( + installed_unit(root.path(), "frontend").unwrap().unwrap().0, + member.pinned_original_body + ); + } else { + assert_eq!(member.preserve_original, Some(true)); + assert_eq!(*live.calls.lock().unwrap(), ["snapshot-original"]); + assert_eq!( + installed_unit(root.path(), &member.original.name) + .unwrap() + .unwrap() + .0, + member.original.body + ); + } + } + } + #[tokio::test] + async fn pre_target_restart_adopts_own_recreation_after_lost_start_ack_without_stopping_it() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(true); + runtime.lose_frontend_start.store(true, Ordering::SeqCst); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + assert_eq!(records(&guard).unwrap()[0].phase, Phase::Restoring); + let before = runtime.members["frontend"] + .observed("frontend") + .await + .unwrap() + .unwrap() + .id; + for member in runtime.members.values() { + member.calls.lock().unwrap().clear(); + } + recover(&guard, &runtime).await.unwrap(); + assert_eq!( + runtime.members["frontend"] + .observed("frontend") + .await + .unwrap() + .unwrap() + .id, + before + ); + for member in runtime.members.values() { + assert!(member.calls.lock().unwrap().is_empty()); + } + assert_eq!(records(&guard).unwrap()[0].phase, Phase::Restored); + } + #[tokio::test] + async fn pre_target_recovery_rejects_running_recreation_of_originally_stopped_member() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(true); + runtime.lose_frontend_start.store(true, Ordering::SeqCst); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + let mut record = records(&guard).unwrap().pop().unwrap(); + // Model a restart journal whose captured operator intent was stopped, + // while an external actor has started its operation-owned recovery. + record + .members + .iter_mut() + .find(|m| m.original.name == "frontend") + .unwrap() + .original + .running = false; + save(&guard, &record).unwrap(); + for member in runtime.members.values() { + member.calls.lock().unwrap().clear(); + } + let error = recover(&guard, &runtime).await.unwrap_err(); + assert!(error.to_string().contains("Originally stopped member")); + assert_eq!(records(&guard).unwrap()[0].phase, Phase::Restoring); + assert!(super::super::update_transaction::is_held(root.path(), "frontend").unwrap()); + for member in runtime.members.values() { + assert!(member.calls.lock().unwrap().is_empty()); + } + } + #[tokio::test] + async fn legacy_pre_target_journal_upgrades_before_preserving_surviving_originals() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(true); + runtime.lose_frontend_start.store(true, Ordering::SeqCst); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + let mut legacy = records(&guard).unwrap().pop().unwrap(); + legacy.schema = 1; + for member in &mut legacy.members { + member.preserve_original = None; + } + save(&guard, &legacy).unwrap(); + for member in runtime.members.values() { + member.calls.lock().unwrap().clear(); + } + recover(&guard, &runtime).await.unwrap(); + let migrated = records(&guard).unwrap().pop().unwrap(); + assert_eq!(migrated.schema, 2); + assert_eq!(migrated.phase, Phase::Restored); + assert!(migrated + .members + .iter() + .all(|m| m.preserve_original.is_some())); + for member in runtime.members.values() { + assert!(member.calls.lock().unwrap().is_empty()); + } + } + #[tokio::test] + async fn changed_preserved_writer_blocks_recovery_and_cannot_be_reclassified_for_recreation() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = DrainFailureStack::new(true); + runtime.lose_frontend_start.store(true, Ordering::SeqCst); + assert!(execute(&guard, "movie", &runtime.targets(), &runtime) + .await + .is_err()); + runtime.members["worker"] + .generation + .store(99, Ordering::SeqCst); + for member in runtime.members.values() { + member.calls.lock().unwrap().clear(); + } + assert!(recover(&guard, &runtime).await.is_err()); + for member in runtime.members.values() { + assert!(member.calls.lock().unwrap().is_empty()); + } + assert!(super::super::update_transaction::is_held(root.path(), "worker").unwrap()); + let mut record = records(&guard).unwrap().pop().unwrap(); + record.target_startup_began = true; + assert!(save(&guard, &record).is_err()); + } #[test] fn administrative_original_capture_is_idempotent_and_uninstall_forgets_it() { let root = tempfile::tempdir().unwrap();