Integrate retained updater candidate for release qualification

This commit is contained in:
archipelago
2026-10-07 13:50:26 -04:00
22 changed files with 5389 additions and 168 deletions
@@ -362,7 +362,11 @@ impl RpcHandler {
set_package_state(
&handler.state_manager,
&package_id_spawn,
if result.get("status").and_then(|v| v.as_str()) == Some("up-to-date") {
if result.get("status").and_then(|v| v.as_str()) == Some("staged") {
PackageState::Stopped
} else if result.get("status").and_then(|v| v.as_str())
== Some("up-to-date")
{
pre_state.clone().unwrap_or(PackageState::Running)
} else {
PackageState::Running
@@ -373,12 +377,34 @@ impl RpcHandler {
Err(e) => {
error!("package.update {} failed: {:#}", package_id_spawn, e);
install_log(&format!("UPDATE FAIL: {} — {:#}", package_id_spawn, e)).await;
// Inner handler already ran rollback_update + cleared
// update state, but be defensive: revert to pre-state
// in case the inner flow died before its cleanup.
if let Some(prev) = pre_state {
set_package_state(&handler.state_manager, &package_id_spawn, prev).await;
}
// Release the transitional overlay before asking the scanner
// for real state. Prior Running is not proof of successful
// rollback, and a failed preflight is not proof of Stopped.
handler
.state_manager
.mutate_data(|data| {
if let Some(entry) = data.package_data.get_mut(&package_id_spawn) {
finish_failed_update(entry);
}
data.notifications.retain(|item| {
item.id != format!("update-failed-{package_id_spawn}")
});
data.notifications.push(crate::data_model::Notification {
id: format!("update-failed-{package_id_spawn}"),
level: crate::data_model::NotificationLevel::Error,
title: format!("Could not update {package_id_spawn}"),
message: format!(
"{e}. Runtime recovery does not roll back database changes."
),
timestamp: chrono::Utc::now().to_rfc3339(),
app_id: Some(package_id_spawn.clone()),
});
while data.notifications.len() > 20 {
data.notifications.remove(0);
}
})
.await;
kick_scanner_and_wait(&handler).await;
}
}
});
@@ -581,3 +607,34 @@ async fn kick_scanner_and_wait(handler: &RpcHandler) {
})
.await;
}
fn finish_failed_update(entry: &mut crate::data_model::PackageDataEntry) {
if entry.state == PackageState::Updating {
entry.state = PackageState::Installed;
}
entry.install_progress = None;
}
#[cfg(test)]
mod update_completion_tests {
use super::*;
#[test]
fn failure_releases_spinner_without_inventing_stopped_or_restored_runtime() {
let mut entry = super::super::progress::create_installing_entry("movie");
entry.state = PackageState::Updating;
finish_failed_update(&mut entry);
assert_eq!(entry.state, PackageState::Installed);
assert!(entry.install_progress.is_none());
for actual in [
PackageState::Running,
PackageState::Stopped,
PackageState::Exited,
] {
entry.state = actual.clone();
finish_failed_update(&mut entry);
assert_eq!(
entry.state, actual,
"Fresh scanner evidence must win over old pre-update intent"
);
}
}
}
+38 -1
View File
@@ -526,7 +526,18 @@ pub(in crate::api::rpc) async fn get_containers_for_app(package_id: &str) -> Res
.await
.context("podman ps timed out while listing containers")?
.context("Failed to list containers")?;
let stdout = String::from_utf8_lossy(&output.stdout);
containers_from_list_output(package_id, &output)
}
fn containers_from_list_output(
package_id: &str,
output: &std::process::Output,
) -> Result<Vec<String>> {
anyhow::ensure!(
output.status.success(),
"podman ps failed while listing containers"
);
let stdout = std::str::from_utf8(&output.stdout).context("Invalid container list response")?;
let all: Vec<&str> = stdout.lines().filter(|s| !s.is_empty()).collect();
let patterns = all_container_names(package_id);
@@ -543,6 +554,32 @@ pub(in crate::api::rpc) async fn get_containers_for_app(package_id: &str) -> Res
mod tests {
use super::{all_container_names, get_data_dirs_for_app, get_health_check_args};
#[test]
fn failed_container_listing_is_not_an_absent_app() {
use std::os::unix::process::ExitStatusExt;
let mut output = std::process::Output {
status: std::process::ExitStatus::from_raw(1 << 8),
stdout: vec![],
stderr: b"store unavailable".to_vec(),
};
assert!(super::containers_from_list_output("node-demo-music", &output).is_err());
output.stdout = b"node-demo-music\n".to_vec();
assert!(super::containers_from_list_output("node-demo-music", &output).is_err());
output.status = std::process::ExitStatus::from_raw(0);
assert_eq!(
super::containers_from_list_output("node-demo-music", &output).unwrap(),
vec!["node-demo-music"]
);
output.stdout.clear();
assert!(
super::containers_from_list_output("node-demo-music", &output)
.unwrap()
.is_empty()
);
output.stdout = vec![0xff];
assert!(super::containers_from_list_output("node-demo-music", &output).is_err());
}
#[test]
fn bitcoin_variant_container_names_are_precise() {
let core = all_container_names("bitcoin-core");
@@ -280,6 +280,10 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?;
let docker_image = params
.get("dockerImage")
@@ -1901,7 +1905,10 @@ autopilot.active=false\n",
.await
.context("DATUM credentials are not available yet; wait for installation to finish")?;
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!({
"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.",
@@ -60,6 +60,10 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?;
// A cuprate node that starts on a too-small disk fills it and takes
// Archipelago down with it (no upstream pruning — see
// dependencies::check_cuprate_disk_compatibility). Fail the start
@@ -104,6 +108,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id);
let data_dir = self.config.data_dir.clone();
tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_start(orchestrator.as_ref(), &to_start).await
@@ -126,7 +131,19 @@ impl RpcHandler {
Err(e) => {
tracing::error!("package.start {} failed: {:#}", package_id_owned, e);
install_log(&format!("START FAIL: {} — {:#}", package_id_owned, e)).await;
if let Some(prev) = pre_state {
if e.downcast_ref::<crate::container::prod_orchestrator::StagedCleanupFailure>()
.is_some()
{
// Installed is neutral and scanner-owned. Restoring the
// prior Stopped state would conceal a failed cleanup;
// keeping Starting would prevent scanner convergence.
set_package_state(
&state_manager,
&package_id_owned,
PackageState::Installed,
)
.await;
} else if let Some(prev) = pre_state {
set_package_state(&state_manager, &package_id_owned, prev).await;
} else {
warn!(
@@ -155,6 +172,10 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?;
let single_orchestrator_app =
self.orchestrator.is_some() && uses_single_orchestrator_app(package_id);
@@ -219,6 +240,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id);
tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_stop(orchestrator.as_ref(), &to_stop).await
@@ -257,6 +279,10 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?;
// Restart is stop + recreate, so on a disk that shrank below the cuprate
// minimum after install it resumes the doomed unprunable sync just like
// start would — same gate, same "fail before clearing user-stopped /
@@ -319,6 +345,7 @@ impl RpcHandler {
let op_lock = app_op_lock(package_id);
let data_dir = self.config.data_dir.clone();
tokio::spawn(async move {
let _lifecycle_guard = lifecycle_guard;
let _op_guard = op_lock.lock().await;
let result = if let Some(orchestrator) = orchestrator.as_ref() {
do_orchestrator_package_restart(orchestrator.as_ref(), &to_restart).await
@@ -362,6 +389,10 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing package id"))?;
validate_app_id(package_id)?;
let lifecycle_guard =
crate::container::update_transaction::Guard::acquire(&self.config.data_dir)?;
lifecycle_guard.require_clear()?;
lifecycle_guard.require_unheld(&super::config::all_container_names(package_id))?;
let preserve_data = params
.get("preserve_data")
.and_then(|v| v.as_bool())
+225 -146
View File
@@ -4,15 +4,15 @@
//! remove old container(s) → recreate (orchestrator-first, legacy fallback) → verify running.
//! Data volumes are preserved (bind mounts, not stored in container).
use super::config::get_containers_for_app;
use super::config::{all_container_names, get_containers_for_app};
use super::install::install_log;
use super::progress::parse_pull_progress;
use super::runtime::stop_timeout_secs;
use super::validation::validate_app_id;
use crate::api::rpc::RpcHandler;
use crate::container::image_versions;
use crate::data_model::{InstallPhase, PackageState};
use anyhow::{Context, Result};
use std::{collections::HashSet, path::Path};
use tokio::io::{AsyncBufReadExt, BufReader};
use tracing::{error, info, warn};
@@ -31,6 +31,9 @@ impl RpcHandler {
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing 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
// button. Fetch and verify first; failure leaves running containers alone.
@@ -60,9 +63,24 @@ impl RpcHandler {
let targets = pinned
.as_ref()
.map(|target| self.resolve_images_to_pull(package_id, target));
// A stopped Quadlet normally removes its --rm container. Absence is
// not an install decision: only a known single managed app with durable
// installed AND stopped evidence may enter the recreate path.
let known_managed =
if should_try_orchestrator_update(package_id, self.orchestrator.is_some()) {
self.orchestrator
.as_ref()
.expect("orchestrator presence checked")
.knows_app(orchestrator_update_app_id(package_id))
.await
} else {
false
};
let markers = UpdateMarkers::load(&self.config.data_dir, package_id).await?;
let installed = inspect_update_images(package_id).await?;
validate_update_presence(package_id, known_managed, !installed.is_empty(), &markers)?;
if let Some(targets) = &targets {
let installed = inspect_update_images(package_id).await?;
if !update_targets_need_change(targets, &installed)? {
if !update_targets_need_change(targets, &installed)? && !markers.stopped {
install_log(&format!(
"UPDATE SKIP: {} — target versions already installed",
package_id
@@ -111,6 +129,19 @@ impl RpcHandler {
if let Some(orchestrator) = self.orchestrator.as_ref() {
match orchestrator.upgrade(orchestrator_app_id).await {
Ok(()) => {
if let Some(image) = orchestrator
.staged_upgrade_image(orchestrator_app_id)
.await?
{
// The orchestrator proved the pinned image exists and
// persisted its reviewed manifest. No running container
// or healthy service is claimed for a stopped update.
self.clear_install_progress(package_id).await;
return Ok(serde_json::json!({
"status": "staged", "state": "stopped",
"package_id": package_id, "image": image,
}));
}
if let Some(targets) = &targets {
verify_update_targets(
targets,
@@ -188,7 +219,7 @@ impl RpcHandler {
match preflighted_stack_update(
&images_to_pull,
|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
{
@@ -212,10 +243,14 @@ impl RpcHandler {
package_id, e
))
.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_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
)))
}
}
}
@@ -260,114 +295,59 @@ 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(
&self,
package_id: &str,
containers: &[String],
images_to_pull: &[(String, String)],
guard: &crate::container::update_transaction::Guard,
) -> Result<()> {
// Phase: Preparing — about to stop the running container(s) so
// we can swap images. Fast.
use crate::container::update_transaction::{self, Podman};
let mut supervised = 0;
for name in containers {
if crate::container::quadlet::unit_exists(name).await {
supervised += 1;
}
}
anyhow::ensure!(
supervised == 0 || supervised == containers.len(),
"Mixed managed/unmanaged stack requires an explicit recovery plan; originals unchanged"
);
let managed = supervised != 0;
anyhow::ensure!(
!managed || package_id == "indeedhub",
"This managed app has no qualified maintenance controller; originals unchanged"
);
let targets = Podman::targets(images_to_pull, !managed).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)
.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()
if managed {
use crate::container::supervised_runtime::{
load_reviewed_plans, LegacyIndeeMaintenance, SystemdSupervisor,
};
let plans = load_reviewed_plans(&self.config.data_dir, package_id)?;
let adapter = SystemdSupervisor::new(
self.config.data_dir.clone(),
plans,
LegacyIndeeMaintenance::new(guard)?,
)
.await?;
let targets = adapter.reviewed_targets(&targets).await?;
crate::container::supervised_update::execute(guard, package_id, &targets, &adapter)
.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
}
} else {
update_transaction::execute(guard, package_id, &targets, &Podman).await
}
// 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(
@@ -526,61 +506,84 @@ impl RpcHandler {
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).
async fn clear_update_state(&self, package_id: &str) {
let (mut data, _) = self.state_manager.get_snapshot().await;
if let Some(entry) = data.package_data.get_mut(package_id) {
// Don't overwrite state from scanner — just clear if still Updating
if entry.state == PackageState::Updating {
entry.state = PackageState::Stopped;
// Unknown is not stopped: the authoritative scanner will
// refresh actual retained runtime immediately in the wrapper.
entry.state = PackageState::Installed;
}
}
self.state_manager.update_data(data).await;
}
}
async fn inspect_update_images(package_id: &str) -> Result<Vec<(String, String)>> {
let containers = get_containers_for_app(package_id).await?;
#[derive(Default)]
struct UpdateMarkers {
installed: bool,
stopped: bool,
uninstalled: bool,
}
impl UpdateMarkers {
async fn load(data_dir: &Path, package_id: &str) -> Result<Self> {
let mut names = all_container_names(package_id);
names.push(package_id.to_string());
names.push(orchestrator_update_app_id(package_id).to_string());
let installed = read_update_markers(data_dir, "installed-apps.json").await?;
let stopped = read_update_markers(data_dir, "user-stopped.json").await?;
let uninstalled = read_update_markers(data_dir, "user-uninstalled.json").await?;
Ok(Self {
installed: names.iter().any(|name| installed.contains(name)),
stopped: names.iter().any(|name| stopped.contains(name)),
uninstalled: names.iter().any(|name| uninstalled.contains(name)),
})
}
}
// The recovery loaders intentionally return empty on malformed state. An update
// must instead distinguish unavailable evidence from a proven absence of an
// uninstall decision before it can recreate a removed Quadlet container.
async fn read_update_markers(data_dir: &Path, filename: &str) -> Result<HashSet<String>> {
match tokio::fs::read(data_dir.join(filename)).await {
Ok(bytes) => serde_json::from_slice(&bytes)
.with_context(|| format!("Cannot read {filename}; update cancelled")),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(HashSet::new()),
Err(error) => {
Err(error).with_context(|| format!("Cannot read {filename}; update cancelled"))
}
}
}
fn validate_update_presence(
package_id: &str,
known_managed: bool,
has_containers: bool,
markers: &UpdateMarkers,
) -> Result<()> {
anyhow::ensure!(
!containers.is_empty(),
"No containers found for {}",
!markers.uninstalled,
"{} was uninstalled; use an explicit install instead of update",
package_id
);
anyhow::ensure!(
has_containers || (known_managed && markers.installed && markers.stopped),
"No containers or durable stopped installation found for {}",
package_id
);
Ok(())
}
async fn inspect_update_images(package_id: &str) -> Result<Vec<(String, String)>> {
let containers = get_containers_for_app(package_id).await?;
if containers.is_empty() {
// The caller decides whether durable installation evidence authorizes
// a missing runtime. Never invoke bare `podman inspect` here.
return Ok(Vec::new());
}
let mut command = tokio::process::Command::new("podman");
command.arg("inspect").args(&containers).kill_on_drop(true);
let output = tokio::time::timeout(std::time::Duration::from_secs(30), command.output())
@@ -626,6 +629,7 @@ fn verify_update_targets(
targets: &[(String, String)],
installed: &[(String, String)],
) -> Result<()> {
anyhow::ensure!(!targets.is_empty(), "No update targets resolved");
for (app_id, target) in targets {
let running = installed_image_for_target(app_id, installed).ok_or_else(|| {
anyhow::anyhow!("Update {}: target container missing after recreate", app_id)
@@ -651,6 +655,7 @@ fn update_targets_need_change(
installed: &[(String, String)],
) -> Result<bool> {
use std::cmp::Ordering;
anyhow::ensure!(!targets.is_empty(), "No update targets resolved");
let mut changed = false;
for (app_id, target) in targets {
let running = installed_image_for_target(app_id, installed);
@@ -771,9 +776,83 @@ mod tests {
use super::{
candidate_app_ids_for_container, immutable_update_image, orchestrator_update_app_id,
should_try_orchestrator_update, update_targets_need_change, uses_legacy_update_flow,
verify_update_targets,
validate_update_presence, verify_update_targets, UpdateMarkers,
};
#[tokio::test]
async fn stopped_installed_managed_app_updates_after_quadlet_container_disappears() {
let root = tempfile::tempdir().unwrap();
crate::crash_recovery::mark_installed(root.path(), "node-demo-music").await;
crate::crash_recovery::mark_user_stopped(root.path(), "node-demo-music").await;
let markers = UpdateMarkers::load(root.path(), "node-demo-music")
.await
.unwrap();
validate_update_presence("node-demo-music", true, false, &markers).unwrap();
let target = vec![("node-demo-music".into(), "localhost/music:2".into())];
assert!(super::update_targets_need_change(&target, &[]).unwrap());
// Success must still prove that upgrade actually created the target.
assert!(verify_update_targets(&target, &[]).is_err());
assert!(verify_update_targets(&target, &target).is_ok());
assert!(crate::crash_recovery::load_user_stopped(root.path())
.await
.contains("node-demo-music"));
}
#[test]
fn absent_catalog_app_or_unmanaged_runtime_is_not_an_update_installation() {
for (managed, installed, stopped) in [
(true, false, false),
(true, false, true),
(true, true, false),
(false, true, true),
] {
let markers = UpdateMarkers {
installed,
stopped,
uninstalled: false,
};
assert!(validate_update_presence("optional", managed, false, &markers).is_err());
}
// Existing legacy containers remain updateable without modern markers.
assert!(validate_update_presence("legacy", false, true, &UpdateMarkers::default()).is_ok());
assert!(super::update_targets_need_change(&[], &[]).is_err());
assert!(verify_update_targets(&[], &[]).is_err());
}
#[tokio::test]
async fn uninstall_tombstone_beats_stale_installed_and_stopped_markers() {
let root = tempfile::tempdir().unwrap();
crate::crash_recovery::mark_installed(root.path(), "node-demo-music").await;
crate::crash_recovery::mark_user_stopped(root.path(), "node-demo-music").await;
crate::crash_recovery::mark_user_uninstalled(root.path(), "archy-node-demo-music").await;
let markers = UpdateMarkers::load(root.path(), "node-demo-music")
.await
.unwrap();
for present in [false, true] {
assert!(validate_update_presence("node-demo-music", true, present, &markers).is_err());
}
}
#[tokio::test]
async fn damaged_lifecycle_markers_cannot_authorize_recreation() {
let root = tempfile::tempdir().unwrap();
let absent = UpdateMarkers::load(root.path(), "optional").await.unwrap();
assert!(validate_update_presence("optional", true, false, &absent).is_err());
for name in [
"installed-apps.json",
"user-stopped.json",
"user-uninstalled.json",
] {
tokio::fs::write(root.path().join(name), b"not-json")
.await
.unwrap();
assert!(UpdateMarkers::load(root.path(), "optional").await.is_err());
tokio::fs::remove_file(root.path().join(name))
.await
.unwrap();
}
}
#[tokio::test]
async fn stack_image_failure_precedes_every_lifecycle_action() {
use std::sync::{Arc, Mutex};
+40
View File
@@ -146,6 +146,16 @@ impl AppGate {
// snapshot when the port momentarily leaves the map mid-refresh.
let live = self.port_map.read().await.gated(app.port).cloned();
let app = live.as_ref().unwrap_or(app);
if maintenance_active(&self.data_dir, &app.app_id) {
return Response::builder()
.status(StatusCode::SERVICE_UNAVAILABLE)
.header(header::CACHE_CONTROL, "no-store")
.header(header::RETRY_AFTER, "60")
.body(Body::from(
"This app is temporarily unavailable while its update is recovered.",
))
.unwrap();
}
let path = req.uri().path().to_string();
// A dashboard same-origin proxy strips `/app/<id>/` before this gate
@@ -838,6 +848,16 @@ dashboard and check {name} under My Apps.</p>"#,
resp
}
fn maintenance_active(data_dir: &std::path::Path, app_id: &str) -> bool {
if app_id != "indeedhub" && !app_id.starts_with("indeedhub-") {
return false;
}
match std::fs::symlink_metadata(data_dir.join("app-maintenance/indeedhub")) {
Ok(_) => true,
Err(error) => error.kind() != std::io::ErrorKind::NotFound,
}
}
fn not_found() -> Response<Body> {
Response::builder()
.status(StatusCode::NOT_FOUND)
@@ -1857,3 +1877,23 @@ mod tests {
)
}
}
#[cfg(test)]
mod maintenance_tests {
#[test]
fn held_indee_ingress_never_reaches_upstream_or_another_app() {
let root = tempfile::tempdir().unwrap();
assert!(!super::maintenance_active(root.path(), "indeedhub"));
std::fs::create_dir(root.path().join("app-maintenance")).unwrap();
std::fs::write(
root.path().join("app-maintenance/indeedhub"),
uuid::Uuid::new_v4().to_string(),
)
.unwrap();
assert!(super::maintenance_active(root.path(), "indeedhub"));
assert!(super::maintenance_active(root.path(), "indeedhub-api"));
assert!(!super::maintenance_active(root.path(), "node-demo-v4v"));
std::fs::remove_file(root.path().join("app-maintenance/indeedhub")).unwrap();
assert!(!super::maintenance_active(root.path(), "indeedhub"));
}
}
+48 -2
View File
@@ -83,6 +83,32 @@ const RUNTIME_ASSETS_DIR: &str = "/opt/archipelago/web-ui/archipelago-runtime";
/// Inserted into every server block of the nginx config that lacks the
/// `/api/app-catalog` proxy. Kept in sync with the canonical block in
/// image-recipe/configs/nginx-archipelago.conf.
const INDEEHUB_MAINTENANCE_GUARD: &str =
"if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }";
/// Patch only recognized literal IndeeHub routes, including asset/WebSocket
/// sublocations; unknown operator routes are not guessed by this repair.
fn heal_indeehub_maintenance_guards(content: &str) -> String {
let route=regex::Regex::new(r"(?m)^([ \t]*)(location[ \t]+(?:\^~[ \t]+|=[ \t]+)?/app/indeedhub(?:/[^\s{]*)?[ \t]*\{)[ \t]*$").unwrap();
let mut result = String::new();
let mut previous = 0;
for capture in route.captures_iter(content) {
let full = capture.get(0).unwrap();
result.push_str(&content[previous..full.end()]);
if !content[full.end()..]
.trim_start()
.starts_with(INDEEHUB_MAINTENANCE_GUARD)
{
result.push_str(&format!(
"\n{} {}",
&capture[1], INDEEHUB_MAINTENANCE_GUARD
));
}
previous = full.end();
}
result.push_str(&content[previous..]);
result
}
const NGINX_APP_CATALOG_BLOCK: &str = "\n # App Store catalog proxy — backend fetches from configured registries\n # so the browser doesn't hit CORS/CSP. Without this block nginx falls\n # through to the SPA index.html and the frontend gets HTML back instead\n # of JSON.\n location ~ ^/api/(?:app-catalog|node-app-catalog)$ {\n proxy_pass http://127.0.0.1:5678;\n proxy_http_version 1.1;\n proxy_set_header Host $host;\n proxy_set_header X-Real-IP $remote_addr;\n proxy_set_header Cookie $http_cookie;\n proxy_connect_timeout 15s;\n proxy_read_timeout 30s;\n proxy_send_timeout 15s;\n error_page 502 503 = @backend_unavailable;\n error_page 504 = @backend_timeout;\n }\n\n";
const NGINX_SOURCE_PROXY_BLOCK: &str = " # GitWorkshop follows the dashboard origin so LAN, Tailscale, FIPS, Tor,\n # hostnames and reverse proxies all use the connection that already works.\n location /app/archipelago-source/ {\n proxy_pass http://127.0.0.2:8337/;\n proxy_http_version 1.1;\n proxy_set_header Host $http_host;\n proxy_set_header Cookie $http_cookie;\n proxy_set_header X-Real-IP $remote_addr;\n proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;\n proxy_set_header X-Forwarded-Proto $scheme;\n proxy_set_header X-Forwarded-Prefix /app/archipelago-source;\n proxy_hide_header X-Frame-Options;\n add_header X-Frame-Options \"SAMEORIGIN\" always;\n add_header X-Content-Type-Options \"nosniff\" always;\n proxy_read_timeout 300s;\n }\n";
@@ -2011,8 +2037,10 @@ async fn patch_nginx_conf(path: &str) -> Result<bool> {
let missing_source_prefix = heal_source_forwarded_prefix(&content).is_some();
let missing_nostr_signer = heal_missing_nostr_signer(&content).is_some();
let missing_rental_playback = heal_rental_playback_route(&content) != content;
let missing_maintenance = heal_indeehub_maintenance_guards(&content) != content;
let legacy_catalog_route = content.contains("location /api/app-catalog {");
if !missing_rental_playback
if !missing_maintenance
&& !missing_rental_playback
&& !missing_app_catalog
&& !legacy_catalog_route
&& !missing_bitcoin_status
@@ -2031,7 +2059,9 @@ async fn patch_nginx_conf(path: &str) -> Result<bool> {
return Ok(false);
}
let mut patched = heal_rental_playback_route(&heal_node_catalog_route(&content));
let mut patched = heal_indeehub_maintenance_guards(&heal_rental_playback_route(
&heal_node_catalog_route(&content),
));
if let Some(p) = heal_stale_web_search_block(&patched) {
patched = p;
@@ -2515,3 +2545,19 @@ pub async fn ensure_restart_policy() {
Err(e) => tracing::warn!(error = %e, "could not repair archipelago.service restart policy"),
}
}
#[cfg(test)]
mod indeehub_maintenance_tests {
#[test]
fn legacy_routes_are_fenced_independently_and_repair_is_idempotent() {
let source="server {\n location /app/indeedhub/ {\n proxy_pass http://127.0.0.1:7778;\n }\n location /app/indeedhub/ws/ {\n proxy_pass http://127.0.0.1:7778;\n }\n location /app/other/ {\n proxy_pass http://127.0.0.1:7777;\n }\n}\n";
let repaired = super::heal_indeehub_maintenance_guards(source);
assert_eq!(
repaired.matches(super::INDEEHUB_MAINTENANCE_GUARD).count(),
2
);
assert_eq!(super::heal_indeehub_maintenance_guards(&repaired), repaired);
assert!(repaired.contains("location /app/other/ {\n proxy_pass"));
assert_eq!(repaired.matches("proxy_pass").count(), 3);
}
}
@@ -298,7 +298,22 @@ impl DockerPackageScanner {
uninstall_stage: None,
};
apply_manifest_presentation(&app_id, &mut package);
match super::staged_update::installed_manifest(data_dir, &app_id, &container.image)
.await
{
Ok(Some(manifest)) => {
apply_manifest_value(&serde_json::to_value(manifest)?, &mut package)
}
Ok(None) => apply_manifest_presentation(&app_id, &mut package),
Err(error) => {
tracing::warn!(app_id, error = %error, "Cannot verify retained installed manifest entry point");
// Do not invent a launch path from a later catalog when its
// binding to the running version cannot be established.
if let Some(installed) = package.installed.as_mut() {
installed.interface_addresses.clear();
}
}
}
packages.insert(app_id.clone(), package);
info!(
"Detected container: {} ({})",
+12
View File
@@ -85,6 +85,18 @@ pub async fn run_post_install(manifest: &AppManifest, container_name: &str, data
}
}
/// Strict completion for an explicitly staged update. Keep the durable pin on failure.
pub(super) async fn run_post_install_strict(
manifest: &AppManifest,
container: &str,
data_dir: &Path,
) -> Result<()> {
for step in &manifest.app.hooks.post_install {
run_step(step, container, &manifest.app.id, data_dir).await?;
}
Ok(())
}
async fn run_step(step: &HookStep, container: &str, app_id: &str, data_dir: &Path) -> Result<()> {
match step {
HookStep::Exec { exec } => {
+9 -2
View File
@@ -1,5 +1,4 @@
pub mod app_catalog;
pub mod node_catalog;
pub mod app_gate_config;
pub mod bitcoin_ui;
pub mod boot_reconciler;
@@ -14,11 +13,12 @@ pub mod image_policy;
pub mod image_versions;
pub mod lnd;
pub mod migration_backup;
pub mod node_catalog;
pub mod npm;
pub mod prod_orchestrator;
pub mod quadlet;
pub mod registry;
pub mod registration_pin;
pub mod registry;
pub mod secrets;
pub mod traits;
pub mod ui_detection;
@@ -29,3 +29,10 @@ pub use dev_orchestrator::DevContainerOrchestrator;
pub use docker_packages::DockerPackageScanner;
pub use prod_orchestrator::ProdContainerOrchestrator;
pub use traits::ContainerOrchestrator;
mod staged_update;
pub(crate) mod update_transaction;
pub(crate) mod supervised_runtime;
pub(crate) mod supervised_update;
@@ -1422,6 +1422,17 @@ fn host_port_collisions<'m>(
out
}
/// A staged start failed and its runtime could not be proven stopped. Callers
/// must not restore a cached Stopped state; the scanner/reconciler owns recovery.
#[derive(Debug)]
pub(crate) struct StagedCleanupFailure(String);
impl std::fmt::Display for StagedCleanupFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Display::fmt(&self.0, f)
}
}
impl std::error::Error for StagedCleanupFailure {}
/// Internal: track a manifest together with the absolute directory it was loaded
/// from, so Build sources can resolve relative `context:` paths.
#[derive(Debug, Clone)]
@@ -1966,6 +1977,33 @@ impl ProdContainerOrchestrator {
}
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 held_names = match _update_guard.held_names() {
Ok(names) => names,
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;
// Durable desired-state signal: the container names that were running at
// the last periodic snapshot. Used below to recreate a previously-running
@@ -1987,6 +2025,7 @@ impl ProdContainerOrchestrator {
let filtered = state
.manifests
.iter()
.filter(|(_, lm)| !held_names.contains(&compute_container_name(&lm.manifest)))
.filter(|(app_id, _)| !state.disabled.contains(*app_id))
.filter(|(app_id, lm)| {
dependency_required.contains(*app_id)
@@ -2043,7 +2082,9 @@ impl ProdContainerOrchestrator {
}
}
}
let mut report = ReconcileReport::default();
// Stopped/disabled apps are excluded from ordinary start reconciliation.
// Pending cleanup must still run for them, including after a restart.
let mut report = self.reconcile_staged_updates().await;
let disk_gb = self.disk_gb().await;
let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB
|| crate::settings::bitcoin_storage::load(&self.data_dir)
@@ -2336,6 +2377,12 @@ impl ProdContainerOrchestrator {
let lock = self.app_lock(&app_id).await;
let _guard = lock.lock().await;
if let Some(record) = super::staged_update::load(&self.data_dir, &app_id).await? {
let name = compute_container_name(&record.manifest);
self.stop_staged_runtime(&name).await?;
return Ok(ReconcileAction::Left("staged-update-awaiting-start".into()));
}
self.ensure_app_secrets(&app_id).await?;
// Don't fight the Bitcoin-implementation switch: bitcoin-core and
@@ -2849,11 +2896,172 @@ impl ProdContainerOrchestrator {
}
/// Build-or-pull, create, start. Assumes the per-app mutex is already held.
async fn reconcile_staged_updates(&self) -> ReconcileReport {
let mut report = ReconcileReport::default();
let ids = match super::staged_update::pending_ids(&self.data_dir).await {
Ok(ids) => ids,
Err(error) => {
report
.failures
.push(("staged-updates".into(), error.to_string()));
return report;
}
};
for app_id in ids {
if crate::app_ops::lifecycle_op_in_flight(&app_id) {
continue;
}
let lock = self.app_lock(&app_id).await;
let _guard = lock.lock().await;
// Re-read after acquiring the lock: a successful explicit start
// may have completed and cleared the stage while this pass waited.
let result = async {
if let Some(record) = super::staged_update::load(&self.data_dir, &app_id).await? {
self.stop_staged_runtime(&compute_container_name(&record.manifest))
.await?;
report.record(
&app_id,
ReconcileAction::Left("staged-update-awaiting-start".into()),
);
}
Ok::<_, anyhow::Error>(())
}
.await;
if let Err(error) = result {
report.failures.push((app_id, error.to_string()));
}
}
report
}
async fn stop_staged_runtime(&self, name: &str) -> Result<()> {
let mut errors = Vec::new();
if let Err(e) = self.remove_quadlet_unit_if_present(name).await {
errors.push(format!("could not disable staged unit: {e:#}"));
}
let inventory = self
.runtime
.list_containers()
.await
.context("Cannot verify staged runtime inventory")?;
if inventory.iter().any(|c| {
c.name.trim_start_matches('/') == name
&& !matches!(
c.state,
ContainerState::Stopped | ContainerState::Exited | ContainerState::Created
)
}) {
if let Err(e) = self.runtime.stop_container(name).await {
errors.push(format!("could not stop staged container: {e:#}"));
}
}
let inventory = self
.runtime
.list_containers()
.await
.context("Cannot verify staged runtime stopped")?;
if inventory.iter().any(|c| {
c.name.trim_start_matches('/') == name
&& !matches!(
c.state,
ContainerState::Stopped | ContainerState::Exited | ContainerState::Created
)
}) {
errors.push("staged container is still active".into());
}
anyhow::ensure!(errors.is_empty(), "{}", errors.join("; "));
Ok(())
}
async fn start_staged_update(&self, app_id: &str) -> Result<bool> {
use super::staged_update;
if staged_update::load(&self.data_dir, app_id).await?.is_none() {
return Ok(false);
}
let lock = self.app_lock(app_id).await;
let _guard = lock.lock().await;
let record = staged_update::load(&self.data_dir, app_id)
.await?
.context("Staged update changed; retry start")?;
let name = compute_container_name(&record.manifest);
anyhow::ensure!(
!staged_update::marked(&self.data_dir, "user-uninstalled.json", app_id, &name).await?,
"App was uninstalled; staged update cannot reinstall it"
);
anyhow::ensure!(
record.ready,
"Update preparation was interrupted; retry update before start"
);
anyhow::ensure!(
self.runtime.image_exists(record.image()?).await?,
"Staged image is unavailable; refusing to start another version"
);
// Clear only this explicit start's stopped intent. The pending manifest
// remains durable until all creation, hooks and readiness checks succeed.
crate::crash_recovery::clear_user_stopped(&self.data_dir, app_id).await;
crate::crash_recovery::clear_user_stopped(&self.data_dir, &name).await;
crate::crash_recovery::clear_user_stopped(&self.data_dir, &format!("archy-{app_id}")).await;
self.state.write().await.disabled.remove(app_id);
let lm = LoadedManifest {
manifest: record.manifest.clone(),
manifest_dir: record.manifest_dir.clone(),
};
let result = async {
self.remove_quadlet_unit_if_present(&name).await?;
if self
.runtime
.list_containers()
.await?
.iter()
.any(|c| c.name.trim_start_matches('/') == name)
{
self.runtime.stop_container(&name).await?;
self.runtime.remove_container(&name).await?;
}
self.install_fresh_with_pin(&lm, true).await?;
let observed = self.runtime.get_container_status(&name).await?;
anyhow::ensure!(
observed.state == ContainerState::Running,
"Pinned start did not become running"
);
anyhow::ensure!(
staged_update::matches_image(record.image()?, &observed.image),
"Pinned start produced a different image"
);
// Installed presentation is durable before pending intent disappears.
// A crash on either side of clear cannot fall back to a later catalog.
staged_update::save_installed(&self.data_dir, &record).await?;
staged_update::clear(&self.data_dir, app_id).await
}
.await;
if let Err(error) = result {
crate::crash_recovery::mark_user_stopped(&self.data_dir, app_id).await;
self.state.write().await.disabled.insert(app_id.into());
let retained = staged_update::save(&self.data_dir, &record).await;
let cleanup = self.stop_staged_runtime(&name).await;
if let Err(cleanup) = cleanup {
return Err(StagedCleanupFailure(format!(
"Staged start failed: {error:#}; cleanup failed: {cleanup:#}; pin persistence: {:?}", retained.err()
)).into());
}
retained.context("Could not retain failed staged start pin")?;
return Err(error);
}
Ok(true)
}
async fn install_fresh(&self, lm: &LoadedManifest) -> Result<()> {
self.install_fresh_with_pin(lm, false).await
}
async fn install_fresh_with_pin(&self, lm: &LoadedManifest, pinned: bool) -> Result<()> {
anyhow::ensure!(super::supervised_update::installed_unit(&self.data_dir, &compute_container_name(&lm.manifest))?.is_none(), "Reviewed managed runtime is missing; restore its saved unit/image explicitly instead of recreating from the current catalog");
self.ensure_app_secrets(&lm.manifest.app.id).await?;
let mut resolved_manifest = lm.manifest.clone();
self.resolve_dynamic_env(&mut resolved_manifest).await?;
resolve_catalog_image(&mut resolved_manifest);
if !pinned {
resolve_catalog_image(&mut resolved_manifest);
}
let resolved = resolved_manifest.app.container.resolve().ok_or_else(|| {
anyhow::anyhow!(
@@ -3008,7 +3216,17 @@ impl ProdContainerOrchestrator {
// freshly created container — exactly when container mutations (e.g.
// indeedhub's nginx X-Frame-Options strip + nostr-provider injection) must
// be re-applied. Best-effort + idempotent: never fails the install.
crate::container::hooks::run_post_install(&resolved_manifest, &name, &self.data_dir).await;
if pinned {
crate::container::hooks::run_post_install_strict(
&resolved_manifest,
&name,
&self.data_dir,
)
.await?;
} else {
crate::container::hooks::run_post_install(&resolved_manifest, &name, &self.data_dir)
.await;
}
if uses_pasta_network(&resolved_manifest) {
if let Err(err) = wait_for_manifest_host_ports(
&resolved_manifest,
@@ -3065,6 +3283,16 @@ impl ProdContainerOrchestrator {
}
async fn prepare_for_start(&self, manifest: &AppManifest) -> Result<()> {
if super::supervised_update::installed_unit(
&self.data_dir,
&compute_container_name(manifest),
)?
.is_some()
{
// Already-reviewed managed configuration is authoritative. Applying
// today's catalog files/mount preparation could mutate restored data.
return Ok(());
}
self.run_pre_start_hooks(&manifest.app.id).await?;
self.ensure_bind_mount_sockets(manifest).await?;
self.ensure_bind_mount_dirs(manifest).await?;
@@ -3137,6 +3365,9 @@ impl ProdContainerOrchestrator {
lm: &LoadedManifest,
name: &str,
) -> Result<Option<ReconcileAction>> {
if super::supervised_update::installed_unit(&self.data_dir, name)?.is_some() {
return Ok(None); // Never remove an original runtime to migrate it from changed catalog data.
}
// Skip companion apps — bitcoin-ui / electrs-ui / lnd-ui have shipped
// via Quadlet since v1.7.41 (companion.rs renders the unit). Running
// migration for them races companion rendering: when migration ran
@@ -3236,6 +3467,18 @@ impl ProdContainerOrchestrator {
/// changes restart the service and retain a durable pending marker until
/// that succeeds, including across daemon restarts and failed reloads.
async fn sync_quadlet_unit(&self, lm: &LoadedManifest, name: &str) -> Result<()> {
if super::update_transaction::is_held(&self.data_dir, name)? {
return Ok(()); // Preserve the recovered unit instead of current catalog drift.
}
if let Some((body, mode)) = super::supervised_update::installed_unit(&self.data_dir, name)?
{
use std::os::unix::fs::PermissionsExt;
let path = quadlet::unit_dir().await?.join(format!("{name}.container"));
let meta = std::fs::symlink_metadata(&path)
.context("Saved managed unit is missing; recovery required")?;
anyhow::ensure!(meta.is_file() && meta.permissions().mode() & 0o777 == mode && std::fs::read_to_string(&path)? == body, "Reviewed managed unit changed; preserve it for explicit reconciliation instead of overwriting operator configuration");
return Ok(());
}
// Companions: same reasoning as migrate_to_quadlet_if_needed —
// companion.rs renders these units with a different shape, syncing
// here would clobber them.
@@ -3352,6 +3595,17 @@ impl ProdContainerOrchestrator {
}
async fn ensure_resolved_source_available(&self, lm: &LoadedManifest) -> Result<()> {
if let Some((body, _)) = super::supervised_update::installed_unit(
&self.data_dir,
&compute_container_name(&lm.manifest),
)? {
let image = body
.lines()
.find_map(|line| line.trim().strip_prefix("Image="))
.context("Saved managed image missing")?;
anyhow::ensure!(self.runtime.image_exists(image).await?, "Saved managed image is unavailable; refusing to pull or recreate from a changed catalog");
return Ok(());
}
let resolved = lm.manifest.app.container.resolve().ok_or_else(|| {
anyhow::anyhow!(
"manifest for {} has invalid container source (neither image nor build)",
@@ -4825,6 +5079,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
}
async fn install(&self, app_id: &str) -> Result<String> {
anyhow::ensure!(super::staged_update::load(&self.data_dir, app_id).await?.is_none(),
"An update is staged; use explicit start, or uninstall before selecting another installation");
let lm = self.loaded(app_id).await?;
// Optional shared-service preconditions are checked before recording
// installation or creating anything. A headless adapter must not claim
@@ -4896,7 +5152,7 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
// point). Just delegate.
self.prepare_for_start(&lm.manifest).await?;
let action = self.ensure_running(&lm).await?;
match action {
let result = match action {
ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed => {
Ok(name)
}
@@ -4915,10 +5171,17 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
self.install_fresh(&lm).await?;
Ok(name)
}
};
if result.is_ok() {
super::staged_update::clear_installed(&self.data_dir, app_id).await?;
}
result
}
async fn start(&self, app_id: &str) -> Result<()> {
if self.start_staged_update(app_id).await? {
return Ok(());
}
if let Some(members) = self.mempool_umbrella_members(app_id).await {
tracing::info!(
app_id,
@@ -5134,8 +5397,11 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
let _guard = lock.lock().await;
for name in [app_id.to_string(), format!("archy-{app_id}")] {
self.remove_quadlet_unit_if_present(&name).await?;
super::supervised_update::forget_installed(&self.data_dir, &name)?;
}
self.state.write().await.disabled.insert(app_id.to_string());
super::staged_update::clear(&self.data_dir, app_id).await?;
super::staged_update::clear_installed(&self.data_dir, app_id).await?;
return Ok(());
}
let lm = self.loaded(app_id).await?;
@@ -5189,15 +5455,84 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
// stale claim behind would let desired-state recovery recreate the very
// app that was just uninstalled.
crate::crash_recovery::clear_installed(&self.data_dir, app_id).await;
super::staged_update::clear(&self.data_dir, app_id).await?;
super::staged_update::clear_installed(&self.data_dir, app_id).await?;
super::supervised_update::forget_installed(&self.data_dir, &name)?;
Ok(())
}
/// Upgrade: stop-remove-reinstall (re-pulls or rebuilds as required).
async fn upgrade(&self, app_id: &str) -> Result<()> {
let lm = self.loaded(app_id).await?;
use super::staged_update;
let lock = self.app_lock(app_id).await;
let _guard = lock.lock().await;
let existing = staged_update::load(&self.data_dir, app_id).await?;
let lm = if let Some(record) = &existing {
LoadedManifest {
manifest: record.manifest.clone(),
manifest_dir: record.manifest_dir.clone(),
}
} else {
self.loaded(app_id).await?
};
let name = compute_container_name(&lm.manifest);
anyhow::ensure!(
!staged_update::marked(&self.data_dir, "user-uninstalled.json", app_id, &name).await?,
"App was uninstalled; update cannot reinstall it"
);
let stopped =
staged_update::marked(&self.data_dir, "user-stopped.json", app_id, &name).await?;
if stopped || existing.is_some() {
anyhow::ensure!(
staged_update::marked(&self.data_dir, "installed-apps.json", app_id, &name).await?,
"Stopped update needs durable installation evidence"
);
let mut record = match existing {
Some(record) => record,
None => {
let mut manifest = lm.manifest.clone();
resolve_catalog_image(&mut manifest);
staged_update::StagedUpdate {
manifest,
manifest_dir: lm.manifest_dir.clone(),
ready: false,
}
}
};
record.image()?;
crate::crash_recovery::mark_user_stopped(&self.data_dir, app_id).await;
anyhow::ensure!(
staged_update::marked(&self.data_dir, "user-stopped.json", app_id, &name).await?,
"Could not persist stopped update intent"
);
// Freeze the approved manifest before any pull. Retries cannot
// silently switch to another catalog revision after an interruption.
record.ready = false;
staged_update::save(&self.data_dir, &record).await?;
self.remove_quadlet_unit_if_present(&name).await?;
if self
.runtime
.list_containers()
.await?
.iter()
.any(|c| c.name.trim_start_matches('/') == name)
{
self.runtime.stop_container(&name).await?;
self.runtime.remove_container(&name).await?;
}
let pinned = LoadedManifest {
manifest: record.manifest.clone(),
manifest_dir: record.manifest_dir.clone(),
};
self.ensure_resolved_source_available(&pinned).await?;
anyhow::ensure!(
self.runtime.image_exists(record.image()?).await?,
"Staged image unavailable after preparation"
);
record.ready = true;
staged_update::save(&self.data_dir, &record).await?;
return Ok(());
}
let mut resolved = lm.manifest.clone();
resolve_catalog_image(&mut resolved);
if resolved.app.container.build.is_none() {
@@ -5213,7 +5548,10 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
running.image,
target
),
Some(std::cmp::Ordering::Equal) => return Ok(()),
Some(std::cmp::Ordering::Equal) => {
staged_update::clear_installed(&self.data_dir, app_id).await?;
return Ok(());
}
_ => {}
}
}
@@ -5221,7 +5559,22 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
}
let _ = self.runtime.stop_container(&name).await;
let _ = self.runtime.remove_container(&name).await;
self.install_fresh(&lm).await
self.install_fresh(&lm).await?;
staged_update::clear_installed(&self.data_dir, app_id).await
}
async fn staged_upgrade_image(&self, app_id: &str) -> Result<Option<String>> {
match super::staged_update::load(&self.data_dir, app_id).await? {
Some(record) => {
anyhow::ensure!(record.ready, "Staged update is not ready");
anyhow::ensure!(
self.runtime.image_exists(record.image()?).await?,
"Staged image unavailable"
);
Ok(Some(record.image()?.to_string()))
}
None => Ok(None),
}
}
async fn status(&self, app_id: &str) -> Result<ContainerStatus> {
@@ -5778,6 +6131,7 @@ mod tests {
fail_image_exists: StdMutex<HashMap<String, String>>,
/// If set, `start_container` for this container fails with this message.
fail_start: StdMutex<HashMap<String, String>>,
fail_stop: StdMutex<HashMap<String, String>>,
}
impl MockRuntime {
@@ -5841,6 +6195,12 @@ mod tests {
) -> Result<String> {
self.record(format!("create_container:{name}:offset={port_offset}"));
self.set_state(name, ContainerState::Created);
if let Some(image) = manifest.app.container.image_ref() {
self.running_images
.lock()
.unwrap()
.insert(name.into(), image);
}
self.created_env
.lock()
.unwrap()
@@ -5862,6 +6222,9 @@ mod tests {
}
async fn stop_container(&self, name: &str) -> Result<()> {
self.record(format!("stop_container:{name}"));
if let Some(error) = self.fail_stop.lock().unwrap().remove(name) {
return Err(anyhow::anyhow!(error));
}
self.set_state(name, ContainerState::Stopped);
Ok(())
}
@@ -6310,6 +6673,333 @@ app:
}
}
async fn stopped_update_fixture() -> (Arc<MockRuntime>, ProdContainerOrchestrator, String) {
let rt = Arc::new(MockRuntime::default());
let orch = orch_with(rt.clone()).await;
let image = format!("docker.io/library/alpine@sha256:{}", "a".repeat(64));
let mut manifest = pull_manifest("staged-music", &image);
manifest
.app
.environment
.push("STAGED_VERSION=original".into());
orch.insert_manifest_for_test(manifest, PathBuf::from("/tmp"))
.await;
crate::crash_recovery::mark_installed(&orch.data_dir, "staged-music").await;
crate::crash_recovery::mark_user_stopped(&orch.data_dir, "staged-music").await;
(rt, orch, image)
}
#[tokio::test]
async fn stopped_upgrade_stages_without_start_then_starts_pinned_manifest() {
let (rt, orch, image) = stopped_update_fixture().await;
let sentinel = orch.data_dir.join("persistent-song-and-secret-fixture");
tokio::fs::write(&sentinel, b"preserved").await.unwrap();
orch.upgrade("staged-music").await.unwrap();
assert_eq!(
orch.staged_upgrade_image("staged-music").await.unwrap(),
Some(image)
);
assert!(!rt
.calls()
.iter()
.any(|c| c.starts_with("create_container:") || c.starts_with("start_container:")));
let replacement = pull_manifest(
"staged-music",
&format!("docker.io/library/alpine@sha256:{}", "b".repeat(64)),
);
orch.insert_manifest_for_test(replacement, PathBuf::from("/tmp/new-catalog"))
.await;
let lm = orch.loaded("staged-music").await.unwrap();
assert_eq!(
orch.ensure_running(&lm).await.unwrap(),
ReconcileAction::Left("staged-update-awaiting-start".into())
);
orch.start("staged-music").await.unwrap();
assert!(rt
.created_env_for("staged-music")
.contains(&"STAGED_VERSION=original".into()));
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Running
);
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_none());
assert!(!crate::crash_recovery::load_user_stopped(&orch.data_dir)
.await
.contains("staged-music"));
assert_eq!(tokio::fs::read(sentinel).await.unwrap(), b"preserved");
}
#[tokio::test]
async fn interrupted_stopped_pull_retains_original_pin_and_requires_retry() {
let (rt, orch, image) = stopped_update_fixture().await;
*rt.fail_pull.lock().unwrap() = Some("interrupted pull".into());
assert!(orch.upgrade("staged-music").await.is_err());
assert!(orch
.start("staged-music")
.await
.unwrap_err()
.to_string()
.contains("interrupted"));
let record = super::super::staged_update::load(&orch.data_dir, "staged-music")
.await
.unwrap()
.unwrap();
assert!(!record.ready);
assert_eq!(record.image().unwrap(), image);
*rt.fail_pull.lock().unwrap() = None;
orch.insert_manifest_for_test(
pull_manifest("staged-music", "docker.io/library/alpine:latest"),
PathBuf::from("/tmp"),
)
.await;
orch.upgrade("staged-music").await.unwrap();
assert_eq!(
orch.staged_upgrade_image("staged-music").await.unwrap(),
Some(image)
);
assert!(!rt.calls().iter().any(|c| c.starts_with("start_container:")));
}
#[tokio::test]
async fn failed_staged_start_preserves_pin_and_stop_until_retry_succeeds() {
let (rt, orch, _) = stopped_update_fixture().await;
orch.upgrade("staged-music").await.unwrap();
rt.fail_start
.lock()
.unwrap()
.insert("staged-music".into(), "injected start failure".into());
assert!(orch.start("staged-music").await.is_err());
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_some());
assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir)
.await
.contains("staged-music"));
assert!(matches!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Stopped | ContainerState::Created
));
orch.start("staged-music").await.unwrap();
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_none());
}
#[tokio::test]
async fn staged_missing_image_uninstall_and_damaged_record_refuse_start() {
let (rt, orch, image) = stopped_update_fixture().await;
orch.upgrade("staged-music").await.unwrap();
rt.images.lock().unwrap().remove(&image);
assert!(orch
.start("staged-music")
.await
.unwrap_err()
.to_string()
.contains("unavailable"));
rt.mark_image_present(&image);
crate::crash_recovery::mark_user_uninstalled(&orch.data_dir, "staged-music").await;
assert!(orch
.start("staged-music")
.await
.unwrap_err()
.to_string()
.contains("uninstalled"));
assert!(orch.upgrade("staged-music").await.is_err());
tokio::fs::write(
orch.data_dir.join("staged-updates/staged-music.json"),
b"damaged",
)
.await
.unwrap();
assert!(orch.start("staged-music").await.is_err());
assert!(!rt.calls().iter().any(|c| c.starts_with("start_container:")));
}
#[tokio::test]
async fn staged_failed_post_install_hook_keeps_pin_and_stops_container() {
let (rt, orch, _) = stopped_update_fixture().await;
let mut lm = orch.loaded("staged-music").await.unwrap();
lm.manifest.app.hooks.post_install = serde_yaml::from_str(
"- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n",
)
.unwrap();
orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir)
.await;
orch.upgrade("staged-music").await.unwrap();
assert!(orch.start("staged-music").await.is_err());
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_some());
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Stopped
);
assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir)
.await
.contains("staged-music"));
}
#[tokio::test]
async fn failed_staged_cleanup_reports_active_runtime_and_reconcile_retries_stop() {
let (rt, orch, _) = stopped_update_fixture().await;
let mut lm = orch.loaded("staged-music").await.unwrap();
lm.manifest.app.hooks.post_install = serde_yaml::from_str(
"- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n",
)
.unwrap();
orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir)
.await;
orch.upgrade("staged-music").await.unwrap();
rt.fail_stop
.lock()
.unwrap()
.insert("staged-music".into(), "stop refused".into());
let error = orch.start("staged-music").await.unwrap_err();
assert!(error.downcast_ref::<StagedCleanupFailure>().is_some());
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Running
);
let report = orch.reconcile_all().await;
assert!(report.failures.is_empty());
assert!(report
.actions
.iter()
.any(|(app, action)| app == "staged-music"
&& action == &ReconcileAction::Left("staged-update-awaiting-start".into())));
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Stopped
);
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_some());
}
#[tokio::test]
async fn pending_stage_enforces_stop_after_restart_when_stop_marker_cannot_be_saved() {
let (rt, orch, _) = stopped_update_fixture().await;
let mut lm = orch.loaded("staged-music").await.unwrap();
lm.manifest.app.hooks.post_install = serde_yaml::from_str(
"- copy_from_host:\n src: nonexistent-staged-hook-file\n dest: /tmp/test\n",
)
.unwrap();
orch.insert_manifest_for_test(lm.manifest.clone(), lm.manifest_dir.clone())
.await;
orch.upgrade("staged-music").await.unwrap();
let marker = orch.data_dir.join("user-stopped.json");
tokio::fs::remove_file(&marker).await.unwrap();
tokio::fs::create_dir(&marker).await.unwrap();
rt.fail_stop
.lock()
.unwrap()
.insert("staged-music".into(), "stop interrupted".into());
let failure = orch.start("staged-music").await.unwrap_err();
assert!(failure.downcast_ref::<StagedCleanupFailure>().is_some());
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Running
);
let starts = rt
.calls()
.iter()
.filter(|c| c.starts_with("start_container:"))
.count();
let mut resumed = orch_with(rt.clone()).await;
resumed.set_data_dir(orch.data_dir.clone());
resumed
.insert_manifest_for_test(lm.manifest, lm.manifest_dir)
.await;
let report = resumed.reconcile_all().await;
assert!(report.failures.is_empty());
assert_eq!(
rt.get_container_status("staged-music").await.unwrap().state,
ContainerState::Stopped
);
assert_eq!(
rt.calls()
.iter()
.filter(|c| c.starts_with("start_container:"))
.count(),
starts
);
assert!(resumed
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_some());
}
#[tokio::test]
async fn staged_configuration_failure_keeps_stop_and_pin_before_create() {
let (rt, mut orch, _) = stopped_update_fixture().await;
orch.secrets_dir = orch.data_dir.join("fixture-secrets");
let mut lm = orch.loaded("staged-music").await.unwrap();
lm.manifest.app.container.secret_env =
serde_yaml::from_str("- key: REQUIRED_SECRET\n secret_file: missing-staged-secret\n")
.unwrap();
orch.insert_manifest_for_test(lm.manifest, lm.manifest_dir)
.await;
orch.upgrade("staged-music").await.unwrap();
assert!(orch.start("staged-music").await.is_err());
assert!(orch
.staged_upgrade_image("staged-music")
.await
.unwrap()
.is_some());
assert!(crate::crash_recovery::load_user_stopped(&orch.data_dir)
.await
.contains("staged-music"));
assert!(!rt
.calls()
.iter()
.any(|c| c.starts_with("create_container:")));
}
#[tokio::test]
async fn staged_record_persistence_failure_precedes_runtime_changes() {
let (rt, orch, _) = stopped_update_fixture().await;
tokio::fs::write(orch.data_dir.join("staged-updates"), b"blocked directory")
.await
.unwrap();
assert!(orch.upgrade("staged-music").await.is_err());
assert!(!rt.calls().iter().any(|c| c.starts_with("pull_image:")
|| c.starts_with("stop_container:")
|| c.starts_with("remove_container:")
|| c.starts_with("create_container:")));
}
#[tokio::test]
async fn stopped_unpinned_update_refuses_without_starting() {
let (rt, orch, _) = stopped_update_fixture().await;
orch.insert_manifest_for_test(
pull_manifest("staged-music", "docker.io/library/alpine:3.20"),
PathBuf::from("/tmp"),
)
.await;
assert!(orch
.upgrade("staged-music")
.await
.unwrap_err()
.to_string()
.contains("digest-pinned"));
assert!(!rt.calls().iter().any(|c| c.starts_with("pull_image:")
|| c.starts_with("create_container:")
|| c.starts_with("start_container:")));
}
#[tokio::test]
async fn install_fresh_pull() {
let rt = Arc::new(MockRuntime::default());
@@ -0,0 +1,314 @@
//! Durable, immutable preparation for updates that must remain stopped.
use anyhow::{Context, Result};
use archipelago_container::AppManifest;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::{
collections::HashSet,
path::{Path, PathBuf},
};
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(super) struct StagedUpdate {
pub manifest: AppManifest,
pub manifest_dir: PathBuf,
pub ready: bool,
}
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Envelope {
schema: u32,
payload: String,
checksum: String,
}
impl StagedUpdate {
pub fn image(&self) -> Result<&str> {
anyhow::ensure!(
self.manifest.app.container.build.is_none(),
"Stopped updates of build-based apps require immutable image resolution"
);
let image = self
.manifest
.app
.container
.image
.as_deref()
.context("Missing staged image")?;
let digest = image.rsplit_once("@sha256:").map(|(_, d)| d);
anyhow::ensure!(
digest.is_some_and(|d| d.len() == 64 && d.bytes().all(|b| b.is_ascii_hexdigit())),
"Stopped update requires a digest-pinned image; app remains stopped"
);
Ok(image)
}
}
fn record_path(root: &Path, directory: &str, app: &str) -> Result<PathBuf> {
anyhow::ensure!(
!app.is_empty()
&& app.len() <= 128
&& app
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_')),
"Invalid staged app id"
);
Ok(root.join(directory).join(format!("{app}.json")))
}
fn path(root: &Path, app: &str) -> Result<PathBuf> {
record_path(root, "staged-updates", app)
}
pub(super) async fn load(root: &Path, app: &str) -> Result<Option<StagedUpdate>> {
load_at(root, "staged-updates", app).await
}
async fn load_at(root: &Path, directory: &str, app: &str) -> Result<Option<StagedUpdate>> {
let bytes = match tokio::fs::read(record_path(root, directory, app)?).await {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(e.into()),
};
let envelope: Envelope = serde_json::from_slice(&bytes)
.context("Damaged staged update; refusing another version")?;
anyhow::ensure!(envelope.schema == 1, "Unsupported staged update schema");
anyhow::ensure!(
envelope.checksum == hex::encode(Sha256::digest(envelope.payload.as_bytes())),
"Staged update checksum mismatch"
);
let record: StagedUpdate =
serde_json::from_str(&envelope.payload).context("Damaged staged manifest")?;
record
.manifest
.validate()
.context("Invalid staged manifest")?;
anyhow::ensure!(
record.manifest.app.id == app,
"Staged update belongs to another app"
);
record.image()?;
Ok(Some(record))
}
pub(super) async fn save(root: &Path, record: &StagedUpdate) -> Result<()> {
save_at(root, "staged-updates", record).await
}
pub(super) async fn save_installed(root: &Path, record: &StagedUpdate) -> Result<()> {
anyhow::ensure!(record.ready, "Cannot publish incomplete installed manifest");
save_at(root, "installed-manifests", record).await
}
pub(super) async fn installed_manifest(
root: &Path,
app: &str,
observed_image: &str,
) -> Result<Option<AppManifest>> {
let Some(record) = load_at(root, "installed-manifests", app).await? else {
return Ok(None);
};
let expected = record.image()?;
let matches = matches_image(expected, observed_image);
Ok((record.ready && matches).then_some(record.manifest))
}
pub(super) fn matches_image(expected: &str, observed: &str) -> bool {
expected == observed
|| match (
expected.rsplit_once("@sha256:"),
observed.rsplit_once("@sha256:"),
) {
(Some((_, a)), Some((_, b))) => a.len() == 64 && a.eq_ignore_ascii_case(b),
_ => false,
}
}
async fn save_at(root: &Path, directory: &str, record: &StagedUpdate) -> Result<()> {
use std::{
io::Write,
os::unix::fs::{DirBuilderExt, OpenOptionsExt, PermissionsExt},
};
record.image()?;
let target = record_path(root, directory, &record.manifest.app.id)?;
record
.manifest
.validate()
.context("Invalid staged manifest")?;
let payload = serde_json::to_string(record)?;
let checksum = hex::encode(Sha256::digest(payload.as_bytes()));
let bytes = serde_json::to_vec(&Envelope {
schema: 1,
payload,
checksum,
})?;
// Bounded metadata commit stays in this poll while the caller owns its
// lifecycle lock; cancellation cannot leave a late detached writer.
(|| -> Result<()> {
let parent = target.parent().context("Missing staged update directory")?;
match std::fs::DirBuilder::new().mode(0o700).create(parent) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {}
Err(e) => return Err(e.into()),
}
std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700))?;
// Persist the directory entry itself before an acknowledged record can
// depend on it surviving a crash (syncing the child alone is insufficient).
std::fs::File::open(parent.parent().context("Missing staged update parent")?)?
.sync_all()?;
let tmp = parent.join(format!(".{}.tmp", uuid::Uuid::new_v4()));
let result = (|| -> Result<()> {
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(&tmp)?;
file.write_all(&bytes)?;
file.sync_all()?;
std::fs::rename(&tmp, &target)?;
std::fs::File::open(parent)?.sync_all()?;
Ok(())
})();
if result.is_err() {
let _ = std::fs::remove_file(tmp);
}
result
})()
}
pub(super) async fn clear(root: &Path, app: &str) -> Result<()> {
clear_at(root, "staged-updates", app).await
}
pub(super) async fn clear_installed(root: &Path, app: &str) -> Result<()> {
clear_at(root, "installed-manifests", app).await
}
async fn clear_at(root: &Path, directory: &str, app: &str) -> Result<()> {
let target = record_path(root, directory, app)?;
// Bounded metadata commit stays in this poll while the caller owns its
// lifecycle lock; cancellation cannot leave a late detached writer.
(|| -> Result<()> {
match std::fs::remove_file(&target) {
Ok(()) => std::fs::File::open(target.parent().unwrap())?.sync_all()?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(e.into()),
}
Ok(())
})()
}
pub(super) async fn marked(root: &Path, file: &str, app: &str, container: &str) -> Result<bool> {
let bytes = match tokio::fs::read(root.join(file)).await {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(false),
Err(e) => return Err(e).with_context(|| format!("Cannot read {file}")),
};
let set: HashSet<String> =
serde_json::from_slice(&bytes).with_context(|| format!("Damaged {file}"))?;
Ok(set.contains(app) || set.contains(container) || set.contains(&format!("archy-{app}")))
}
#[cfg(test)]
mod tests {
use super::*;
use futures_util::FutureExt;
use std::os::unix::fs::PermissionsExt;
#[tokio::test]
async fn envelope_rejects_corruption_schema_and_invalid_manifest_and_is_private() {
let root = tempfile::tempdir().unwrap();
let manifest = AppManifest::parse(&format!("app:\n id: music\n name: Music\n version: 1.0.0\n container:\n image: docker.io/library/alpine@sha256:{}\n", "a".repeat(64))).unwrap();
let record = StagedUpdate {
manifest,
manifest_dir: root.path().into(),
ready: true,
};
save(root.path(), &record)
.now_or_never()
.expect("Journal commit must finish before lifecycle lock cancellation is possible")
.unwrap();
let target = path(root.path(), "music").unwrap();
assert_eq!(
std::fs::metadata(target.parent().unwrap())
.unwrap()
.permissions()
.mode()
& 0o777,
0o700
);
assert_eq!(
std::fs::metadata(&target).unwrap().permissions().mode() & 0o777,
0o600
);
assert!(load(root.path(), "music").await.unwrap().unwrap().ready);
save_installed(root.path(), &record).await.unwrap();
clear(root.path(), "music")
.now_or_never()
.expect("Journal clear must not leave a detached late deletion")
.unwrap();
assert!(load(root.path(), "music").await.unwrap().is_none());
assert!(
installed_manifest(root.path(), "music", record.image().unwrap())
.await
.unwrap()
.is_some()
);
assert!(installed_manifest(
root.path(),
"music",
&format!("docker.io/library/alpine@sha256:{}", "b".repeat(64))
)
.await
.unwrap()
.is_none());
save(root.path(), &record).await.unwrap();
let original = tokio::fs::read(&target).await.unwrap();
let mut envelope: Envelope = serde_json::from_slice(&original).unwrap();
envelope.payload = envelope.payload.replace("Music", "Changed Music");
tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap())
.await
.unwrap();
assert!(load(root.path(), "music")
.await
.unwrap_err()
.to_string()
.contains("checksum"));
envelope = serde_json::from_slice(&original).unwrap();
envelope.schema = 2;
tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap())
.await
.unwrap();
assert!(load(root.path(), "music").await.is_err());
envelope = serde_json::from_slice(&original).unwrap();
let mut value: serde_json::Value = serde_json::from_str(&envelope.payload).unwrap();
value["manifest"]["app"]["container"]["image"] = serde_json::Value::Null;
envelope.payload = serde_json::to_string(&value).unwrap();
envelope.checksum = hex::encode(Sha256::digest(envelope.payload.as_bytes()));
tokio::fs::write(&target, serde_json::to_vec(&envelope).unwrap())
.await
.unwrap();
assert!(load(root.path(), "music").await.is_err());
}
}
pub(super) async fn pending_ids(root: &Path) -> Result<Vec<String>> {
let mut dir = match tokio::fs::read_dir(root.join("staged-updates")).await {
Ok(dir) => dir,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(vec![]),
Err(e) => return Err(e.into()),
};
let mut ids = Vec::new();
while let Some(entry) = dir.next_entry().await? {
if let Some(name) = entry
.file_name()
.to_str()
.and_then(|name| name.strip_suffix(".json"))
{
ids.push(name.to_string());
}
}
Ok(ids)
}
@@ -0,0 +1,672 @@
//! 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, Completion, PreparedTarget, RecoveryImage, Supervisor, Unit},
update_transaction::{Observed, Podman, Runtime, Target},
};
use anyhow::{Context, Result};
use sha2::{Digest, Sha256};
use std::{
collections::HashMap,
future::Future,
io::Write,
os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt},
path::{Path, PathBuf},
time::Duration,
};
pub(crate) trait DrainBarrier: Sync {
/// 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,
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)]
#[serde(deny_unknown_fields)]
pub(crate) struct MigrationPlan {
original_sha256: String,
prepared: PreparedTarget,
}
impl MigrationPlan {
/// The old render must be reproducible from an image-bound retained original
/// manifest. A current target catalog is not an installed-version receipt.
pub(crate) fn from_installed_render(
original: &str,
old_render: &str,
target_render: &str,
manifest: archipelago_container::AppManifest,
) -> Result<Self> {
manifest.validate()?;
let body = supervised_update::merge_reviewed_unit(original, old_render, target_render)?;
Ok(Self {
original_sha256: hex::encode(Sha256::digest(original.as_bytes())),
prepared: PreparedTarget { body, manifest },
})
}
/// Legacy overrides whose provenance cannot be reconstructed need a reviewed
/// exact-original-hash-bound migration, prepared by node administration code.
/// No browser RPC accepts unit bodies or this type.
pub(crate) fn from_reviewed_legacy(
original_sha256: &str,
body: String,
manifest: archipelago_container::AppManifest,
) -> Result<Self> {
anyhow::ensure!(
original_sha256.len() == 64 && original_sha256.bytes().all(|v| v.is_ascii_hexdigit()),
"Invalid original unit commitment"
);
manifest.validate()?;
Ok(Self {
original_sha256: original_sha256.into(),
prepared: PreparedTarget { body, manifest },
})
}
}
#[derive(serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct PlanFile {
schema: u8,
package: String,
plans: HashMap<String, MigrationPlan>,
}
/// Administrative preparation is separate from the Update RPC. The RPC never
/// accepts unit bodies, hook commands, or a browser-chosen preparation path.
pub(crate) fn load_reviewed_plans(
data_dir: &Path,
package: &str,
) -> Result<HashMap<String, MigrationPlan>> {
anyhow::ensure!(
!package.is_empty()
&& package.len() <= 128
&& package
.bytes()
.all(|v| v.is_ascii_alphanumeric() || matches!(v, b'-' | b'_')),
"Invalid managed package"
);
let dir = data_dir.join("managed-update-plans");
let directory = std::fs::symlink_metadata(&dir)
.context("No original-bound managed update plan; existing app remains unchanged")?;
anyhow::ensure!(
directory.is_dir()
&& !directory.file_type().is_symlink()
&& directory.uid() == unsafe { libc::geteuid() }
&& directory.mode() & 0o077 == 0,
"Managed update plan directory must be private and node-owned"
);
let path = dir.join(format!("{package}.json"));
let metadata = std::fs::symlink_metadata(&path)?;
anyhow::ensure!(
metadata.is_file()
&& !metadata.file_type().is_symlink()
&& metadata.uid() == unsafe { libc::geteuid() }
&& metadata.mode() & 0o077 == 0
&& metadata.len() <= 4 * 1024 * 1024,
"Managed update plan must be private and node-owned"
);
let record: PlanFile = serde_json::from_slice(&std::fs::read(path)?)?;
anyhow::ensure!(
record.schema == 1
&& record.package == package
&& !record.plans.is_empty()
&& record.plans.len() <= 32,
"Managed update plan does not match this package"
);
for plan in record.plans.values() {
anyhow::ensure!(
plan.original_sha256.len() == 64
&& plan.original_sha256.bytes().all(|v| v.is_ascii_hexdigit()),
"Invalid original unit commitment"
);
plan.prepared.manifest.validate()?;
}
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,
plans: HashMap<String, MigrationPlan>,
barrier: B,
}
impl<B: DrainBarrier> SystemdSupervisor<B> {
pub(crate) async fn new(
data_dir: PathBuf,
plans: HashMap<String, MigrationPlan>,
barrier: B,
) -> Result<Self> {
let unit_dir = super::quadlet::unit_dir().await?;
let meta = std::fs::symlink_metadata(&unit_dir)?;
anyhow::ensure!(
meta.is_dir()
&& !meta.file_type().is_symlink()
&& meta.uid() == unsafe { libc::geteuid() }
&& meta.mode() & 0o022 == 0,
"Quadlet directory ownership changed"
);
Ok(Self {
data_dir,
unit_dir,
plans,
barrier,
})
}
/// A reviewed plan may replace a mutable catalog spelling with its exact
/// locally verified digest, but never select different image bytes.
pub(crate) async fn reviewed_targets(&self, targets: &[Target]) -> Result<Vec<Target>> {
let mut refs = Vec::new();
for target in targets {
let plan = self
.plans
.get(&target.name)
.context("Missing reviewed managed target")?;
let reference = plan
.prepared
.manifest
.app
.container
.image
.as_ref()
.context("Reviewed image missing")?;
anyhow::ensure!(
reference.rsplit_once("@sha256:").is_some_and(
|(_, hash)| hash.len() == 64 && hash.bytes().all(|b| b.is_ascii_hexdigit())
),
"Managed target must use an immutable digest"
);
refs.push((target.name.clone(), reference.clone()));
}
let resolved = Podman::targets(&refs, false).await?;
anyhow::ensure!(
resolved
.iter()
.zip(targets)
.all(|(reviewed, catalog)| reviewed.name == catalog.name
&& reviewed.image == catalog.image),
"Reviewed managed image differs from prepared catalog image"
);
Ok(resolved)
}
fn path(&self, name: &str) -> Result<PathBuf> {
anyhow::ensure!(
!name.is_empty()
&& name.len() <= 128
&& name
.bytes()
.all(|v| v.is_ascii_alphanumeric() || matches!(v, b'-' | b'_')),
"Invalid supervised unit name"
);
Ok(self.unit_dir.join(format!("{name}.container")))
}
async fn manager(args: &[&str]) -> Result<String> {
let output = tokio::time::timeout(
Duration::from_secs(120),
tokio::process::Command::new("systemctl")
.arg("--user")
.args(args)
.kill_on_drop(true)
.output(),
)
.await
.context("User service manager timed out")??;
anyhow::ensure!(
output.status.success(),
"User service manager rejected operation"
);
Ok(String::from_utf8(output.stdout)?)
}
async fn owned_file(&self, name: &str) -> Result<(PathBuf, std::fs::Metadata)> {
let expected = self.path(name)?;
let service = format!("{name}.service");
let output = Self::manager(&[
"show",
&service,
"--property=SourcePath",
"--property=DropInPaths",
"--property=LoadState",
])
.await?;
let values: HashMap<_, _> = output
.lines()
.filter_map(|line| line.split_once('='))
.collect();
anyhow::ensure!(
values.get("LoadState") == Some(&"loaded")
&& values.get("DropInPaths") == Some(&"")
&& values
.get("SourcePath")
.is_some_and(|path| Path::new(path) == expected),
"Service is not owned by the exact original source Quadlet or has external overrides"
);
let meta = std::fs::symlink_metadata(&expected)?;
anyhow::ensure!(
meta.is_file()
&& !meta.file_type().is_symlink()
&& meta.uid() == unsafe { libc::geteuid() }
&& meta.mode() & 0o022 == 0
&& meta.len() <= 1024 * 1024,
"Original Quadlet ownership changed"
);
Ok((expected, meta))
}
}
impl<B: DrainBarrier> Supervisor for SystemdSupervisor<B> {
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, outcome: Completion) -> Result<()> {
self.barrier.release(operation, outcome).await
}
async fn prepare_target(&self, target: &Target, original: &Unit) -> Result<PreparedTarget> {
let plan = self
.plans
.get(&target.name)
.context("No reviewed original-bound managed migration; app unchanged")?;
anyhow::ensure!(
hex::encode(Sha256::digest(original.body.as_bytes())) == plan.original_sha256
&& plan.prepared.manifest.app.container.image.as_deref()
== Some(target.reference.as_str())
&& super::prod_orchestrator::compute_container_name(&plan.prepared.manifest)
== target.name,
"Original unit or reviewed immutable target changed before update"
);
if plan
.prepared
.manifest
.app
.container
.media_registration_identity
{
let identity =
crate::identity::NodeIdentity::load_existing(&self.data_dir.join("identity"))
.await?;
let pin = super::registration_pin::load_existing(
&self.data_dir,
&plan.prepared.manifest.app.id,
&identity,
)?;
let mut expected = plan.prepared.manifest.clone();
super::registration_pin::apply_environment(&mut expected, &pin)?;
for entry in expected
.app
.environment
.iter()
.filter(|entry| entry.starts_with("ARCHIPELAGO_REGISTRATION_"))
{
let key = entry
.split_once('=')
.context("Malformed registration binding")?
.0;
let in_manifest: Vec<_> = plan
.prepared
.manifest
.app
.environment
.iter()
.filter(|value| value.split_once('=').is_some_and(|(name, _)| name == key))
.collect();
let in_unit: Vec<_> = plan
.prepared
.body
.lines()
.filter_map(|line| line.trim().strip_prefix("Environment="))
.map(|value| value.trim_matches('"'))
.filter(|value| value.split_once('=').is_some_and(|(name, _)| name == key))
.collect();
anyhow::ensure!(
in_manifest.len() == 1
&& in_manifest[0] == entry
&& in_unit.len() == 1
&& in_unit[0] == entry,
"Reviewed registration identity does not match the existing installer pin"
);
}
}
Ok(plan.prepared.clone())
}
async fn target_hooks(
&self,
name: &str,
manifest: &archipelago_container::AppManifest,
) -> Result<()> {
super::hooks::run_post_install_strict(manifest, name, &self.data_dir).await
}
async fn capture(&self, name: &str) -> Result<Unit> {
let (path, meta) = self.owned_file(name).await?;
let body = std::fs::read_to_string(path)?;
let observed = Podman
.inspect(name)
.await?
.context("Original supervised runtime missing")?;
Ok(Unit {
name: name.into(),
body,
image: observed.image,
container_id: observed.id,
file_mode: meta.mode() & 0o777,
running: observed.running,
config_sha256: observed.config_sha256,
})
}
async fn validate_original_file(&self, original: &Unit) -> Result<()> {
let (_, meta) = self.owned_file(&original.name).await?;
anyhow::ensure!(
meta.mode() & 0o777 == original.file_mode,
"Original Quadlet mode changed"
);
Ok(())
}
async fn read(&self, name: &str) -> Result<String> {
let (path, _) = self.owned_file(name).await?;
Ok(std::fs::read_to_string(path)?)
}
async fn write(&self, original: &Unit, expected: &[String], body: &str) -> Result<()> {
let (path, meta) = self.owned_file(&original.name).await?;
anyhow::ensure!(
meta.mode() & 0o777 == original.file_mode
&& expected.contains(&std::fs::read_to_string(&path)?),
"Original Quadlet changed; refusing to replace a foreign edit"
);
let temporary = self
.unit_dir
.join(format!(".archy-update-{}.tmp", uuid::Uuid::new_v4()));
// Commit point is synchronous while the transaction guard is held;
// cancelled futures cannot rename over a later recovery after unlock.
let result = (|| -> Result<()> {
let mut file = std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.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()?;
Ok(())
})();
if result.is_err() {
let _ = std::fs::remove_file(&temporary);
}
result
}
async fn snapshot(&self, original: &Unit, operation: &str, tag: &str) -> Result<RecoveryImage> {
// 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<()> {
anyhow::ensure!(
image.len() == 64
&& image.bytes().all(|v| v.is_ascii_hexdigit())
&& tag.starts_with("localhost/archy-update-recovery:"),
"Invalid recovery image pin"
);
let output = tokio::time::timeout(
Duration::from_secs(30),
tokio::process::Command::new("podman")
.args(["tag", &format!("sha256:{image}"), tag])
.kill_on_drop(true)
.output(),
)
.await??;
anyhow::ensure!(
output.status.success(),
"Original recovery image unavailable"
);
Ok(())
}
async fn stop(&self, name: &str) -> Result<()> {
self.owned_file(name).await?;
super::quadlet::stop_service_with_timeout(
&format!("{name}.service"),
Duration::from_secs(archipelago_container::runtime::stop_grace_secs_for(name) + 30),
)
.await
}
async fn reload(&self) -> Result<()> {
super::quadlet::daemon_reload_user().await
}
async fn start(&self, name: &str) -> Result<()> {
self.owned_file(name).await?;
super::quadlet::enable_now(&format!("{name}.service")).await
}
async fn observed(&self, name: &str) -> Result<Option<Observed>> {
Podman.inspect(name).await
}
async fn healthy(&self, name: &str) -> Result<bool> {
Podman.healthy(name).await
}
}
#[cfg(test)]
mod tests {
use super::*;
struct UnusedBarrier;
impl DrainBarrier for UnusedBarrier {
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, _: Completion) -> Result<()> {
Ok(())
}
}
#[tokio::test]
async fn original_bound_plan_rejects_changed_unit_or_image_without_service_calls() {
let image = format!("localhost/movie@sha256:{}", "b".repeat(64));
let manifest = archipelago_container::AppManifest::parse(&format!(
"app:\n id: movie\n name: Movie\n version: 2.0.0\n container:\n image: {image}\n")).unwrap();
let original_body = "[Container]\nContainerName=movie\nImage=old\n";
let next_body = format!("[Container]\nContainerName=movie\nImage={image}\n");
let plan = MigrationPlan::from_reviewed_legacy(
&hex::encode(Sha256::digest(original_body)),
next_body.clone(),
manifest,
)
.unwrap();
let adapter = SystemdSupervisor {
data_dir: PathBuf::from("/unused"),
unit_dir: PathBuf::from("/unused"),
plans: HashMap::from([("movie".into(), plan)]),
barrier: UnusedBarrier,
};
let mut original = Unit {
name: "movie".into(),
body: original_body.into(),
image: "a".repeat(64),
container_id: "c".repeat(64),
file_mode: 0o600,
running: true,
config_sha256: "d".repeat(64),
};
let mut target = Target {
name: "movie".into(),
reference: image,
image: "b".repeat(64),
};
assert_eq!(
adapter
.prepare_target(&target, &original)
.await
.unwrap()
.body,
next_body
);
original.body.push_str("Environment=OPERATOR=changed\n");
assert!(adapter.prepare_target(&target, &original).await.is_err());
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], false)
.await
.is_err());
}
}
File diff suppressed because it is too large Load Diff
+5
View File
@@ -62,6 +62,11 @@ pub trait ContainerOrchestrator: Send + Sync {
/// Pull/rebuild the image and recreate the container from scratch.
async fn upgrade(&self, app_id: &str) -> Result<()>;
/// Exact image prepared by a stopped upgrade; no running container is claimed.
async fn staged_upgrade_image(&self, _app_id: &str) -> Result<Option<String>> {
Ok(None)
}
/// Current state of a single container.
async fn status(&self, app_id: &str) -> Result<ContainerStatus>;
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,114 @@
# Legacy IndeeHub maintenance controller
Status: isolated source implementation. Fourteen pure Python fake-runtime regressions
pass; no live invocation or production qualification. The controller is not part
of the already signed private app candidate and needs no new app image/API.
The supervised updater owns the lifecycle flock, seven durable holds, original
Quadlets and private writable-layer recovery images. It records its destructive
obligation before invoking the fixed controller with bounded JSON over stdin:
```
python3 /opt/archipelago/scripts/indeehub-maintenance-controller.py acquire
{ "operation_id": "<uuid>", "original_members": [
{ "name": "indeedhub", "container_id": "<64hex>", "image_id": "<64hex>",
"unit_sha256": "<64hex>", "config_sha256": "<64hex>", "running": true }
// All seven exact members; JSON does not include this illustrative comment.
], "recovery": false }
```
Other actions are `verify` with operation_id, and `release` with operation_id and
outcome committed/restored/aborted. Replies are <=4KiB and report drained, held,
released, or recovering for explicit recovery acquire. Inherited
ARCHY_UPDATE_LOCK_FD stays open and is passed to child commands; the script never
unlocks it. Journal: data/update-transactions/indeehub-maintenance/<uuid>/journal.json.
## Forward sequence
- Validate exact original IDs/source-unit hashes and known port exposure. Only
frontend127.0.0.1:7778 is supported; direct backend/S3 ports refuse before stop.
- Require deployed native AppGate and legacy nginx maintenance guards. Inspect
every known legacy sublocation and any direct7778 proxy; unknown routes refuse.
Save an operation-owned readable sentinel, then verify local ingress returns503.
- Record prior BullMQ transcode pause state, globally pause future job admission,
retain queued/delayed/failed jobs. Gracefully stop frontend ingress; bounded
polling waits for active transcodes to finish before stopping worker and API.
- Require successful systemd shutdown plus an exact original Podman died event
with exit0. Forced exits and missing event evidence retain the hold and are
never labelled completed writes.
- While PostgreSQL remains running, capture a fresh custom dump. Cleanly stop
MinIO/Redis/relay/Postgres, then archive all four complete quiescent volumes
(including SQLite WAL and Redis persistence) with metadata. No volume deletion
or migration rollback. Archive hashes/size and per-step obligations are durable.
- Keep admission closed while the updater renders, starts and verifies targets.
## Interrupted recovery
The node first records phase Restoring with boolean target_startup_began, then
calls acquire with recovery:true. That path preserves the original failure and
fence; it does not retry a killed original into a fictitious successful drain or
claim missing backups exist. The node restores exact saved old runtime under the
same hold. Release before any target startup can state only that original runtime
was restored. If target startup/migration began, recorded data-compatibility
verification is required before restored release; an old image alone does not
prove compatibility with newly changed data. No automatic DB/media restore exists.
## Qualification and remaining integration
`python3 tests/regression/test_indeehub_maintenance_controller.py` passes fourteen
fake-runtime cases in temporary directories, without services/network/containers.
Source nginx template guard coverage also passes its parser check. Production
adapter compilation, actual Podman event format/systemd clean-exit behavior,
application writer shutdown, interrupted backup and supervised restart still need
isolated lifecycle fixtures and then coordinated node acceptance. A long-lived
WebSocket or active upload can exceed graceful-stop deadlines; the current code
refuses completion and preserves recovery obligations rather than silently
calling interrupted work finished.
The deployment must install the exact qualified controller script and record its
hash alongside the backend artifact. Binary-only deployment does not install it.
The backend must refuse missing/mismatched prerequisites before snapshots/stops.
Native AppGate + nginx guards are separate node source changes owned by the
supervised updater agent. The signed app catalog/private image receipts remain
unchanged. Existing live stop/uninstall intent must not be rewritten as maintenance.
A pre-acquire snapshot/preflight failure may leave no controller journal. An
Aborted node journal with target_startup_began=false then permits idempotent
no-op acknowledgement, without touching any other operation’s admission fence.
A matching fence without its controller journal requires recovery investigation.
Read-only source evidence from actual old API/ffmpeg shows neither has SIGTERM
shutdown hooks. The controller permits worker143 only after a paused queue has
zero active jobs. Legacy API143 additionally requires closed/stopped frontend,
stopped worker, and a fresh empty projects/contents/payments/shareholders/
subscriptions/library_items store with no other active DB transaction. This is
a narrow first-upgrade compatibility path, not evidence populated work completed.
Populated or ambiguous legacy state remains a refused forward cutover.
### Operation-bound rollback data verification
Restored release after any target startup now performs its own PostgreSQL
compatibility check; it does not accept a manually asserted verification boolean.
Before the coherent backup, the controller captures a read-only, repeatable-read
transaction containing every original public table's columns, constraints,
indexes, triggers, row-security policies, row count and canonical row SHA-256,
plus the exact migration history. The private operation journal binds this
baseline to the original operation UUID.
After original runtime recovery, while ingress and worker admission remain
closed, the controller captures the same observations again. It requires original
tables and definitions unchanged, original non-migration rows unchanged and the
original migration-history prefix intact. Additional migration records must be
the exact ordered three migrations qualified for API commit `3b09b81`, and only
their five named new tables may appear, all empty. Any unexpected data or schema
change keeps ingress closed. Successful proof records before/after commitment
hashes and the operation UUID before release. No down migration or automatic
volume restoration is performed.
Sixteen pure Python regressions pass, including altered rows/schema/history,
foreign operation, nonempty added tables and durable proof before fence release.
The SQL transaction and Podman lifecycle still require isolated integration
qualification. These are table-level compatibility checks, not a claim that
arbitrary database extensions/functions, other writers, or changed application
code are safe. The candidate images, migration scope and admission barrier must
also match the reviewed operation.
@@ -0,0 +1,77 @@
# Managed update runtime recovery
Status: isolated source implementation; Rust and real Podman/systemd qualification
pending. No live update, snapshot, stop, backup or rollback has been performed by
this work. Active deployed source is unchanged.
The managed path captures original source Quadlet bytes, mode, immutable image,
container identity, launch configuration and running intent. It requires an
original-hash-bound reviewed forward plan and verifies every planned image against
the already prepared catalog image. New manifest configuration/hooks are applied
on the forward path; rollback uses the captured original recipe and a private,
local-only writable-layer image. AutoRemove rollback recreates containers and does
not claim to restore their original IDs. Stopped supervised stacks currently fail
before mutation; the retained-container and separate stopped-stage paths cover
only their respective supported cases.
The legacy IndeeHub controller runs under the same inherited lifecycle lock and
operation-owned reconciliation holds. Original writable layers are captured
before any destructive stop. The controller fences ingress, drains supported
legacy work, takes coherent quiescent volume/database backups and retains the
fence through cutover or recovery. Its exact source hash must match the separately
installed script; a backend binary alone does not install the controller.
Completed updates and verified runtime restorations publish exact unit recipes
before releasing holds. Routine drift reconciliation validates those recipes;
it does not regenerate them from a newer catalog. Catalog-driven pre-start file
and mount mutations are skipped for these pinned installations. Missing saved
units/images are explicit recovery failures, never permission to reconstruct a
different runtime. Explicit uninstall removes the installed recipe, and completed
old journals cannot recreate it. A new reviewed transaction replaces the recipe.
Opted-in API registration environments must match an already provisioned pin
and the existing node identity in both manifest and exact Quadlet. Administrative
plan preparation must use the existing installer provisioning code. Execution
never invents a node identity or accepts a browser-supplied unit or hook.
Runtime restoration does not establish database compatibility. The controller's
restored-release verifier compares original table schemas and row commitments,
permitting only the exact reviewed additive migration prefix and empty new
application tables. Any other data/schema change keeps ingress closed. This is
not automatic database rollback or a promise that arbitrary migrations are
reversible.
Qualification required before integration/activation:
- Isolated backend compile and injectable lifecycle/fault tests, including lost
replies, daemon interruption, foreign units/holds and preflight failures.
- Disposable real Podman/systemd and PostgreSQL execution of the controller and
adapter. Sixteen pure controller tests currently pass; SQL/runtime behavior is
not yet qualified.
- Final seven-member private plan with installer-resolved identity environment,
verified local images and source-unit provenance; review required changes and
retained operator configuration without exposing secret values.
- Disk-capacity and recovery-image retention checks, backup integrity and a
documented recovery path for missing runtime artifacts.
- Actual-node controlled deployment and acceptance, preserving persistent data.
## Resumed qualification — 2026-10-07
Backup verification now requires the full database/four-volume artifact set and
rechecks SHA256, including same-size corruption. Nineteen pure controller tests
pass. The owned, network-none PostgreSQL fixture passes unchanged/additive
commitments, rejects four data/schema/history mutations, restores a real custom
dump with matching original commitments, and rejects a truncated dump. It mounts
no live volume and removes only its own container. This does not qualify actual
application writer drain or the complete supervised systemd cutover.
The updater compiled and its full isolated suite ran: 1,958 passed, one failed,
five ignored. The failure is the old snake_case rental receipt JSON fixture;
`c2c4d915` already corrects that exact test on the release candidate branch.
Do not duplicate or suppress it here. Integrate and rerun the complete candidate
suite before claiming a green backend gate. The earlier interrupted compile
and PostgreSQL timeout remain failed/incomplete attempts, not acceptance.
Evidence: `/tmp/archy-resumed-20261007-updater-full-backend.log`,
`/tmp/archy-resumed-20261007-indeehub-controller-tests-final.log`, and
`/tmp/archy-resumed-20261007-indeehub-postgres-restore.log`.
@@ -691,6 +691,7 @@ server {
sub_filter '</head>' '<script src="/nostr-provider.js"></script></head>';
}
location /app/indeedhub/_next/ {
if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }
proxy_pass http://127.0.0.1:7778/_next/;
proxy_http_version 1.1;
proxy_set_header Host $host;
@@ -699,6 +700,7 @@ server {
}
# IndeeHub WebSocket proxy
location /app/indeedhub/ws/ {
if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }
proxy_pass http://127.0.0.1:7778/ws/;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
@@ -708,6 +710,7 @@ server {
proxy_read_timeout 86400s;
}
location /app/indeedhub/ {
if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }
proxy_pass http://127.0.0.1:7778/;
proxy_http_version 1.1;
proxy_set_header Host $host;
+337
View File
@@ -0,0 +1,337 @@
#!/usr/bin/env python3
"""Operation-owned legacy IndeeHub maintenance. Called only by supervised updater.
No live execution is part of source qualification. Original writable-layer images
must already be durable. Never unlock ARCHY_UPDATE_LOCK_FD or release another hold.
"""
import datetime, hashlib, json, os, pathlib, re, shutil, subprocess, sys, time, uuid
NAMES = ('indeedhub','indeedhub-api','indeedhub-ffmpeg','indeedhub-minio','indeedhub-postgres','indeedhub-redis','indeedhub-relay')
VOLUMES = ('indeedhub-minio-data','indeedhub-postgres-data','indeedhub-redis-data','indeedhub-relay-data')
DATA = pathlib.Path('/var/lib/archipelago')
QUEUE_SCRIPT = r'''const {Queue}=require('bullmq');
(async()=>{const q=new Queue('transcode',{connection:{host:process.env.QUEUE_HOST,port:Number(process.env.QUEUE_PORT||6379),password:process.env.QUEUE_PASSWORD,maxRetriesPerRequest:1}});
try{const action=process.argv[1];if(action==='pause')await q.pause();else if(action==='resume')await q.resume();else if(action!=='status')throw Error('action');
console.log(JSON.stringify({paused:await q.isPaused(),counts:await q.getJobCounts('active','waiting','paused','delayed','failed','completed')}));}
finally{await q.close()}})().catch(()=>process.exit(1));'''
# The exact three migrations in the privately qualified API candidate. This is
# an allowlist of additive schema history, never permission to discard app data.
ADDITIVE_MIGRATIONS = {
'AddArchipelagoPublicationsAndRentals1791288000000':1791288000000,
'AddMediaRegistrationIntents1791374400000':1791374400000,
'AddMediaRegistrationRetirements1791374401000':1791374401000,
}
ADDITIVE_TABLES = {'archipelago_media_registrations','archipelago_publications','archipelago_publication_outbox','archipelago_rental_entitlements','archipelago_registration_intents'}
DB_COMMITMENTS_SQL = r'''
BEGIN TRANSACTION ISOLATION LEVEL REPEATABLE READ READ ONLY;
SELECT format($query$
SELECT jsonb_build_object('table',%L,'schema',
jsonb_build_object(
'columns',(SELECT coalesce(jsonb_agg(jsonb_build_array(a.attnum,a.attname,format_type(a.atttypid,a.atttypmod),a.attnotnull,a.attidentity,a.attgenerated,pg_get_expr(d.adbin,d.adrelid)) ORDER BY a.attnum),'[]'::jsonb) FROM pg_attribute a LEFT JOIN pg_attrdef d ON d.adrelid=a.attrelid AND d.adnum=a.attnum WHERE a.attrelid=%s AND a.attnum>0 AND NOT a.attisdropped),
'constraints',(SELECT coalesce(jsonb_agg(jsonb_build_array(conname,pg_get_constraintdef(oid,true)) ORDER BY conname),'[]'::jsonb) FROM pg_constraint WHERE conrelid=%s),
'indexes',(SELECT coalesce(jsonb_agg(pg_get_indexdef(indexrelid) ORDER BY indexrelid::regclass::text),'[]'::jsonb) FROM pg_index WHERE indrelid=%s),
'triggers',(SELECT coalesce(jsonb_agg(pg_get_triggerdef(oid,true) ORDER BY tgname),'[]'::jsonb) FROM pg_trigger WHERE tgrelid=%s AND NOT tgisinternal),
'rls',%L,'policies',(SELECT coalesce(jsonb_agg(to_jsonb(p) ORDER BY policyname),'[]'::jsonb) FROM pg_policies p WHERE schemaname='public' AND tablename=%L)),
'rows',count(*),'rows_sha256',encode(sha256(convert_to(coalesce(string_agg(to_jsonb(t)::text,E'\n' ORDER BY to_jsonb(t)::text),''),'UTF8')),'hex')) FROM public.%I t;
$query$,c.relname,c.oid,c.oid,c.oid,c.oid,c.relrowsecurity::text||':'||c.relforcerowsecurity::text,c.relname,c.relname)
FROM pg_class c JOIN pg_namespace n ON n.oid=c.relnamespace WHERE n.nspname='public' AND c.relkind IN ('r','p') ORDER BY c.relname
\gexec
SELECT jsonb_build_object('migration_rows',coalesce(jsonb_agg(to_jsonb(m) ORDER BY id),'[]'::jsonb)) FROM public.migrations m;
COMMIT;
'''
def verify_database_compatibility(before, after):
require(before.get('operation_id')==after.get('operation_id'),'Data compatibility operation changed')
old=before['tables'];new=after['tables'];require(set(old)<=set(new),'Data compatibility lost original tables')
extra=set(new)-set(old);require(extra<=ADDITIVE_TABLES,'Data compatibility contains unreviewed tables')
for name in old:
require(old[name]['schema']==new[name]['schema'],'Data compatibility changed original table schema')
if name!='migrations':
require(old[name]['rows']==new[name]['rows'] and old[name]['rows_sha256']==new[name]['rows_sha256'],'Data compatibility changed original rows')
for name in extra:require(new[name]['rows']==0,'Data compatibility contains new application data')
previous=before['migrations'];current=after['migrations']
require(current[:len(previous)]==previous,'Data compatibility changed original migration history')
added=current[len(previous):];seen=set()
for row in added:
require(set(row)=={'id','timestamp','name'} and type(row['id']) is int and row['name'] not in seen and ADDITIVE_MIGRATIONS.get(row['name'])==row['timestamp'],'Data compatibility has unreviewed migration history')
seen.add(row['name'])
ordered=list(ADDITIVE_MIGRATIONS)
require([row['name'] for row in added]==ordered[:len(added)],'Data compatibility migration order changed')
expected_extra=(ADDITIVE_TABLES-{'archipelago_registration_intents'}) if added else set()
if len(added)>=2:expected_extra=ADDITIVE_TABLES
require(extra==expected_extra-set(old),'Data compatibility additions do not match reviewed migration evidence')
return {'original_tables':len(old),'new_empty_tables':sorted(extra),'reviewed_migrations':[r['name'] for r in added]}
def require(condition, message):
if not condition: raise RuntimeError(message)
def atomic(path, value):
path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
require(not path.is_symlink(), 'Refuse symbolic journal path')
temporary=path.with_name('.'+path.name+'.'+str(uuid.uuid4()))
with temporary.open('x') as stream:
json.dump(value,stream,separators=(',',':'));stream.flush();os.fsync(stream.fileno())
os.chmod(temporary,0o600);os.replace(temporary,path)
descriptor=os.open(path.parent,os.O_RDONLY);os.fsync(descriptor);os.close(descriptor)
def sha(path):
with path.open('rb') as stream:return hashlib.file_digest(stream,'sha256').hexdigest()
def validate_members(members):
require(isinstance(members,list) and len(members)==7,'Seven exact original members required')
require({m.get('name') for m in members}==set(NAMES),'IndeeHub member scope changed')
for m in members:
require(set(m)=={'name','container_id','image_id','unit_sha256','config_sha256','running'},'Unexpected member fields')
for key in ('container_id','image_id','unit_sha256','config_sha256'):
require(bool(re.fullmatch('[0-9a-f]{64}',m[key])),'Invalid original identity/hash')
require(m['running'] is True,'Legacy barrier currently supports an originally running complete stack only')
return sorted(members,key=lambda m:m['name'])
def validate_nginx_guards(config):
guard='if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }'
lines=config.splitlines();matched=0;index=0
while index<len(lines):
header=lines[index]
if not re.match(r'^\s*location\b.*\{\s*$',header):index+=1;continue
block=[];depth=0
while index<len(lines):
line=lines[index];block.append(line)
unquoted=re.sub(r"(['\"])(?:\\.|(?!\1).)*\1",'',line).split('#',1)[0]
depth+=unquoted.count('{')-unquoted.count('}');index+=1
if depth==0:break
body='\n'.join(block)
if '/app/indeedhub' in header or re.search(r'proxy_pass\s+https?://127[.]0[.]0[.]1:7778(?:/|;)',body):
require(bool(re.match(r'^\s*location\s+(?:\^~\s+|=\s+)?/app/indeedhub(?:/[^\s{]*)?\s*\{\s*$',header)),'Unrecognized direct IndeeHub proxy exposure')
require(guard in body,'An IndeeHub proxy route is missing its maintenance guard')
matched+=1
require(matched>=3,'Expected complete legacy IndeeHub route guards')
return matched
class Controller:
def __init__(self, data, operation, lock_fd, runner=None):
self.data=pathlib.Path(data);self.operation=operation;self.lock_fd=lock_fd;self.runner=runner
require(str(uuid.UUID(operation))==operation,'Invalid operation UUID')
self.root=self.data/'update-transactions'/'indeehub-maintenance'/operation
self.path=self.root/'journal.json';self.fence=self.data/'app-maintenance'/'indeedhub'
self.record=json.loads(self.path.read_text()) if self.path.exists() else None
if self.record:require(self.record['operation_id']==operation,'Maintenance journal changed')
def save(self): atomic(self.path,self.record)
def run(self, argv, timeout=30, output=None, input_bytes=None):
if self.runner:return self.runner(argv,timeout,output)
self.root.mkdir(mode=0o700,parents=True,exist_ok=True)
with (self.root/'commands.private.log').open('ab') as errors:
result=subprocess.run(argv,stdout=output or subprocess.PIPE,stderr=errors,timeout=timeout,check=True,pass_fds=(self.lock_fd,),input=input_bytes)
if output:return b''
require(len(result.stdout)<=2*1024*1024,'Command response exceeds bound')
return result.stdout
def inspect(self, name):
rows=json.loads(self.run(['podman','inspect',name]));require(len(rows)==1,'Unexpected container inspection');return rows[0]
def holds(self):
for name in NAMES:
path=self.data/'update-transactions'/'holds'/name
require(path.is_file() and not path.is_symlink() and path.read_text()==self.operation,'Matching durable lifecycle hold required')
def fence_matches(self):
require(self.fence.is_file() and not self.fence.is_symlink() and self.fence.read_text()==self.operation,'Admission fence ownership changed')
def close_ingress(self):
# The deployed native AppGate and legacy nginx guards consume this exact
# sentinel. This code never edits arbitrary nginx configuration.
config=self.run(['sudo','-n','nginx','-T']).decode()
validate_nginx_guards(config)
self.fence.parent.mkdir(mode=0o755,exist_ok=True)
self.fence.parent.chmod(0o755)
require(not self.fence.parent.is_symlink(),'Admission directory is a symlink')
if self.fence.exists():self.fence_matches()
else:
with self.fence.open('x') as stream:stream.write(self.operation);stream.flush();os.fsync(stream.fileno())
self.fence.chmod(0o644)
fd=os.open(self.fence.parent,os.O_RDONLY);os.fsync(fd);os.close(fd)
# A local legacy probe must be rejected without entering the old app.
import urllib.request,urllib.error
try:
urllib.request.urlopen('http://127.0.0.1/app/indeedhub/__maintenance_probe',timeout=5)
raise RuntimeError('Legacy ingress was not fenced')
except urllib.error.HTTPError as error:
require(error.code==503,'Legacy ingress guard did not return maintenance status')
self.record['ingress_closed']=True;self.save()
def queue(self, action):
result=json.loads(self.run(['podman','exec','indeedhub-api','node','-e',QUEUE_SCRIPT,action]))
require(type(result.get('paused')) is bool and isinstance(result.get('counts'),dict),'Invalid queue observation')
for value in result['counts'].values():require(type(value) is int and value>=0,'Invalid job count')
return result
def pause_queue(self):
if 'queue_was_paused' not in self.record:
original=self.queue('status');self.record['queue_was_paused']=original['paused'];self.record['queue_original_counts']=original['counts'];self.save()
state=self.queue('pause');require(state['paused'],'Worker admission did not close')
self.record['queue_pause_confirmed']=True;self.save()
def legacy_api_idle(self):
# Narrow first-upgrade compatibility for the observed legacy API which
# has no SIGTERM hooks. Existing customer/business work is never inferred
# completed: this path requires a fresh empty store behind closed ingress.
require(self.record.get('stopped',{}).get('indeedhub',{}).get('confirmed'),'Frontend ingress must already be stopped')
require(self.record.get('stopped',{}).get('indeedhub-ffmpeg',{}).get('confirmed'),'Transcode worker must already be stopped')
require(self.record.get('queue_pause_confirmed') is True and self.record.get('last_queue_counts',{}).get('active')==0,'Worker queue is not proven idle')
tables=('projects','contents','payments','shareholders','subscriptions','library_items')
fields=','.join("'%s',(SELECT count(*) FROM public.%s)"%(name,name) for name in tables)
sql="SELECT json_build_object("+fields+",'other_active_transactions',(SELECT count(*) FROM pg_stat_activity WHERE datname=current_database() AND pid<>pg_backend_pid() AND state<>'idle'))"
counts=json.loads(self.run(['podman','exec','indeedhub-postgres','psql','-XAt','-U','indeedhub','-d','indeedhub','-c',sql]))
require(set(counts)==set(tables)|{'other_active_transactions'},'Legacy API business-state observation incomplete')
require(all(type(value) is int and value==0 for value in counts.values()),'Legacy API has business work or active transactions; completion cannot be inferred')
self.record['legacy_api_empty_state']=counts;self.save()
def graceful_stop(self, name):
# Save the obligation before systemd can remove an AutoRemove container.
stopped=self.record.setdefault('stopped',{})
if stopped.get(name,{}).get('confirmed'):return
member=next(m for m in self.record['original_members'] if m['name']==name)
if name not in stopped:
actual=self.inspect(name);require(actual['Id']==member['container_id'] and actual['Image']==member['image_id'],'Original container changed before stop')
stopped[name]={'intent_at':time.time(),'container_id':actual['Id']};self.save()
self.run(['systemctl','--user','stop',name+'.service'],timeout=180)
properties=self.run(['systemctl','--user','show',name+'.service','--property=ActiveState,SubState,Result,ExecMainStatus']).decode()
require('ActiveState=inactive' in properties and 'Result=success' in properties,'Service did not stop successfully')
# --rm removes inspect state. Require a persisted Podman died event for
# this exact original ID; a forced SIGKILL is never called completed work.
events=self.run(['podman','events','--stream=false','--since',str(int(stopped[name]['intent_at'])-1),'--filter','container='+member['container_id'],'--filter','event=died','--format','json']).decode().splitlines()
matching=[json.loads(line) for line in events if line.strip()]
matching=[event for event in matching if event.get('ID',event.get('id'))==member['container_id']]
require(matching,'Original process exit evidence unavailable; hold retained')
code=matching[-1].get('ContainerExitCode',matching[-1].get('containerExitCode'))
idle_worker = name=='indeedhub-ffmpeg' and self.record.get('queue_pause_confirmed') is True and self.record.get('last_queue_counts',{}).get('active')==0
empty_api = name=='indeedhub-api' and self.record.get('legacy_api_empty_state') is not None
require(str(code)=='0' or (str(code)=='143' and (idle_worker or empty_api)),'Original process did not exit cleanly; active work is not claimed completed')
classification=('idle-worker-terminated-after-queue-drain' if idle_worker else 'empty-business-store-legacy-api-terminated') if str(code)=='143' else 'clean-process-exit'
stopped[name].update(confirmed=True,exit_code=int(code),classification=classification,confirmed_at=time.time());self.save()
def volume_sources(self):
expected=VOLUMES
rows=json.loads(self.run(['podman','volume','inspect',*expected]))
require({row['Name'] for row in rows}==set(expected),'Persistent volume scope changed')
return {row['Name']:row['Mountpoint'] for row in rows}
def database_commitments(self):
raw=self.run(['podman','exec','-i','indeedhub-postgres','psql','-XqAt','--set=ON_ERROR_STOP=1','-U','indeedhub','-d','indeedhub'],timeout=300,input_bytes=DB_COMMITMENTS_SQL.encode())
rows=[json.loads(line) for line in raw.decode().splitlines() if line.strip()]
tables={};migrations=None
for row in rows:
if 'migration_rows' in row:
require(migrations is None,'Duplicate database migration observation');migrations=row['migration_rows']
else:
name=row.pop('table');require(name not in tables and re.fullmatch('[a-zA-Z_][a-zA-Z0-9_]*',name),'Invalid database table observation');tables[name]=row
require(tables and 'migrations' in tables and isinstance(migrations,list),'Database compatibility observation incomplete')
return {'operation_id':self.operation,'tables':tables,'migrations':migrations}
def verify_restored_data(self):
baseline=self.record.get('database_before')
require(baseline and baseline.get('operation_id')==self.operation,'Data compatibility baseline missing')
current=self.database_commitments();proof=verify_database_compatibility(baseline,current)
self.record['recovery_data_verification']={'operation_id':self.operation,'checked_at':time.time(),'before_sha256':hashlib.sha256(json.dumps(baseline,sort_keys=True).encode()).hexdigest(),'after_sha256':hashlib.sha256(json.dumps(current,sort_keys=True).encode()).hexdigest(),**proof}
self.record['recovery_data_verified']=True;self.save()
def backup(self):
if self.record.get('backup_complete'):return
if 'database_before' not in self.record:
self.record['database_before']=self.database_commitments();self.save()
sources=self.volume_sources();self.record['volume_sources']=sources;self.save()
backup=self.root/'backup';backup.mkdir(mode=0o700,exist_ok=True)
if 'database.dump' not in self.record.setdefault('artifacts',{}):
path=backup/'database.dump.partial'
with path.open('wb') as output:self.run(['podman','exec','indeedhub-postgres','pg_dump','-U','indeedhub','-d','indeedhub','--format=custom','--no-owner','--no-acl'],timeout=300,output=output);output.flush();os.fsync(output.fileno())
final=backup/'database.dump';os.replace(path,final);self.record['artifacts']['database.dump']={'bytes':final.stat().st_size,'sha256':sha(final)};self.save()
# Redis stop flushes persisted queue state; its clean process exit is
# checked exactly as every other service. SQLite WAL is archived with DB.
for name in ('indeedhub-minio','indeedhub-redis','indeedhub-relay','indeedhub-postgres'):self.graceful_stop(name)
for volume,source in sources.items():
name=volume+'.tar'
if name in self.record['artifacts']:continue
require(pathlib.Path(source).is_absolute() and source.endswith('/_data'),'Invalid volume mountpoint')
available=shutil.disk_usage(backup).free
measured=int(self.run(['podman','unshare','du','-sb',source]).decode().split()[0])
require(available>measured+512*1024*1024,'Insufficient durable backup space')
partial=backup/(name+'.partial')
with partial.open('wb') as output:self.run(['podman','unshare','tar','--xattrs','--acls','--numeric-owner','-C',source,'-cpf','-','.'],timeout=1800,output=output);output.flush();os.fsync(output.fileno())
final=backup/name;os.replace(partial,final);self.record['artifacts'][name]={'bytes':final.stat().st_size,'sha256':sha(final)};self.save()
self.record['backup_complete']=True;self.record['phase']='Drained';self.save()
def acquire(self, members, recovery=False):
members=validate_members(members);self.holds()
if self.record:require(self.record['original_members']==members,'Original operation terms changed')
else:
self.record={'operation_id':self.operation,'original_members':members,'phase':'Prepared','created_at':time.time()};self.save()
require(self.record['phase']!='Released','Completed maintenance must not be reacquired')
if not recovery and not self.record.get('originals_validated'):
for member in members:
actual=self.inspect(member['name']);require(actual['Id']==member['container_id'] and actual['Image']==member['image_id'],'Original member changed')
bindings=actual['HostConfig'].get('PortBindings') or {}
if member['name']=='indeedhub':require(bindings=={'7777/tcp':[{'HostIp':'127.0.0.1','HostPort':'7778'}]},'Unsupported direct frontend exposure')
else:require(not bindings,'Unsupported direct writer exposure')
source=pathlib.Path(self.run(['systemctl','--user','show',member['name']+'.service','--property=SourcePath','--value']).decode().strip())
require(source.is_file() and not source.is_symlink() and source.suffix=='.container','Original unit source missing')
require(source.stat().st_uid==os.getuid() and sha(source)==member['unit_sha256'],'Original unit source changed')
self.record['originals_validated']=True;self.save()
self.close_ingress()
if recovery:
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text())
require(runtime.get('phase')=='Restoring' and type(runtime.get('target_startup_began')) is bool,'Durable explicit restoring obligation required')
self.record['phase']='Recovering';self.record['target_startup_began']=runtime['target_startup_began'];self.save()
return {'operation_id':self.operation,'state':'recovering'}
if not self.record.get('stopped',{}).get('indeedhub-api',{}).get('confirmed'):
self.pause_queue()
self.graceful_stop('indeedhub')
deadline=time.monotonic()+300
while True:
state=self.queue('status');require(state['paused'],'Worker admission reopened')
self.record['last_queue_counts']=state['counts'];self.save()
if state['counts'].get('active',0)==0:break
require(time.monotonic()<deadline,'Transcodes still active; retained job state, no forced completion')
time.sleep(1)
self.graceful_stop('indeedhub-ffmpeg');self.legacy_api_idle();self.graceful_stop('indeedhub-api')
self.backup();self.verify();return {'operation_id':self.operation,'state':'drained'}
def verify(self):
self.holds();self.fence_matches()
if self.record and self.record.get('phase')=='Recovering':
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text())
require(runtime.get('phase')=='Restoring','Recovery ownership changed')
return {'operation_id':self.operation,'state':'held'}
require(self.record and self.record.get('backup_complete'),'Drain not complete')
require(self.record['phase'] in ('Drained','Released'),'Invalid maintenance phase')
# Verification remains possible when API/storage endpoints are stopped.
# The native adapter separately validates target/original runtime identity.
for name in NAMES:require(self.record.get('stopped',{}).get(name,{}).get('confirmed'),'Original writer stop evidence missing')
expected_artifacts={'database.dump',*(volume+'.tar' for volume in VOLUMES)}
require(set(self.record.get('artifacts',{}))==expected_artifacts,'Backup artifact inventory incomplete or unexpected')
for name,record in self.record['artifacts'].items():
path=self.root/'backup'/name;require(path.is_file() and not path.is_symlink() and path.stat().st_size==record['bytes'],'Backup artifact missing or changed')
require(sha(path)==record['sha256'],'Backup artifact checksum changed')
return {'operation_id':self.operation,'state':'held'}
def release(self, outcome):
require(outcome in ('committed','restored','aborted'),'Invalid release outcome')
if self.record is None and outcome=='aborted':
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text())
require(runtime.get('phase')=='Aborted' and runtime.get('target_startup_began') is False,'Untouched abort evidence required')
if self.fence.exists():
require(not self.fence.is_symlink() and self.fence.read_text()!=self.operation,'Matching fence without journal requires recovery')
return {'operation_id':self.operation,'state':'released'}
require(self.record is not None,'Unknown maintenance operation')
if self.record['phase']=='Released':
require(self.record.get('outcome')==outcome,'Maintenance outcome changed')
if self.fence.exists():
self.fence_matches();self.fence.unlink()
return {'operation_id':self.operation,'state':'released'}
self.holds();self.fence_matches()
runtime=json.loads((self.data/'update-transactions'/'supervised'/(self.operation+'.json')).read_text());phase=runtime['phase']
require(phase=={'committed':'Committed','restored':'Restored','aborted':'Restored'}[outcome],'Runtime outcome not durably verified')
# Restoring old runtime over a changed database is not sufficient to open
# admission. Node must verify same-schema/additive migration compatibility.
if outcome!='committed':
require(type(runtime.get('target_startup_began')) is bool,'Target-start obligation unavailable')
if runtime['target_startup_began']:
self.verify_restored_data()
else:
self.record['rollback_data_claim']='No target startup/migration began; only original runtime restored.'
if not self.record.get('queue_was_paused',True):
state=self.queue('resume');require(not state['paused'],'Could not restore queue admission')
self.record['phase']='Released';self.record['outcome']=outcome;self.save()
self.fence.unlink();fd=os.open(self.fence.parent,os.O_RDONLY);os.fsync(fd);os.close(fd)
return {'operation_id':self.operation,'state':'released'}
def main():
os.umask(0o077);require(len(sys.argv)==2 and sys.argv[1] in ('acquire','verify','release'),'Unsupported maintenance action')
raw=sys.stdin.buffer.read(65537);require(len(raw)<=65536,'Maintenance request too large');request=json.loads(raw)
allowed={'operation_id','original_members'} if sys.argv[1]=='acquire' else {'operation_id','outcome'} if sys.argv[1]=='release' else {'operation_id'}
require(set(request) in (allowed, allowed|{'recovery'}) if sys.argv[1]=='acquire' else set(request)==allowed,'Unexpected maintenance fields')
if 'recovery' in request:require(type(request['recovery']) is bool,'Invalid recovery flag')
require(os.getuid()==1000,'Expected node service user');fd=int(os.environ['ARCHY_UPDATE_LOCK_FD']);actual=os.fstat(fd);expected=(DATA/'update-transactions'/'lock').stat();require((actual.st_dev,actual.st_ino)==(expected.st_dev,expected.st_ino),'Inherited lifecycle lock is not the expected file')
controller=Controller(DATA,request['operation_id'],fd)
result=controller.acquire(request['original_members'],request.get('recovery',False)) if sys.argv[1]=='acquire' else controller.verify() if sys.argv[1]=='verify' else controller.release(request['outcome'])
encoded=json.dumps(result);require(len(encoded)<=4096,'Maintenance response exceeds bound');print(encoded)
if __name__=='__main__':
try:main()
except Exception as error:
# Command/env details remain in private journal, never RPC/UI stdout.
print(json.dumps({'error':'Maintenance remains held; inspect its private operation journal','reason':type(error).__name__}),file=sys.stderr);sys.exit(1)
@@ -0,0 +1,171 @@
"""Pure fixture coverage; imports the controller without invoking main/live tools."""
import importlib.util,json,pathlib,tempfile,unittest,uuid
MODULE=pathlib.Path(__file__).resolve().parents[2]/'scripts/indeehub-maintenance-controller.py'
spec=importlib.util.spec_from_file_location('maintenance_controller',MODULE);module=importlib.util.module_from_spec(spec);spec.loader.exec_module(module)
def members():
return [{'name':name,'container_id':f'{n+1:064x}','image_id':'a'*64,'unit_sha256':'b'*64,'config_sha256':'c'*64,'running':True} for n,name in enumerate(module.NAMES)]
class MaintenanceTests(unittest.TestCase):
def setUp(self):
self.tmp=tempfile.TemporaryDirectory();self.addCleanup(self.tmp.cleanup);self.root=pathlib.Path(self.tmp.name);self.operation=str(uuid.uuid4());self.calls=[]
holds=self.root/'update-transactions'/'holds';holds.mkdir(parents=True)
for name in module.NAMES:(holds/name).write_text(self.operation)
self.controller=module.Controller(self.root,self.operation,0,self.command_runner)
def command_runner(self,argv,timeout,output):
self.calls.append(argv)
if argv[:2]==['podman','inspect']:
member=next(m for m in members() if m['name']==argv[2]);return json.dumps([{'Id':member['container_id'],'Image':member['image_id'],'HostConfig':{'PortBindings':{}}}]).encode()
raise AssertionError('Unexpected fixture command '+str(argv))
def test_exact_seven_member_identity_is_required(self):
self.assertEqual(len(module.validate_members(members())),7)
for bad in [members()[:-1],members()+[members()[0]],[dict(members()[0],container_id='invalid'),*members()[1:]]]:
with self.assertRaises(RuntimeError):module.validate_members(bad)
def test_foreign_operation_cannot_release_or_overwrite_hold(self):
path=self.root/'update-transactions'/'holds'/'indeedhub';path.write_text(str(uuid.uuid4()))
with self.assertRaises(RuntimeError):self.controller.holds()
self.assertEqual(self.calls,[])
def test_failed_original_exposure_validation_is_not_skipped_on_retry(self):
for _ in range(2):
with self.assertRaisesRegex(RuntimeError,'frontend exposure'):self.controller.acquire(members())
self.controller=module.Controller(self.root,self.operation,0,self.command_runner)
self.assertFalse(self.controller.record.get('originals_validated',False))
self.assertEqual(len(self.calls),2)
def test_changed_original_terms_cannot_resume_saved_operation(self):
self.controller.record={'operation_id':self.operation,'original_members':module.validate_members(members()),'phase':'Prepared'};self.controller.save()
changed=members();changed[0]['config_sha256']='d'*64
with self.assertRaisesRegex(RuntimeError,'terms changed'):self.controller.acquire(changed)
self.assertEqual(self.calls,[])
def test_release_lost_reply_is_idempotent_but_never_removes_foreign_fence(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Released','outcome':'committed'};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
self.assertEqual(c.release('committed')['state'],'released');self.assertFalse(c.fence.exists())
self.assertEqual(c.release('committed')['state'],'released')
c.fence.write_text(str(uuid.uuid4()))
with self.assertRaises(RuntimeError):c.release('committed')
self.assertTrue(c.fence.exists())
def test_backup_file_existence_never_substitutes_for_drain_evidence(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','backup_complete':False};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
(c.root/'backup').mkdir();(c.root/'backup'/'database.dump').write_bytes(b'not evidence')
with self.assertRaisesRegex(RuntimeError,'Drain not complete'):c.verify()
def completed_backup(self):
c=self.controller
c.record={'operation_id':self.operation,'phase':'Drained','backup_complete':True,
'stopped':{name:{'confirmed':True} for name in module.NAMES},'artifacts':{}}
c.save();c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
(c.root/'backup').mkdir()
for name in ['database.dump',*(v+'.tar' for v in module.VOLUMES)]:
path=c.root/'backup'/name;path.write_bytes(b'original')
c.record['artifacts'][name]={'bytes':path.stat().st_size,'sha256':module.sha(path)}
c.save()
return c
def test_complete_backup_checksums_allow_verification(self):
self.assertEqual(self.completed_backup().verify()['state'],'held')
def test_same_size_corruption_keeps_admission_closed(self):
c=self.completed_backup();(c.root/'backup'/'database.dump').write_bytes(b'corrupt!')
with self.assertRaisesRegex(RuntimeError,'checksum changed'):c.verify()
self.assertEqual(c.fence.read_text(),self.operation);self.assertEqual(self.calls,[])
def test_incomplete_or_unexpected_artifact_inventory_cannot_pass(self):
c=self.completed_backup();original=dict(c.record['artifacts'])
for artifacts in [{}, {k:v for k,v in original.items() if k!='database.dump'},
{**original,'../foreign':original['database.dump']}]:
c.record['artifacts']=artifacts
with self.assertRaisesRegex(RuntimeError,'inventory'):c.verify()
self.assertEqual(c.fence.read_text(),self.operation)
def test_forced_original_exit_never_marks_writer_completed(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','original_members':module.validate_members(members())};c.save()
def command(argv,timeout,output):
if argv[:2]==['podman','inspect']:return self.command_runner(argv,timeout,output)
if argv[:3]==['systemctl','--user','stop']:return b''
if argv[:3]==['systemctl','--user','show']:return b'ActiveState=inactive\nResult=success\n'
if argv[:2]==['podman','events']:return json.dumps({'ID':members()[0]['container_id'],'ContainerExitCode':137}).encode()
raise AssertionError(argv)
c.runner=command
with self.assertRaisesRegex(RuntimeError,'did not exit cleanly'):c.graceful_stop('indeedhub')
self.assertFalse(c.record['stopped']['indeedhub'].get('confirmed',False))
self.assertTrue((self.root/'update-transactions'/'holds'/'indeedhub').exists())
def test_failed_stop_can_restore_without_fabricating_completed_drain(self):
c=self.controller;c.record={'operation_id':self.operation,'original_members':module.validate_members(members()),'phase':'Prepared','stopped':{'indeedhub':{'intent_at':1,'exit_code':137}}};c.save()
runtime=c.data/'update-transactions'/'supervised'/(self.operation+'.json');module.atomic(runtime,{'phase':'Restoring','target_startup_began':False})
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
c.close_ingress=lambda: c.fence_matches()
self.assertEqual(c.acquire(members(),recovery=True)['state'],'recovering')
self.assertFalse(c.record.get('backup_complete',False))
self.assertEqual(c.verify()['state'],'held')
module.atomic(runtime,{'phase':'Restored','target_startup_began':False})
self.assertEqual(c.release('aborted')['state'],'released')
self.assertIn('No target startup',c.record['rollback_data_claim'])
def test_target_started_rollback_requires_data_compatibility(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Recovering'};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
module.atomic(c.data/'update-transactions'/'supervised'/(self.operation+'.json'),{'phase':'Restored','target_startup_began':True})
with self.assertRaisesRegex(RuntimeError,'Data compatibility'):c.release('restored')
self.assertTrue(c.fence.exists())
def test_every_legacy_sublocation_must_be_fenced(self):
guard='if (-f /var/lib/archipelago/app-maintenance/indeedhub) { return 503; }'
blocks=[f'location /app/indeedhub/{suffix} {{\n {guard}\n proxy_pass http://127.0.0.1:7778/;\n}}' for suffix in ('','_next/','ws/')]
self.assertEqual(module.validate_nginx_guards('\n'.join(blocks)),3)
with self.assertRaisesRegex(RuntimeError,'missing its maintenance guard'):module.validate_nginx_guards('\n'.join(blocks).replace(guard,'',1))
with self.assertRaisesRegex(RuntimeError,'Unrecognized direct'):module.validate_nginx_guards('\n'.join(blocks)+'\nlocation /other/ {\n proxy_pass http://127.0.0.1:7778/;\n}')
def test_pre_acquire_abort_acknowledges_without_mutating_foreign_fence(self):
c=self.controller;runtime=c.data/'update-transactions'/'supervised'/(self.operation+'.json')
module.atomic(runtime,{'phase':'Aborted','target_startup_began':False})
self.assertEqual(c.release('aborted')['state'],'released')
c.fence.parent.mkdir(parents=True);foreign=str(uuid.uuid4());c.fence.write_text(foreign)
self.assertEqual(c.release('aborted')['state'],'released');self.assertEqual(c.fence.read_text(),foreign)
module.atomic(runtime,{'phase':'Aborted','target_startup_began':True})
with self.assertRaisesRegex(RuntimeError,'Untouched abort'):c.release('aborted')
def test_matching_fence_without_journal_is_not_an_untouched_abort(self):
c=self.controller;module.atomic(c.data/'update-transactions'/'supervised'/(self.operation+'.json'),{'phase':'Aborted','target_startup_began':False})
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
with self.assertRaisesRegex(RuntimeError,'without journal'):c.release('aborted')
self.assertTrue(c.fence.exists())
def test_legacy_worker_sigterm_requires_proven_paused_idle_queue(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','original_members':module.validate_members(members())};c.save()
worker=next(m for m in members() if m['name']=='indeedhub-ffmpeg')
def command(argv,timeout,output):
if argv[:2]==['podman','inspect']:return self.command_runner(argv,timeout,output)
if argv[:3]==['systemctl','--user','stop']:return b''
if argv[:3]==['systemctl','--user','show']:return b'ActiveState=inactive\nResult=success\n'
if argv[:2]==['podman','events']:return json.dumps({'ID':worker['container_id'],'ContainerExitCode':143}).encode()
raise AssertionError(argv)
c.runner=command
with self.assertRaises(RuntimeError):c.graceful_stop(worker['name'])
c.record['queue_pause_confirmed']=True;c.record['last_queue_counts']={'active':1}
with self.assertRaises(RuntimeError):c.graceful_stop(worker['name'])
c.record['last_queue_counts']['active']=0;c.graceful_stop(worker['name'])
self.assertEqual(c.record['stopped'][worker['name']]['classification'],'idle-worker-terminated-after-queue-drain')
def test_legacy_api_compatibility_requires_fresh_empty_business_state(self):
c=self.controller;c.record={'operation_id':self.operation,'phase':'Prepared','stopped':{'indeedhub':{'confirmed':True},'indeedhub-ffmpeg':{'confirmed':True}},'queue_pause_confirmed':True,'last_queue_counts':{'active':0}};c.save()
counts={name:0 for name in ('projects','contents','payments','shareholders','subscriptions','library_items','other_active_transactions')}
c.runner=lambda argv,timeout,output:json.dumps(counts).encode()
counts['payments']=1
with self.assertRaisesRegex(RuntimeError,'business work'):c.legacy_api_idle()
self.assertNotIn('legacy_api_empty_state',c.record)
counts['payments']=0;counts['other_active_transactions']=1
with self.assertRaisesRegex(RuntimeError,'business work'):c.legacy_api_idle()
counts['other_active_transactions']=0;c.legacy_api_idle()
self.assertEqual(c.record['legacy_api_empty_state'],counts)
def test_rollback_compatibility_binds_operation_preserves_rows_and_allows_only_empty_additions(self):
import copy
table={'schema':{'columns':['original']},'rows':0,'rows_sha256':'a'*64}
before={'operation_id':self.operation,'tables':{'migrations':copy.deepcopy(table),'contents':copy.deepcopy(table)},'migrations':[{'id':1,'timestamp':1,'name':'Original1'}]}
after=copy.deepcopy(before)
self.assertEqual(module.verify_database_compatibility(before,after)['original_tables'],2)
names=list(module.ADDITIVE_MIGRATIONS)
after['migrations'] += [{'id':i+2,'timestamp':module.ADDITIVE_MIGRATIONS[name],'name':name} for i,name in enumerate(names)]
after['tables']['migrations']['rows']=4;after['tables']['migrations']['rows_sha256']='b'*64
for name in module.ADDITIVE_TABLES:after['tables'][name]=copy.deepcopy(table)
self.assertEqual(len(module.verify_database_compatibility(before,after)['new_empty_tables']),5)
for mutate in [lambda d:d.update(operation_id=str(uuid.uuid4())),lambda d:d['tables']['contents'].update(rows_sha256='c'*64),lambda d:d['tables']['contents']['schema'].update(columns=['changed']),lambda d:d['tables']['archipelago_publications'].update(rows=1),lambda d:d['migrations'][0].update(name='Altered'),lambda d:d['migrations'][-1].update(name='Unreviewed'),lambda d:d['tables'].update(unreviewed=copy.deepcopy(table))]:
damaged=copy.deepcopy(after);mutate(damaged)
with self.assertRaisesRegex(RuntimeError,'Data compatibility'):module.verify_database_compatibility(before,damaged)
def test_verified_rollback_records_operation_proof_before_releasing_fence(self):
c=self.controller;table={'schema':{},'rows':0,'rows_sha256':'a'*64};baseline={'operation_id':self.operation,'tables':{'migrations':table},'migrations':[]}
c.record={'operation_id':self.operation,'phase':'Recovering','database_before':baseline};c.save()
c.fence.parent.mkdir(parents=True);c.fence.write_text(self.operation)
module.atomic(c.data/'update-transactions'/'supervised'/(self.operation+'.json'),{'phase':'Restored','target_startup_began':True})
c.database_commitments=lambda:baseline
self.assertEqual(c.release('restored')['state'],'released')
self.assertEqual(c.record['recovery_data_verification']['operation_id'],self.operation)
self.assertEqual(c.record['recovery_data_verification']['before_sha256'],c.record['recovery_data_verification']['after_sha256'])
if __name__=='__main__':unittest.main()
@@ -0,0 +1,114 @@
#!/usr/bin/env python3
"""Exercise rollback commitments on an owned, network-isolated PostgreSQL.
Requires an already imported image: --image IMAGE. Never mounts node volumes,
publishes ports, or invokes the maintenance entrypoint against installed apps.
"""
import argparse
import importlib.util
import json
from pathlib import Path
import subprocess
import tempfile
import time
import uuid
MODULE = Path(__file__).resolve().parents[2] / 'scripts/indeehub-maintenance-controller.py'
spec = importlib.util.spec_from_file_location('maintenance', MODULE)
maintenance = importlib.util.module_from_spec(spec)
spec.loader.exec_module(maintenance)
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('--image', required=True)
args = parser.parse_args()
image = subprocess.check_output(
['podman', 'image', 'inspect', '--format', '{{.Id}}', args.image], text=True,
).strip()
name = 'archy-maintenance-sql-' + uuid.uuid4().hex
container = None
try:
container = subprocess.check_output([
'podman', 'run', '-d', '--pull=never', '--network=none', '--name', name,
'--tmpfs', '/var/lib/postgresql/data:rw',
'-e', 'POSTGRES_HOST_AUTH_METHOD=trust', '-e', 'POSTGRES_USER=indeedhub',
'-e', 'POSTGRES_DB=indeedhub', image,
], text=True).strip()
deadline = time.monotonic() + 60
while subprocess.run(['podman', 'exec', container, 'pg_isready', '-U', 'indeedhub'],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL).returncode:
if time.monotonic() > deadline:
raise RuntimeError('Disposable PostgreSQL did not become ready')
time.sleep(0.5)
def sql(statement, database='indeedhub'):
return subprocess.check_output([
'podman', 'exec', '-i', container, 'psql', '-XqAt',
'--set=ON_ERROR_STOP=1', '-U', 'indeedhub', '-d', database,
], input=statement.encode(), timeout=60)
sql('CREATE TABLE migrations(id serial PRIMARY KEY, timestamp bigint NOT NULL, name text NOT NULL);'
"INSERT INTO migrations(timestamp,name) VALUES(1,'Original1');"
'CREATE TABLE contents(id int PRIMARY KEY, title text NOT NULL);'
"INSERT INTO contents VALUES(1,'retained original');")
with tempfile.TemporaryDirectory(prefix=name) as root:
controller = maintenance.Controller(root, str(uuid.uuid4()), 0)
database = 'indeedhub'
def fixture_run(argv, timeout=30, output=None, input_bytes=None):
assert argv[:4] == ['podman', 'exec', '-i', 'indeedhub-postgres']
assert output is None and input_bytes is not None
return sql(input_bytes.decode(), database)
controller.run = fixture_run
before = controller.database_commitments()
dump = subprocess.check_output([
'podman', 'exec', container, 'pg_dump', '-U', 'indeedhub',
'-d', 'indeedhub', '--format=custom', '--no-owner', '--no-acl',
], timeout=60)
sql('CREATE DATABASE restore_check')
restore_command = ['podman', 'exec', '-i', container, 'pg_restore',
'-U', 'indeedhub', '-d', 'restore_check',
'--exit-on-error', '--no-owner', '--no-acl']
subprocess.run(restore_command, input=dump, check=True, timeout=60)
database = 'restore_check'
maintenance.verify_database_compatibility(before, controller.database_commitments())
database = 'indeedhub'
rejected_dump = subprocess.run(restore_command, input=dump[:32], timeout=60,
stdout=subprocess.PIPE, stderr=subprocess.PIPE)
assert rejected_dump.returncode != 0, 'Truncated dump incorrectly accepted'
maintenance.verify_database_compatibility(before, controller.database_commitments())
for table in sorted(maintenance.ADDITIVE_TABLES):
sql(f'CREATE TABLE {table}(id int PRIMARY KEY);')
for migration, timestamp in maintenance.ADDITIVE_MIGRATIONS.items():
sql(f"INSERT INTO migrations(timestamp,name) VALUES({timestamp},'{migration}');")
proof = maintenance.verify_database_compatibility(before, controller.database_commitments())
assert len(proof['new_empty_tables']) == 5
rejected = 0
for mutation, undo in [
("UPDATE contents SET title='changed'", "UPDATE contents SET title='retained original'"),
('ALTER TABLE contents ADD COLUMN unexpected text', 'ALTER TABLE contents DROP COLUMN unexpected'),
('INSERT INTO archipelago_publications VALUES(1)', 'DELETE FROM archipelago_publications'),
("UPDATE migrations SET name='changed' WHERE id=1", "UPDATE migrations SET name='Original1' WHERE id=1"),
]:
sql(mutation)
try:
maintenance.verify_database_compatibility(before, controller.database_commitments())
except RuntimeError:
rejected += 1
else:
raise AssertionError('Changed database incorrectly accepted')
sql(undo)
maintenance.verify_database_compatibility(before, controller.database_commitments())
print(json.dumps({'postgres_commitments': 'passed', 'rejected_mutations': rejected,
'network': 'none', 'live_volumes_mounted': False,
'custom_dump_restored': True, 'truncated_dump_rejected': True}))
finally:
if container:
subprocess.run(['podman', 'rm', '-f', container], check=True, stdout=subprocess.DEVNULL)
if __name__ == '__main__':
main()