Draft retained-container update journal and original-runtime recovery

This commit is contained in:
archipelago
2026-10-07 01:24:50 -04:00
parent 1bbf85e0d4
commit 45579e53c7
6 changed files with 944 additions and 140 deletions
@@ -280,6 +280,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
let docker_image = params let docker_image = params
.get("dockerImage") .get("dockerImage")
@@ -1901,7 +1904,10 @@ autopilot.active=false\n",
.await .await
.context("DATUM credentials are not available yet; wait for installation to finish")?; .context("DATUM credentials are not available yet; wait for installation to finish")?;
let password = password.trim(); let password = password.trim();
anyhow::ensure!(!password.is_empty(), "DATUM administrator password is empty"); anyhow::ensure!(
!password.is_empty(),
"DATUM administrator password is empty"
);
return Ok(serde_json::json!({ return Ok(serde_json::json!({
"title": "DATUM Gateway login", "title": "DATUM Gateway login",
"description": "Use this password when DATUM asks you to unlock configuration. In Config, set your Bitcoin payout address. Point miners at this node's IP address on Stratum port 23334 (stratum+tcp://NODE-IP:23334). Gashboard connects automatically.", "description": "Use this password when DATUM asks you to unlock configuration. In Config, set your Bitcoin payout address. Point miners at this node's IP address on Stratum port 23334 (stratum+tcp://NODE-IP:23334). Gashboard connects automatically.",
@@ -60,6 +60,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
// A cuprate node that starts on a too-small disk fills it and takes // A cuprate node that starts on a too-small disk fills it and takes
// Archipelago down with it (no upstream pruning — see // Archipelago down with it (no upstream pruning — see
// dependencies::check_cuprate_disk_compatibility). Fail the start // dependencies::check_cuprate_disk_compatibility). Fail the start
@@ -104,6 +107,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id); let op_lock = app_op_lock(package_id);
let data_dir = self.config.data_dir.clone(); let data_dir = self.config.data_dir.clone();
tokio::spawn(async move { tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await; let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() { let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_start(orchestrator.as_ref(), &to_start).await do_orchestrator_package_start(orchestrator.as_ref(), &to_start).await
@@ -167,6 +171,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
let single_orchestrator_app = let single_orchestrator_app =
self.orchestrator.is_some() && uses_single_orchestrator_app(package_id); self.orchestrator.is_some() && uses_single_orchestrator_app(package_id);
@@ -231,6 +238,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id); let op_lock = app_op_lock(package_id);
tokio::spawn(async move { tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await; let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() { let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_stop(orchestrator.as_ref(), &to_stop).await do_orchestrator_package_stop(orchestrator.as_ref(), &to_stop).await
@@ -269,6 +277,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
// Restart is stop + recreate, so on a disk that shrank below the cuprate // Restart is stop + recreate, so on a disk that shrank below the cuprate
// minimum after install it resumes the doomed unprunable sync just like // minimum after install it resumes the doomed unprunable sync just like
// start would — same gate, same "fail before clearing user-stopped / // start would — same gate, same "fail before clearing user-stopped /
@@ -331,6 +342,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id); let op_lock = app_op_lock(package_id);
let data_dir = self.config.data_dir.clone(); let data_dir = self.config.data_dir.clone();
tokio::spawn(async move { tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await; let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() { let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_restart(orchestrator.as_ref(), &to_restart).await do_orchestrator_package_restart(orchestrator.as_ref(), &to_restart).await
@@ -374,6 +386,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
let preserve_data = params let preserve_data = params
.get("preserve_data") .get("preserve_data")
.and_then(|v| v.as_bool()) .and_then(|v| v.as_bool())
+24 -139
View File
@@ -7,7 +7,6 @@
use super::config::{all_container_names, get_containers_for_app}; use super::config::{all_container_names, get_containers_for_app};
use super::install::install_log; use super::install::install_log;
use super::progress::parse_pull_progress; use super::progress::parse_pull_progress;
use super::runtime::stop_timeout_secs;
use super::validation::validate_app_id; use super::validation::validate_app_id;
use crate::api::rpc::RpcHandler; use crate::api::rpc::RpcHandler;
use crate::container::image_versions; use crate::container::image_versions;
@@ -32,6 +31,9 @@ impl RpcHandler {
.and_then(|v| v.as_str()) .and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?; .ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?; validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
// An Update click must not act on an hourly cache that predates the // An Update click must not act on an hourly cache that predates the
// button. Fetch and verify first; failure leaves running containers alone. // button. Fetch and verify first; failure leaves running containers alone.
@@ -217,7 +219,7 @@ impl RpcHandler {
match preflighted_stack_update( match preflighted_stack_update(
&images_to_pull, &images_to_pull,
|image| async move { self.pull_update_image(package_id, &image).await }, |image| async move { self.pull_update_image(package_id, &image).await },
|| self.execute_update(package_id, &containers, &images_to_pull), || self.execute_update(package_id, &containers, &images_to_pull, &lifecycle_guard),
) )
.await .await
{ {
@@ -241,10 +243,14 @@ impl RpcHandler {
package_id, e package_id, e
)) ))
.await; .await;
self.rollback_update(package_id, &containers).await; // Transaction executor already recovered original identities or
// returned an explicit unresolved state. Never guess/reinstall.
self.clear_install_progress(package_id).await; self.clear_install_progress(package_id).await;
self.clear_update_state(package_id).await; self.clear_update_state(package_id).await;
Err(e.context(format!("Update {} failed, rolled back", package_id))) Err(e.context(format!(
"Update {} failed; see retained-runtime recovery result",
package_id
)))
} }
} }
} }
@@ -289,114 +295,28 @@ impl RpcHandler {
} }
} }
/// Images are prepared first; then stop → remove → recreate → verify. /// Images are prepared first. The transaction preserves original identities,
/// creates stopped replacements, and restores the exact old states on failure.
async fn execute_update( async fn execute_update(
&self, &self,
package_id: &str, package_id: &str,
containers: &[String], containers: &[String],
images_to_pull: &[(String, String)], images_to_pull: &[(String, String)],
guard: &crate::container::update_transaction::Guard,
) -> Result<()> { ) -> Result<()> {
// Phase: Preparing — about to stop the running container(s) so use crate::container::update_transaction::{self, Podman};
// we can swap images. Fast. let targets = Podman::targets(images_to_pull).await?;
let names: HashSet<_> = containers.iter().map(String::as_str).collect();
anyhow::ensure!(
targets.len() == names.len()
&& targets
.iter()
.all(|target| names.contains(target.name.as_str())),
"Update target membership differs from installed stack; no containers changed"
);
self.set_install_phase(package_id, InstallPhase::Preparing) self.set_install_phase(package_id, InstallPhase::Preparing)
.await; .await;
update_transaction::execute(guard, package_id, &targets, &Podman).await
// 1. Graceful stop all containers (reverse order for dependencies)
info!(
"Update {}: stopping {} containers",
package_id,
containers.len()
);
for name in containers.iter().rev() {
let timeout = stop_timeout_secs(name);
info!(
"Update {}: stopping {} (timeout: {}s)",
package_id, name, timeout
);
let out = tokio::process::Command::new("podman")
.args(["stop", "-t", timeout, name])
.output()
.await
.context(format!("Failed to stop {}", name))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
warn!(
"Update {}: stop {} failed: {}",
package_id,
name,
stderr.trim()
);
// Continue — container might already be stopped
}
}
// 3. Remove old containers
info!("Update {}: removing old containers", package_id);
for name in containers {
let out = tokio::process::Command::new("podman")
.args(["rm", name])
.output()
.await
.context(format!("Failed to remove {}", name))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
// Force remove as fallback
warn!(
"Update {}: rm {} failed ({}), forcing",
package_id,
name,
stderr.trim()
);
let _ = tokio::process::Command::new("podman")
.args(["rm", "-f", name])
.output()
.await;
}
}
// Phase: CreatingContainer — about to recreate each container.
self.set_install_phase(package_id, InstallPhase::CreatingContainer)
.await;
// 4. Recreate containers (orchestrator-first, reconcile fallback)
info!("Update {}: recreating containers", package_id);
for name in containers {
self.recreate_container_for_update(package_id, name).await?;
// Brief delay between containers for dependency initialization
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
}
// Phase: WaitingHealthy — reconcile has started every container,
// now verifying each reached running state.
self.set_install_phase(package_id, InstallPhase::WaitingHealthy)
.await;
// 5. Verify containers reached running state
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
for name in containers {
let status = tokio::process::Command::new("podman")
.args(["inspect", name, "--format", "{{.State.Status}}"])
.output()
.await;
if let Ok(o) = status {
let state = String::from_utf8_lossy(&o.stdout).trim().to_string();
anyhow::ensure!(
o.status.success() && state == "running",
"Update {}: container {} is not running after recreate",
package_id,
name
);
} else {
anyhow::bail!(
"Update {}: cannot inspect recreated container {}",
package_id,
name
);
}
}
verify_update_targets(images_to_pull, &inspect_update_images(package_id).await?)?;
Ok(())
} }
async fn recreate_container_for_update( async fn recreate_container_for_update(
@@ -555,41 +475,6 @@ impl RpcHandler {
stack_images stack_images
} }
/// Rollback: restart old containers if they still exist.
/// Called when update fails partway through.
async fn rollback_update(&self, package_id: &str, containers: &[String]) {
warn!("Rolling back update for {}", package_id);
for name in containers {
// Try to start — works if container still exists (wasn't removed yet)
let out = tokio::process::Command::new("podman")
.args(["start", name])
.output()
.await;
match out {
Ok(o) if o.status.success() => {
info!("Rollback: restarted {}", name);
}
Ok(o) => {
let stderr = String::from_utf8_lossy(&o.stderr);
warn!("Rollback: could not restart {}: {}", name, stderr.trim());
// Container was already removed (forward path ran `podman rm`).
// Recreate via orchestrator-first path with legacy fallback.
if let Err(recreate_err) =
self.recreate_container_for_update(package_id, name).await
{
error!(
"Rollback: failed to recreate {} during rollback of {}: {}",
name, package_id, recreate_err
);
}
}
Err(e) => {
error!("Rollback: failed to restart {}: {}", name, e);
}
}
}
}
/// Clear the Updating state (used on failure/rollback). /// Clear the Updating state (used on failure/rollback).
async fn clear_update_state(&self, package_id: &str) { async fn clear_update_state(&self, package_id: &str) {
let (mut data, _) = self.state_manager.get_snapshot().await; let (mut data, _) = self.state_manager.get_snapshot().await;
+2
View File
@@ -31,3 +31,5 @@ pub use prod_orchestrator::ProdContainerOrchestrator;
pub use traits::ContainerOrchestrator; pub use traits::ContainerOrchestrator;
mod staged_update; mod staged_update;
pub(crate) mod update_transaction;
@@ -1977,6 +1977,23 @@ impl ProdContainerOrchestrator {
} }
async fn reconcile_all_with_mode(&self, mode: ReconcileMode) -> ReconcileReport { async fn reconcile_all_with_mode(&self, mode: ReconcileMode) -> ReconcileReport {
// Resolve exact retained update identities before normal desired-state
// reconciliation can recreate or start a member. Keep ownership for pass.
let _update_guard = match super::update_transaction::recover(
&self.data_dir,
&super::update_transaction::Podman,
)
.await
{
Ok(guard) => guard,
Err(error) => {
let mut report = ReconcileReport::default();
report
.failures
.push(("update-recovery".into(), format!("{error:#}")));
return report;
}
};
let user_stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await; let user_stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await;
// Durable desired-state signal: the container names that were running at // Durable desired-state signal: the container names that were running at
// the last periodic snapshot. Used below to recreate a previously-running // the last periodic snapshot. Used below to recreate a previously-running
@@ -0,0 +1,879 @@
//! Retained-container updates. This is runtime recovery, never database rollback.
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use std::{
collections::HashSet,
future::Future,
io::Write,
os::fd::AsRawFd,
os::unix::fs::{DirBuilderExt, OpenOptionsExt},
path::{Path, PathBuf},
};
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
pub(crate) struct Observed {
pub id: String,
pub name: String,
pub image: String,
pub running: bool,
pub config_sha256: String,
/// False for auto-remove, external supervision, pods, paused/unknown states.
pub retainable: bool,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub(crate) struct Target {
pub name: String,
pub reference: String,
pub image: String,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
struct Member {
original: Observed,
target: Target,
backup: String,
replacement_name: String,
replacement_id: Option<String>,
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
enum Phase {
Prepared,
Replacing,
Verified,
Committed,
Restored,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Record {
schema: u8,
operation: String,
package: String,
phase: Phase,
members: Vec<Member>,
}
/// An adapter must address destructive commands by immutable ID. `create` must
/// not start the new member, including when the original was stopped.
pub(crate) trait Runtime: Sync {
fn inspect(&self, id_or_name: &str) -> impl Future<Output = Result<Option<Observed>>> + Send;
fn stop(&self, id: &str) -> impl Future<Output = Result<()>> + Send;
fn start(&self, id: &str) -> impl Future<Output = Result<()>> + Send;
fn rename(&self, id: &str, name: &str) -> impl Future<Output = Result<()>> + Send;
fn create(
&self,
original: &str,
name: &str,
target: &str,
) -> impl Future<Output = Result<String>> + Send;
fn remove(&self, id: &str) -> impl Future<Output = Result<()>> + Send;
fn healthy(&self, id: &str) -> impl Future<Output = Result<bool>> + Send;
}
pub(crate) struct Guard {
_file: std::fs::File,
root: PathBuf,
}
impl Guard {
pub(crate) fn acquire(data: &Path) -> Result<Self> {
let root = data.join("update-transactions");
match std::fs::DirBuilder::new().mode(0o700).create(&root) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {}
Err(e) => return Err(e.into()),
}
anyhow::ensure!(
!std::fs::symlink_metadata(&root)?.file_type().is_symlink(),
"Update journal directory is a symlink"
);
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.mode(0o600)
.custom_flags(libc::O_NOFOLLOW)
.open(root.join("lock"))?;
anyhow::ensure!(
unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } == 0,
"An app lifecycle operation is active; retry when it completes"
);
Ok(Self { _file: file, root })
}
fn path(&self, operation: &str) -> Result<PathBuf> {
anyhow::ensure!(
uuid::Uuid::parse_str(operation)?.to_string() == operation,
"Invalid update operation"
);
Ok(self.root.join(format!("{operation}.json")))
}
fn save(&self, record: &Record) -> Result<()> {
record.validate()?;
let target = self.path(&record.operation)?;
let temp = self.root.join(format!(".{}.tmp", uuid::Uuid::new_v4()));
let bytes = serde_json::to_vec(record)?;
anyhow::ensure!(bytes.len() <= 1024 * 1024, "Update journal too large");
// No await/detached writer after the ownership lock can be released.
let result = (|| -> Result<()> {
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(&temp)?;
file.write_all(&bytes)?;
file.sync_all()?;
std::fs::rename(&temp, &target)?;
std::fs::File::open(&self.root)?.sync_all()?;
std::fs::File::open(self.root.parent().unwrap())?.sync_all()?;
Ok(())
})();
if result.is_err() {
let _ = std::fs::remove_file(temp);
}
result
}
fn records(&self) -> Result<Vec<Record>> {
let mut records = Vec::new();
for entry in std::fs::read_dir(&self.root)? {
let entry = entry?;
if entry.path().extension().and_then(|v| v.to_str()) != Some("json") {
continue;
}
anyhow::ensure!(records.len() < 256, "Too many retained update journals");
let meta = entry.metadata()?;
anyhow::ensure!(
entry.file_type()?.is_file() && meta.len() <= 1024 * 1024,
"Invalid update journal file"
);
let record: Record = serde_json::from_slice(&std::fs::read(entry.path())?)?;
record.validate()?;
anyhow::ensure!(
self.path(&record.operation)? == entry.path(),
"Update journal name changed"
);
records.push(record);
}
Ok(records)
}
pub(crate) fn require_clear(&self) -> Result<()> {
anyhow::ensure!(
self.records()?
.iter()
.all(|r| matches!(r.phase, Phase::Committed | Phase::Restored)),
"An interrupted app update needs recovery before another lifecycle action"
);
Ok(())
}
}
fn name_ok(value: &str) -> bool {
!value.is_empty()
&& value.len() <= 200
&& value
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.'))
}
fn id_ok(value: &str) -> bool {
value.len() == 64 && value.bytes().all(|b| b.is_ascii_hexdigit())
}
impl Record {
fn validate(&self) -> Result<()> {
anyhow::ensure!(
self.schema == 1
&& name_ok(&self.package)
&& !self.members.is_empty()
&& self.members.len() <= 32,
"Invalid update journal"
);
let operation = uuid::Uuid::parse_str(&self.operation)?;
anyhow::ensure!(
operation.to_string() == self.operation,
"Invalid update UUID"
);
let mut names = HashSet::new();
let mut ids = HashSet::new();
for (i, m) in self.members.iter().enumerate() {
anyhow::ensure!(
name_ok(&m.original.name)
&& m.target.name == m.original.name
&& id_ok(&m.original.id)
&& id_ok(&m.original.config_sha256)
&& ids.insert(&m.original.id)
&& names.insert(&m.original.name)
&& m.original.retainable
&& m.backup == format!("archy-update-{}-{i}-old", operation.simple())
&& m.replacement_name == format!("archy-update-{}-{i}-new", operation.simple())
&& m.replacement_id.as_deref().is_none_or(id_ok),
"Invalid retained update member"
);
}
Ok(())
}
}
async fn exact(runtime: &impl Runtime, id: &str) -> Result<Observed> {
let found = runtime
.inspect(id)
.await?
.context("Retained container is missing; runtime recovery unresolved")?;
anyhow::ensure!(
found.id == id,
"Container identity changed; recovery unresolved"
);
Ok(found)
}
/// All original inventory is proven before the first stop/rename. Unsupported
/// runtimes stay untouched. The caller has already prepared every target image.
pub(crate) async fn execute(
guard: &Guard,
package: &str,
targets: &[Target],
runtime: &impl Runtime,
) -> Result<()> {
guard.require_clear()?;
anyhow::ensure!(
!targets.is_empty() && targets.len() <= 32,
"Invalid update members"
);
let operation = uuid::Uuid::new_v4();
let mut members = Vec::new();
for (i, target) in targets.iter().enumerate() {
let original = runtime
.inspect(&target.name)
.await?
.context("Original update member is missing")?;
anyhow::ensure!(
original.retainable,
"{} cannot retain its original runtime safely; no containers changed",
target.name
);
members.push(Member {
original,
target: target.clone(),
backup: format!("archy-update-{}-{i}-old", operation.simple()),
replacement_name: format!("archy-update-{}-{i}-new", operation.simple()),
replacement_id: None,
});
}
let mut record = Record {
schema: 1,
operation: operation.to_string(),
package: package.into(),
phase: Phase::Prepared,
members,
};
guard.save(&record)?;
match replace(&guard,&mut record,runtime).await {
Ok(())=>Ok(()),
Err(error)=>match restore(&guard,&mut record,runtime).await {
Ok(())=>Err(error.context("Update failed; original runtime states restored (persistent data was not rolled back)")),
Err(recovery)=>Err(error.context(format!("Update failed; runtime recovery remains unresolved: {recovery:#}"))),
}
}
}
async fn replace(guard: &Guard, record: &mut Record, runtime: &impl Runtime) -> Result<()> {
for member in record.members.iter().rev() {
let current = exact(runtime, &member.original.id).await?;
anyhow::ensure!(
current.name == member.original.name
&& current.image == member.original.image
&& current.config_sha256 == member.original.config_sha256
&& current.retainable,
"Original changed before update"
);
if current.running {
runtime.stop(&current.id).await?;
}
anyhow::ensure!(
!exact(runtime, &current.id).await?.running,
"Original did not stop"
);
}
record.phase = Phase::Replacing;
guard.save(record)?;
for index in 0..record.members.len() {
let member = &record.members[index];
anyhow::ensure!(
runtime.inspect(&member.backup).await?.is_none()
&& runtime.inspect(&member.replacement_name).await?.is_none(),
"Update backup name collision"
);
runtime.rename(&member.original.id, &member.backup).await?;
let id = runtime
.create(
&member.original.id,
&member.replacement_name,
&member.target.reference,
)
.await?;
anyhow::ensure!(id_ok(&id), "Invalid replacement identity");
record.members[index].replacement_id = Some(id.clone());
guard.save(record)?;
let member = &record.members[index];
let created = exact(runtime, &id).await?;
anyhow::ensure!(
created.name == member.replacement_name
&& !created.running
&& created.image == member.target.image
&& created.config_sha256 == member.original.config_sha256,
"Replacement configuration/state mismatch"
);
runtime.rename(&id, &member.original.name).await?;
}
for member in &record.members {
let id = member
.replacement_id
.as_deref()
.context("Missing replacement")?;
if member.original.running {
runtime.start(id).await?;
}
let observed = exact(runtime, id).await?;
anyhow::ensure!(
observed.running == member.original.running && observed.image == member.target.image,
"Replacement state mismatch"
);
if member.original.running {
anyhow::ensure!(
runtime.healthy(id).await?,
"Updated member failed health verification"
);
}
}
record.phase = Phase::Verified;
guard.save(record)?;
// Original backups are retained. Committing does not remove volumes or
// silently discard the only recoverable original after a schema migration.
record.phase = Phase::Committed;
guard.save(record)
}
async fn restore(guard: &Guard, record: &mut Record, runtime: &impl Runtime) -> Result<()> {
// Preflight every original before mutating any replacement. Never infer
// ownership of a new container from a reused public name.
for member in &record.members {
let original = exact(runtime, &member.original.id).await?;
anyhow::ensure!(
original.image == member.original.image
&& original.config_sha256 == member.original.config_sha256
&& (original.name == member.original.name || original.name == member.backup),
"Original rollback identity/config changed"
);
if member.replacement_id.is_none() {
anyhow::ensure!(
runtime.inspect(&member.replacement_name).await?.is_none(),
"Unacknowledged replacement exists; preserve it for explicit recovery"
);
}
if let Some(current) = runtime.inspect(&member.original.name).await? {
anyhow::ensure!(
current.id == member.original.id
|| Some(&current.id) == member.replacement_id.as_ref(),
"Foreign container occupies original name"
);
}
}
for member in record.members.iter().rev() {
if let Some(id) = &member.replacement_id {
if let Some(current) = runtime.inspect(id).await? {
anyhow::ensure!(
current.id == *id && current.image == member.target.image,
"Replacement identity changed"
);
if current.running {
runtime.stop(id).await?;
}
runtime.remove(id).await?;
}
}
}
for member in &record.members {
let current = exact(runtime, &member.original.id).await?;
if current.name != member.original.name {
runtime.rename(&current.id, &member.original.name).await?;
}
if member.original.running && !current.running {
runtime.start(&current.id).await?;
} else if !member.original.running && current.running {
runtime.stop(&current.id).await?;
}
let observed = exact(runtime, &current.id).await?;
anyhow::ensure!(
observed.name == member.original.name
&& observed.image == member.original.image
&& observed.running == member.original.running,
"Original state was not restored"
);
}
record.phase = Phase::Restored;
guard.save(record)
}
/// Called before ordinary reconciliation. An unresolved recovery must prevent
/// reconciliation from guessing a replacement or starting stopped originals.
pub(crate) async fn recover(data: &Path, runtime: &impl Runtime) -> Result<Guard> {
let guard = Guard::acquire(data)?;
for mut record in guard.records()? {
if !matches!(record.phase, Phase::Committed | Phase::Restored) {
restore(&guard, &mut record, runtime).await?;
}
}
Ok(guard)
}
pub(crate) struct Podman;
impl Podman {
async fn output(args: &[&str]) -> Result<std::process::Output> {
tokio::time::timeout(
std::time::Duration::from_secs(120),
tokio::process::Command::new("podman")
.args(args)
.kill_on_drop(true)
.output(),
)
.await
.context("Container operation timed out; inspect original operation before retry")?
.context("Container runtime unavailable")
}
async fn command(args: &[&str]) -> Result<String> {
let output = Self::output(args).await?;
anyhow::ensure!(
output.status.success(),
"Container runtime rejected {} (exit {:?})",
args.first().unwrap_or(&"operation"),
output.status.code()
);
Ok(std::str::from_utf8(&output.stdout)?.trim().to_string())
}
pub(crate) async fn targets(images: &[(String, String)]) -> Result<Vec<Target>> {
// Confirm clone support before any stop. Never silently switch to a
// latest-catalog install if this runtime cannot create without running.
Self::command(&["container", "clone", "--help"]).await?;
let mut targets = Vec::new();
for (name, reference) in images {
let raw = Self::command(&["image", "inspect", reference]).await?;
let rows: Vec<serde_json::Value> = serde_json::from_str(&raw)?;
let image = rows
.first()
.and_then(|v| v.get("Id").or_else(|| v.get("ID")))
.and_then(|v| v.as_str())
.context("Prepared image identity unavailable")?;
targets.push(Target {
name: name.clone(),
reference: reference.clone(),
image: image.strip_prefix("sha256:").unwrap_or(image).into(),
});
}
Ok(targets)
}
}
impl Runtime for Podman {
async fn inspect(&self, value: &str) -> Result<Option<Observed>> {
let exists = Self::output(&["container", "exists", value]).await?;
match exists.status.code() {
Some(0) => {}
Some(1) => return Ok(None),
_ => anyhow::bail!("Container inventory unavailable"),
}
let raw = Self::command(&["inspect", value]).await?;
let rows: Vec<serde_json::Value> = serde_json::from_str(&raw)?;
let row = rows.first().context("Empty container inspection")?;
let string = |key: &str| {
row.get(key)
.and_then(|v| v.as_str())
.context("Incomplete container identity")
};
let state = row
.pointer("/State/Status")
.and_then(|v| v.as_str())
.context("Missing runtime state")?;
let labels = row.pointer("/Config/Labels");
let supervised = labels.and_then(|v| v.get("PODMAN_SYSTEMD_UNIT")).is_some()
|| labels
.and_then(|v| v.get("io.containers.systemd.unit"))
.is_some();
let supervised = supervised
|| row
.pointer("/Config/Env")
.and_then(|v| v.as_array())
.is_some_and(|env| {
env.iter().any(|v| {
v.as_str()
.is_some_and(|v| v.starts_with("PODMAN_SYSTEMD_UNIT="))
})
});
let restart = row
.pointer("/HostConfig/RestartPolicy/Name")
.and_then(|v| v.as_str())
.unwrap_or("");
let retainable = row
.pointer("/HostConfig/AutoRemove")
.and_then(|v| v.as_bool())
== Some(false)
&& !supervised
&& matches!(restart, "" | "no")
&& row
.get("Pod")
.and_then(|v| v.as_str())
.is_none_or(str::is_empty)
&& row
.get("Dependencies")
.and_then(|v| v.as_array())
.is_none_or(Vec::is_empty)
&& matches!(state, "running" | "stopped" | "exited" | "created");
let config = serde_json::json!({
"mounts":row.get("Mounts"),"env":row.pointer("/Config/Env"),
"command":row.pointer("/Config/Cmd"),"entrypoint":row.pointer("/Config/Entrypoint"),
"user":row.pointer("/Config/User"),"directory":row.pointer("/Config/WorkingDir"),
"ports":row.pointer("/HostConfig/PortBindings"),"network":row.pointer("/HostConfig/NetworkMode"),
"privileged":row.pointer("/HostConfig/Privileged"),"devices":row.pointer("/HostConfig/Devices"),
"security":row.pointer("/HostConfig/SecurityOpt"),"cap_add":row.pointer("/HostConfig/CapAdd"),"cap_drop":row.pointer("/HostConfig/CapDrop")
});
use sha2::Digest;
Ok(Some(Observed {
id: string("Id")?.into(),
name: string("Name")?.trim_start_matches('/').into(),
image: string("Image")?
.strip_prefix("sha256:")
.unwrap_or(string("Image")?)
.into(),
running: state == "running",
retainable,
config_sha256: hex::encode(sha2::Sha256::digest(serde_json::to_vec(&config)?)),
}))
}
async fn stop(&self, id: &str) -> Result<()> {
let original = exact(self, id).await?;
let grace = archipelago_container::runtime::stop_grace_secs_for(&original.name);
let timeout = grace.to_string();
let result = tokio::time::timeout(
std::time::Duration::from_secs(grace + 30),
tokio::process::Command::new("podman")
.args(["stop", "--time", &timeout, id])
.kill_on_drop(true)
.output(),
)
.await
.context("Graceful stop is unresolved; original runtime retained")??;
anyhow::ensure!(result.status.success(), "Graceful container stop failed");
Ok(())
}
async fn start(&self, id: &str) -> Result<()> {
Self::command(&["start", id]).await?;
Ok(())
}
async fn rename(&self, id: &str, name: &str) -> Result<()> {
Self::command(&["rename", id, name]).await?;
Ok(())
}
async fn create(&self, id: &str, name: &str, target: &str) -> Result<String> {
Self::command(&["container", "clone", id, name, target]).await
}
async fn remove(&self, id: &str) -> Result<()> {
Self::command(&["rm", id]).await?;
Ok(())
}
async fn healthy(&self, id: &str) -> Result<bool> {
// Missing healthcheck is not fabricated health. Running-only containers
// can pass runtime readiness; healthchecked members must report healthy.
for _ in 0..30 {
let raw = Self::command(&["inspect", id]).await?;
let rows: Vec<serde_json::Value> = serde_json::from_str(&raw)?;
let row = rows.first().context("Missing updated container")?;
if row.pointer("/State/Status").and_then(|v| v.as_str()) != Some("running") {
return Ok(false);
}
match row.pointer("/State/Health/Status").and_then(|v| v.as_str()) {
None | Some("") | Some("healthy") => return Ok(true),
Some("unhealthy") => return Ok(false),
_ => {}
}
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
Ok(false)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{
atomic::{AtomicBool, AtomicUsize, Ordering},
Mutex,
};
struct Mock {
rows: Mutex<Vec<Observed>>,
calls: Mutex<Vec<String>>,
fail_at: AtomicUsize,
sequence: AtomicUsize,
lost_create: AtomicBool,
fail_health: AtomicBool,
}
impl Mock {
fn new() -> Self {
Self {
rows: Mutex::new(vec![
Observed {
id: format!("{:064x}", 1),
name: "db".into(),
image: "old-db".into(),
running: true,
retainable: true,
config_sha256: "a".repeat(64),
},
Observed {
id: format!("{:064x}", 2),
name: "web".into(),
image: "old-web".into(),
running: false,
retainable: true,
config_sha256: "b".repeat(64),
},
]),
calls: Default::default(),
fail_at: AtomicUsize::new(0),
sequence: AtomicUsize::new(0),
lost_create: AtomicBool::new(false),
fail_health: AtomicBool::new(false),
}
}
fn event(&self, action: String) -> Result<()> {
self.calls.lock().unwrap().push(action);
let step = self.sequence.fetch_add(1, Ordering::SeqCst) + 1;
anyhow::ensure!(
step != self.fail_at.load(Ordering::SeqCst),
"Injected lifecycle failure"
);
Ok(())
}
fn targets() -> Vec<Target> {
vec![
Target {
name: "db".into(),
reference: "new-db".into(),
image: "new-db".into(),
},
Target {
name: "web".into(),
reference: "new-web".into(),
image: "new-web".into(),
},
]
}
fn originals_restored(&self) {
let rows = self.rows.lock().unwrap();
assert_eq!(rows.len(), 2);
for (i, name, running) in [(1, "db", true), (2, "web", false)] {
let row = rows.iter().find(|r| r.id == format!("{i:064x}")).unwrap();
assert_eq!(row.name, name);
assert_eq!(row.running, running);
assert_eq!(row.image, format!("old-{name}"));
}
}
}
impl Runtime for Mock {
async fn inspect(&self, value: &str) -> Result<Option<Observed>> {
Ok(self
.rows
.lock()
.unwrap()
.iter()
.find(|r| r.id == value || r.name == value)
.cloned())
}
async fn stop(&self, id: &str) -> Result<()> {
self.event(format!("stop:{id}"))?;
self.rows
.lock()
.unwrap()
.iter_mut()
.find(|r| r.id == id)
.unwrap()
.running = false;
Ok(())
}
async fn start(&self, id: &str) -> Result<()> {
self.event(format!("start:{id}"))?;
self.rows
.lock()
.unwrap()
.iter_mut()
.find(|r| r.id == id)
.unwrap()
.running = true;
Ok(())
}
async fn rename(&self, id: &str, name: &str) -> Result<()> {
self.event(format!("rename:{id}:{name}"))?;
let mut rows = self.rows.lock().unwrap();
anyhow::ensure!(!rows.iter().any(|r| r.name == name), "Name occupied");
rows.iter_mut().find(|r| r.id == id).unwrap().name = name.into();
Ok(())
}
async fn create(&self, original: &str, name: &str, target: &str) -> Result<String> {
self.event(format!("create-stopped:{name}"))?;
let mut rows = self.rows.lock().unwrap();
let mut row = rows.iter().find(|r| r.id == original).unwrap().clone();
row.id = format!("{:064x}", 100 + rows.len());
row.name = name.into();
row.image = target.into();
row.running = false;
let id = row.id.clone();
rows.push(row);
anyhow::ensure!(
!self.lost_create.swap(false, Ordering::SeqCst),
"Lost create acknowledgement"
);
Ok(id)
}
async fn remove(&self, id: &str) -> Result<()> {
self.event(format!("remove:{id}"))?;
self.rows.lock().unwrap().retain(|r| r.id != id);
Ok(())
}
async fn healthy(&self, _id: &str) -> Result<bool> {
Ok(!self.fail_health.load(Ordering::SeqCst))
}
}
#[tokio::test]
async fn mixed_stack_update_never_starts_stopped_member_and_retains_original_ids() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
let guard = Guard::acquire(root.path()).unwrap();
execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.unwrap();
let record = guard.records().unwrap().pop().unwrap();
assert_eq!(record.phase, Phase::Committed);
let rows = runtime.rows.lock().unwrap();
assert_eq!(rows.len(), 4);
for member in &record.members {
let original = rows.iter().find(|r| r.id == member.original.id).unwrap();
assert_eq!(original.name, member.backup);
assert!(!original.running);
let new = rows
.iter()
.find(|r| Some(&r.id) == member.replacement_id.as_ref())
.unwrap();
assert_eq!(new.running, member.original.running);
}
assert!(!runtime.calls.lock().unwrap().iter().any(|call| call
== &format!(
"start:{}",
record.members[1].replacement_id.as_ref().unwrap()
)));
assert!(!runtime
.calls
.lock()
.unwrap()
.iter()
.any(|call| call.starts_with("remove:")));
}
#[tokio::test]
async fn every_mutation_failure_restores_exact_originals_without_catalog_recreate() {
// Eight mutating calls in this two-member update, including starting db.
for fail in 1..=8 {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
runtime.fail_at.store(fail, Ordering::SeqCst);
let guard = Guard::acquire(root.path()).unwrap();
let error = execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.unwrap_err();
assert!(
error
.to_string()
.contains("original runtime states restored"),
"failure {fail}: {error:#}"
);
runtime.originals_restored();
assert_eq!(guard.records().unwrap()[0].phase, Phase::Restored);
}
}
#[tokio::test]
async fn unsupported_member_refuses_before_any_runtime_mutation() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
runtime.rows.lock().unwrap()[1].retainable = false;
let guard = Guard::acquire(root.path()).unwrap();
assert!(execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.is_err());
assert!(runtime.calls.lock().unwrap().is_empty());
assert!(guard.records().unwrap().is_empty());
runtime.originals_restored();
}
#[tokio::test]
async fn failed_health_restores_running_and_stopped_states() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
runtime.fail_health.store(true, Ordering::SeqCst);
let guard = Guard::acquire(root.path()).unwrap();
assert!(execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.is_err());
runtime.originals_restored();
}
#[tokio::test]
async fn lost_create_reply_is_unresolved_and_never_deletes_unowned_container() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
runtime.lost_create.store(true, Ordering::SeqCst);
let guard = Guard::acquire(root.path()).unwrap();
let error = execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.unwrap_err();
assert!(error.to_string().contains("unresolved"));
assert!(guard.require_clear().is_err());
assert!(!runtime
.calls
.lock()
.unwrap()
.iter()
.any(|call| call.starts_with("remove:")));
drop(guard);
assert!(recover(root.path(), &runtime).await.is_err());
assert_eq!(runtime.rows.lock().unwrap().len(), 3);
}
#[tokio::test]
async fn daemon_restart_restores_interrupted_rename_before_normal_reconciliation() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
let guard = Guard::acquire(root.path()).unwrap();
execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.unwrap();
let mut record = guard.records().unwrap().pop().unwrap();
// Simulate interruption after replacements existed but before durable commit.
record.phase = Phase::Replacing;
guard.save(&record).unwrap();
drop(guard);
let recovered = recover(root.path(), &runtime).await.unwrap();
runtime.originals_restored();
assert_eq!(recovered.records().unwrap()[0].phase, Phase::Restored);
}
#[tokio::test]
async fn missing_original_does_not_fabricate_rollback_or_start_other_members() {
let root = tempfile::tempdir().unwrap();
let runtime = Mock::new();
let guard = Guard::acquire(root.path()).unwrap();
execute(&guard, "stack", &Mock::targets(), &runtime)
.await
.unwrap();
let mut record = guard.records().unwrap().pop().unwrap();
record.phase = Phase::Replacing;
guard.save(&record).unwrap();
runtime
.rows
.lock()
.unwrap()
.retain(|r| r.id != format!("{:064x}", 1));
runtime.calls.lock().unwrap().clear();
drop(guard);
assert!(recover(root.path(), &runtime).await.is_err());
assert!(runtime.calls.lock().unwrap().is_empty());
}
#[test]
fn lifecycle_lock_excludes_competing_commands() {
let root = tempfile::tempdir().unwrap();
let guard = Guard::acquire(root.path()).unwrap();
assert!(Guard::acquire(root.path()).is_err());
drop(guard);
assert!(Guard::acquire(root.path()).is_ok());
}
}