Preserve intact originals when pre-target Indee drain refuses

This commit is contained in:
archipelago
2026-10-07 22:35:54 -04:00
parent 72df1d17aa
commit 07c7eb0f14
@@ -39,6 +39,10 @@ struct Member {
pinned_original_body: String,
original_tag: String,
recovery_image: Option<RecoveryImage>,
/// 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<bool>,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
enum Phase {
@@ -423,7 +427,7 @@ fn root(guard: &Guard) -> Result<PathBuf> {
}
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<String, Mock>,
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<Target> {
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<PreparedTarget> {
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<RecoveryImage> {
self.members[&original.name]
.snapshot(original, operation, tag)
.await
}
async fn capture(&self, name: &str) -> Result<Unit> {
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<String> {
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<Option<Observed>> {
self.members[name].observed(name).await
}
async fn healthy(&self, name: &str) -> Result<bool> {
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();