Integrate legacy managed update maintenance and fenced recovery before reconciliation

This commit is contained in:
archipelago
2026-10-07 02:34:51 -04:00
parent 46fdc2764c
commit f79ecd11ae
7 changed files with 399 additions and 43 deletions
@@ -1,7 +1,7 @@
//! Production systemd/Podman adapter. Application write admission/drain is an
//! explicit dependency: neither process pause nor a filesystem receipt is drain.
use super::{
supervised_update::{self, PreparedTarget, RecoveryImage, Supervisor, Unit},
supervised_update::{self, Completion, PreparedTarget, RecoveryImage, Supervisor, Unit},
update_transaction::{Observed, Podman, Runtime, Target},
};
use anyhow::{Context, Result};
@@ -10,24 +10,155 @@ use std::{
collections::HashMap,
future::Future,
io::Write,
os::unix::fs::{MetadataExt, OpenOptionsExt},
os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt},
path::{Path, PathBuf},
time::Duration,
};
pub(crate) trait DrainBarrier: Sync {
/// Persist ownership before blocking admissions. Return only after API,
/// direct uploads and workers have drained and the coherent backup finished.
/// Called after the node persisted original writable recovery images.
/// Persist ownership before blocking admissions or stopping writers. Return
/// only after API/direct uploads/workers drained and coherent backup finished.
fn acquire(
&self,
operation: &str,
originals: &[Unit],
recovery: bool,
) -> impl Future<Output = Result<()>> + Send;
/// Must inspect the live admission fence and same-operation ownership, not
/// merely trust an old 'backup complete' file. Keep it through snapshot/stop.
fn verify(&self, operation: &str) -> impl Future<Output = Result<()>> + Send;
/// Idempotent; an old recovery must never clear a newer operation's fence.
fn release(&self, operation: &str) -> impl Future<Output = Result<()>> + Send;
fn release(
&self,
operation: &str,
outcome: Completion,
) -> impl Future<Output = Result<()>> + Send;
}
/// The controller is deployed from reviewed platform source. Its durable
/// journal owns legacy ingress/drain/backup across API/container shutdown.
pub(crate) struct LegacyIndeeMaintenance {
lock: std::sync::Arc<std::fs::File>,
}
impl LegacyIndeeMaintenance {
pub(crate) fn new(guard: &super::update_transaction::Guard) -> Result<Self> {
Self::controller_path()?;
Ok(Self {
lock: std::sync::Arc::new(guard.clone_lifecycle_lock()?),
})
}
fn controller_path() -> Result<&'static Path> {
let path = Path::new("/opt/archipelago/scripts/indeehub-maintenance-controller.py");
let metadata = std::fs::symlink_metadata(path)
.context("Legacy maintenance controller is not installed; original runtime retained")?;
anyhow::ensure!(
metadata.is_file()
&& !metadata.file_type().is_symlink()
&& (metadata.uid() == 0 || metadata.uid() == unsafe { libc::geteuid() }),
"Maintenance controller ownership changed"
);
anyhow::ensure!(
metadata.mode() & 0o022 == 0,
"Maintenance controller is writable by another user"
);
anyhow::ensure!(metadata.len() <= 128*1024
&& Sha256::digest(std::fs::read(path)?) == Sha256::digest(include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/../../scripts/indeehub-maintenance-controller.py"))),
"Maintenance controller does not match this qualified backend; install the exact companion script first");
Ok(path)
}
async fn invoke(&self, action: &'static str, request: serde_json::Value) -> Result<()> {
use std::os::{fd::AsRawFd, unix::process::CommandExt};
use tokio::io::AsyncWriteExt;
let path = Self::controller_path()?;
let operation = request["operation_id"]
.as_str()
.context("Missing maintenance operation")?
.to_owned();
let recovery = request.get("recovery").and_then(|v| v.as_bool()) == Some(true);
let bytes = serde_json::to_vec(&request)?;
let lock = self.lock.try_clone()?;
// Detaching the RPC future must not detach mutation from its lifecycle
// lock. This task and the child retain the same flock open description.
tokio::spawn(async move {
let fd = lock.as_raw_fd();
let mut command = tokio::process::Command::new("/usr/bin/python3");
command
.arg(path)
.arg(action)
.env("ARCHY_UPDATE_LOCK_FD", fd.to_string())
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::null())
.kill_on_drop(false);
unsafe {
command.as_std_mut().pre_exec(move || {
if libc::fcntl(fd, libc::F_SETFD, 0) == -1 {
return Err(std::io::Error::last_os_error());
}
Ok(())
});
}
let mut child = command
.spawn()
.context("Could not start legacy maintenance controller")?;
if let Some(mut stdin) = child.stdin.take() {
stdin.write_all(&bytes).await?;
stdin.shutdown().await?;
}
let result = child.wait_with_output().await?;
anyhow::ensure!(
result.status.success() && result.stdout.len() <= 4096,
"Legacy maintenance remains unresolved; inspect its original operation journal"
);
let reply: serde_json::Value = serde_json::from_slice(&result.stdout)?;
let wanted = match action {
"acquire" if recovery => "recovering",
"acquire" => "drained",
"verify" => "held",
"release" => "released",
_ => unreachable!(),
};
anyhow::ensure!(
reply["operation_id"].as_str() == Some(operation.as_str())
&& reply["state"].as_str() == Some(wanted),
"Maintenance controller did not acknowledge this operation"
);
drop(lock);
Ok::<_, anyhow::Error>(())
})
.await
.context("Maintenance completion task interrupted")?
}
}
impl DrainBarrier for LegacyIndeeMaintenance {
async fn acquire(&self, operation: &str, originals: &[Unit], recovery: bool) -> Result<()> {
let members: Vec<_> = originals
.iter()
.map(|unit| {
serde_json::json!({
"name":unit.name,"container_id":unit.container_id,"image_id":unit.image,
"unit_sha256":hex::encode(Sha256::digest(unit.body.as_bytes())),
"config_sha256":unit.config_sha256,"running":unit.running})
})
.collect();
self.invoke(
"acquire",
serde_json::json!({"operation_id":operation,"original_members":members,"recovery":recovery}),
)
.await
}
async fn verify(&self, operation: &str) -> Result<()> {
self.invoke("verify", serde_json::json!({"operation_id":operation}))
.await
}
async fn release(&self, operation: &str, outcome: Completion) -> Result<()> {
self.invoke(
"release",
serde_json::json!({"operation_id":operation,"outcome":outcome}),
)
.await
}
}
#[derive(Clone, serde::Serialize, serde::Deserialize)]
@@ -132,6 +263,22 @@ pub(crate) fn load_reviewed_plans(
Ok(record.plans)
}
pub(crate) async fn recover_before_reconcile(
data_dir: &Path,
guard: &super::update_transaction::Guard,
) -> Result<()> {
if !supervised_update::needs_recovery(guard)? {
return Ok(());
}
let adapter = SystemdSupervisor::new(
data_dir.into(),
HashMap::new(),
LegacyIndeeMaintenance::new(guard)?,
)
.await?;
supervised_update::recover(guard, &adapter).await
}
pub(crate) struct SystemdSupervisor<B> {
data_dir: PathBuf,
unit_dir: PathBuf,
@@ -224,14 +371,19 @@ impl<B: DrainBarrier> SystemdSupervisor<B> {
}
}
impl<B: DrainBarrier> Supervisor for SystemdSupervisor<B> {
async fn begin_barrier(&self, operation: &str, originals: &[Unit]) -> Result<()> {
self.barrier.acquire(operation, originals).await
async fn begin_barrier(
&self,
operation: &str,
originals: &[Unit],
recovery: bool,
) -> Result<()> {
self.barrier.acquire(operation, originals, recovery).await
}
async fn verify_barrier(&self, operation: &str) -> Result<()> {
self.barrier.verify(operation).await
}
async fn release_barrier(&self, operation: &str) -> Result<()> {
self.barrier.release(operation).await
async fn release_barrier(&self, operation: &str, outcome: Completion) -> Result<()> {
self.barrier.release(operation, outcome).await
}
async fn prepare_target(&self, target: &Target, original: &Unit) -> Result<PreparedTarget> {
let plan = self
@@ -303,6 +455,7 @@ impl<B: DrainBarrier> Supervisor for SystemdSupervisor<B> {
.mode(original.file_mode)
.open(&temporary)?;
file.write_all(body.as_bytes())?;
file.set_permissions(std::fs::Permissions::from_mode(original.file_mode))?;
file.sync_all()?;
std::fs::rename(&temporary, &path)?;
std::fs::File::open(&self.unit_dir)?.sync_all()?;
@@ -314,7 +467,8 @@ impl<B: DrainBarrier> Supervisor for SystemdSupervisor<B> {
result
}
async fn snapshot(&self, original: &Unit, operation: &str, tag: &str) -> Result<RecoveryImage> {
self.barrier.verify(operation).await?;
// Pre-maintenance image capture preserves code and hook mutations;
// mounted-data consistency is established separately by acquire().
supervised_update::capture_local_recovery_image(original, operation, tag).await
}
async fn pin(&self, image: &str, tag: &str) -> Result<()> {
@@ -366,13 +520,13 @@ mod tests {
use super::*;
struct UnusedBarrier;
impl DrainBarrier for UnusedBarrier {
async fn acquire(&self, _: &str, _: &[Unit]) -> Result<()> {
async fn acquire(&self, _: &str, _: &[Unit], _: bool) -> Result<()> {
anyhow::bail!("No application barrier installed")
}
async fn verify(&self, _: &str) -> Result<()> {
anyhow::bail!("No application barrier installed")
}
async fn release(&self, _: &str) -> Result<()> {
async fn release(&self, _: &str, _: Completion) -> Result<()> {
Ok(())
}
}
@@ -422,6 +576,9 @@ mod tests {
original.body = original_body.into();
target.reference = format!("localhost/movie@sha256:{}", "e".repeat(64));
assert!(adapter.prepare_target(&target, &original).await.is_err());
assert!(adapter.begin_barrier("unused", &[original]).await.is_err());
assert!(adapter
.begin_barrier("unused", &[original], false)
.await
.is_err());
}
}
@@ -46,6 +46,7 @@ enum Phase {
Aborted,
Editing,
Starting,
Restoring,
Committed,
Restored,
}
@@ -57,19 +58,39 @@ struct Journal {
package: String,
phase: Phase,
members: Vec<Member>,
#[serde(default)]
cleanup_done: bool,
#[serde(default)]
target_startup_began: bool,
}
#[derive(Clone, Copy, Debug, serde::Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum Completion {
Committed,
Restored,
Aborted,
}
pub(crate) trait Supervisor: Sync {
/// Admission must be fenced and in-flight application work drained under
/// this durable operation before snapshots. A paused process is not proof.
/// Called only after every original writable recovery image is durable.
/// Legacy acquisition may gracefully stop AutoRemove writers, so its
/// destructive obligation is journaled before entering the controller.
/// It must fence/drain writers and finish coherent mounted-data backup;
/// image snapshots alone never establish application consistency.
/// Acquisition/release are idempotent and operation-owned across restart.
fn begin_barrier(
&self,
operation: &str,
originals: &[Unit],
recovery: bool,
) -> impl Future<Output = Result<()>> + Send;
fn verify_barrier(&self, operation: &str) -> impl Future<Output = Result<()>> + Send;
fn release_barrier(&self, operation: &str) -> impl Future<Output = Result<()>> + Send;
fn release_barrier(
&self,
operation: &str,
outcome: Completion,
) -> impl Future<Output = Result<()>> + Send;
/// Only an internal reviewed signed-manifest planner may supply this value;
/// browser parameters must never become a unit body or hook recipe.
fn prepare_target(
@@ -518,9 +539,8 @@ fn records(guard: &Guard) -> Result<Vec<Journal>> {
}
pub(crate) fn require_clear(guard: &Guard) -> Result<()> {
anyhow::ensure!(
records(guard)?
.iter()
.all(|r| matches!(r.phase, Phase::Committed | Phase::Restored | Phase::Aborted)),
records(guard)?.iter().all(|r| r.cleanup_done
&& matches!(r.phase, Phase::Committed | Phase::Restored | Phase::Aborted)),
"A supervised update needs recovery first"
);
Ok(())
@@ -574,6 +594,8 @@ pub(crate) async fn execute(
package: package.into(),
phase: Phase::Prepared,
members,
cleanup_done: false,
target_startup_began: false,
};
save(guard, &record)?;
for member in &record.members {
@@ -593,12 +615,6 @@ pub(crate) async fn execute(
Ok(())
}
async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> {
let originals: Vec<_> = record
.members
.iter()
.map(|member| member.original.clone())
.collect();
supervisor.begin_barrier(&record.id, &originals).await?;
for index in 0..record.members.len() {
let member = &record.members[index];
supervisor.validate_original_file(&member.original).await?;
@@ -606,7 +622,6 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor
supervisor.read(&member.original.name).await? == member.original.body,
"Unit edited before update; originals retained"
);
supervisor.verify_barrier(&record.id).await?;
let image = supervisor
.snapshot(&member.original, &record.id, &member.original_tag)
.await?;
@@ -626,9 +641,19 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor
// The image acknowledgement becomes durable before any original stop.
save(guard, record)?;
}
supervisor.verify_barrier(&record.id).await?;
record.phase = Phase::Editing;
save(guard, record)?;
// The legacy controller may stop original writers. Every writable layer
// is already recoverable and this obligation survives cancellation.
let originals: Vec<_> = record
.members
.iter()
.map(|member| member.original.clone())
.collect();
supervisor
.begin_barrier(&record.id, &originals, false)
.await?;
supervisor.verify_barrier(&record.id).await?;
for member in record.members.iter().rev() {
supervisor.verify_barrier(&record.id).await?;
supervisor.stop(&member.original.name).await?;
@@ -644,6 +669,8 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor
}
supervisor.reload().await?;
record.phase = Phase::Starting;
record.target_startup_began = true;
record.cleanup_done = false;
save(guard, record)?;
for member in &record.members {
supervisor.start(&member.original.name).await?;
@@ -663,11 +690,14 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor
}
record.phase = Phase::Committed;
save(guard, record)?;
supervisor.release_barrier(&record.id).await?;
supervisor
.release_barrier(&record.id, Completion::Committed)
.await?;
for member in &record.members {
guard.release_hold(&member.original.name, &record.id)?;
}
Ok(())
record.cleanup_done = true;
save(guard, record)
}
async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> {
if record.phase == Phase::Prepared {
@@ -689,12 +719,25 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis
}
record.phase = Phase::Aborted;
save(guard, record)?;
supervisor.release_barrier(&record.id).await?;
supervisor
.release_barrier(&record.id, Completion::Aborted)
.await?;
for member in &record.members {
guard.release_hold(&member.original.name, &record.id)?;
}
return Ok(());
record.cleanup_done = true;
return save(guard, record);
}
let originals: Vec<_> = record
.members
.iter()
.map(|member| member.original.clone())
.collect();
record.phase = Phase::Restoring;
save(guard, record)?;
supervisor
.begin_barrier(&record.id, &originals, true)
.await?;
supervisor.verify_barrier(&record.id).await?;
// Refuse to overwrite a foreign edit before stopping any surviving member.
for member in &record.members {
@@ -767,25 +810,44 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis
}
record.phase = Phase::Restored;
save(guard, record)?;
supervisor.release_barrier(&record.id).await
supervisor
.release_barrier(&record.id, Completion::Restored)
.await?;
record.cleanup_done = true;
save(guard, record)
}
pub(crate) async fn recover(guard: &Guard, supervisor: &impl Supervisor) -> Result<()> {
for mut record in records(guard)? {
if record.cleanup_done {
continue;
}
match record.phase {
Phase::Committed | Phase::Aborted => {
supervisor.release_barrier(&record.id).await?;
let outcome = if record.phase == Phase::Committed {
Completion::Committed
} else {
Completion::Aborted
};
supervisor.release_barrier(&record.id, outcome).await?;
for member in &record.members {
guard.release_hold(&member.original.name, &record.id)?;
}
}
Phase::Restored => {
supervisor.release_barrier(&record.id).await?;
supervisor
.release_barrier(&record.id, Completion::Restored)
.await?;
} // Never release a newer owner.
_ => restore(guard, &mut record, supervisor).await?,
}
record.cleanup_done = true;
save(guard, &record)?;
}
Ok(())
}
pub(crate) fn needs_recovery(guard: &Guard) -> Result<bool> {
Ok(records(guard)?.iter().any(|record| !record.cleanup_done))
}
#[cfg(test)]
mod tests {
@@ -837,7 +899,12 @@ mod tests {
}
}
impl Supervisor for Mock {
async fn begin_barrier(&self, operation: &str, _originals: &[Unit]) -> Result<()> {
async fn begin_barrier(
&self,
operation: &str,
_originals: &[Unit],
_recovery: bool,
) -> Result<()> {
let mut held = self.barrier.lock().unwrap();
anyhow::ensure!(
held.as_deref().is_none_or(|id| id == operation),
@@ -853,7 +920,7 @@ mod tests {
);
Ok(())
}
async fn release_barrier(&self, operation: &str) -> Result<()> {
async fn release_barrier(&self, operation: &str, _outcome: Completion) -> Result<()> {
anyhow::ensure!(
!self.fail_barrier_release.load(Ordering::SeqCst),
"Barrier release unavailable"
@@ -1062,7 +1129,9 @@ mod tests {
execute(&guard, "movie", &[Mock::target()], &runtime)
.await
.unwrap();
let record = records(&guard).unwrap().pop().unwrap();
let mut record = records(&guard).unwrap().pop().unwrap();
record.cleanup_done = false;
save(&guard, &record).unwrap();
guard.hold("movie", &record.id).unwrap();
runtime.calls.lock().unwrap().clear();
recover(&guard, &runtime).await.unwrap();
@@ -1125,6 +1194,7 @@ mod tests {
.unwrap();
let mut record = records(&guard).unwrap().pop().unwrap();
record.phase = Phase::Starting;
record.cleanup_done = false;
*runtime.barrier.lock().unwrap() = Some(record.id.clone());
save(&guard, &record).unwrap();
runtime.calls.lock().unwrap().clear();
@@ -1150,6 +1220,7 @@ mod tests {
.unwrap();
let mut record = records(&guard).unwrap().pop().unwrap();
record.phase = Phase::Starting;
record.cleanup_done = false;
*runtime.barrier.lock().unwrap() = Some(record.id.clone());
save(&guard, &record).unwrap();
*runtime.body.lock().unwrap() = "operator replaced unit".into();
@@ -74,6 +74,9 @@ pub(crate) struct Guard {
root: PathBuf,
}
impl Guard {
pub(crate) fn clone_lifecycle_lock(&self) -> Result<std::fs::File> {
Ok(self._file.try_clone()?)
}
pub(crate) fn directory(&self) -> &Path {
&self.root
}
@@ -500,6 +503,7 @@ pub(crate) async fn recover(data: &Path, runtime: &impl Runtime) -> Result<Guard
}
}
}
super::supervised_runtime::recover_before_reconcile(data, &guard).await?;
Ok(guard)
}
@@ -544,10 +548,15 @@ impl Podman {
);
Ok(std::str::from_utf8(&output.stdout)?.trim().to_string())
}
pub(crate) async fn targets(images: &[(String, String)]) -> Result<Vec<Target>> {
pub(crate) async fn targets(
images: &[(String, String)],
require_clone: bool,
) -> 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?;
if require_clone {
Self::command(&["container", "clone", "--help"]).await?;
}
let mut targets = Vec::new();
for (name, reference) in images {
let raw = Self::command(&["image", "inspect", reference]).await?;