From 1bbf85e0d438efaf15ba194626f9f452ac517e98 Mon Sep 17 00:00:00 2001
From: archipelago
Date: Wed, 7 Oct 2026 01:14:49 -0400
Subject: [PATCH 01/13] Recover stopped update staging draft on current
private-image preflight
---
.../src/api/rpc/package/async_lifecycle.rs | 6 +-
.../archipelago/src/api/rpc/package/config.rs | 39 +-
.../src/api/rpc/package/runtime.rs | 14 +-
.../archipelago/src/api/rpc/package/update.rs | 177 ++++-
.../src/container/docker_packages.rs | 17 +-
core/archipelago/src/container/hooks.rs | 12 +
core/archipelago/src/container/mod.rs | 6 +-
.../src/container/prod_orchestrator.rs | 637 +++++++++++++++++-
.../src/container/staged_update.rs | 314 +++++++++
core/archipelago/src/container/traits.rs | 5 +
10 files changed, 1206 insertions(+), 21 deletions(-)
create mode 100644 core/archipelago/src/container/staged_update.rs
diff --git a/core/archipelago/src/api/rpc/package/async_lifecycle.rs b/core/archipelago/src/api/rpc/package/async_lifecycle.rs
index 80678e92..63ad037b 100644
--- a/core/archipelago/src/api/rpc/package/async_lifecycle.rs
+++ b/core/archipelago/src/api/rpc/package/async_lifecycle.rs
@@ -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
diff --git a/core/archipelago/src/api/rpc/package/config.rs b/core/archipelago/src/api/rpc/package/config.rs
index 20cf592d..28586bfd 100644
--- a/core/archipelago/src/api/rpc/package/config.rs
+++ b/core/archipelago/src/api/rpc/package/config.rs
@@ -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> {
+ 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");
diff --git a/core/archipelago/src/api/rpc/package/runtime.rs b/core/archipelago/src/api/rpc/package/runtime.rs
index ca081db9..9a46c547 100644
--- a/core/archipelago/src/api/rpc/package/runtime.rs
+++ b/core/archipelago/src/api/rpc/package/runtime.rs
@@ -126,7 +126,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::()
+ .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!(
diff --git a/core/archipelago/src/api/rpc/package/update.rs b/core/archipelago/src/api/rpc/package/update.rs
index f6daf30e..19957959 100644
--- a/core/archipelago/src/api/rpc/package/update.rs
+++ b/core/archipelago/src/api/rpc/package/update.rs
@@ -4,7 +4,7 @@
//! 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;
@@ -13,6 +13,7 @@ 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};
@@ -60,9 +61,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 +127,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,
@@ -574,13 +603,69 @@ impl RpcHandler {
}
}
-async fn inspect_update_images(package_id: &str) -> Result> {
- 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 {
+ 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> {
+ 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> {
+ 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 +711,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 +737,7 @@ fn update_targets_need_change(
installed: &[(String, String)],
) -> Result {
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 +858,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};
diff --git a/core/archipelago/src/container/docker_packages.rs b/core/archipelago/src/container/docker_packages.rs
index 2b341b80..1fe2f1d0 100644
--- a/core/archipelago/src/container/docker_packages.rs
+++ b/core/archipelago/src/container/docker_packages.rs
@@ -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: {} ({})",
diff --git a/core/archipelago/src/container/hooks.rs b/core/archipelago/src/container/hooks.rs
index 4541792c..0ca32aa1 100644
--- a/core/archipelago/src/container/hooks.rs
+++ b/core/archipelago/src/container/hooks.rs
@@ -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 } => {
diff --git a/core/archipelago/src/container/mod.rs b/core/archipelago/src/container/mod.rs
index 60a71a1a..0153fd0d 100644
--- a/core/archipelago/src/container/mod.rs
+++ b/core/archipelago/src/container/mod.rs
@@ -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,5 @@ pub use dev_orchestrator::DevContainerOrchestrator;
pub use docker_packages::DockerPackageScanner;
pub use prod_orchestrator::ProdContainerOrchestrator;
pub use traits::ContainerOrchestrator;
+
+mod staged_update;
diff --git a/core/archipelago/src/container/prod_orchestrator.rs b/core/archipelago/src/container/prod_orchestrator.rs
index 2dfcbbbc..2c547f16 100644
--- a/core/archipelago/src/container/prod_orchestrator.rs
+++ b/core/archipelago/src/container/prod_orchestrator.rs
@@ -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)]
@@ -2043,7 +2054,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 +2349,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 +2868,171 @@ 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 {
+ 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<()> {
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 +3187,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,
@@ -4825,6 +5014,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
}
async fn install(&self, app_id: &str) -> Result {
+ 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 +5087,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 +5106,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,
@@ -5136,6 +5334,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
self.remove_quadlet_unit_if_present(&name).await?;
}
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 +5389,83 @@ 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?;
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 +5481,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 +5492,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
"#,
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 {
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"));
+ }
+}
diff --git a/core/archipelago/src/bootstrap.rs b/core/archipelago/src/bootstrap.rs
index 34d068cb..0b521061 100644
--- a/core/archipelago/src/bootstrap.rs
+++ b/core/archipelago/src/bootstrap.rs
@@ -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 {
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 {
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);
+ }
+}
diff --git a/core/archipelago/src/container/supervised_runtime.rs b/core/archipelago/src/container/supervised_runtime.rs
index 21f51705..68a79287 100644
--- a/core/archipelago/src/container/supervised_runtime.rs
+++ b/core/archipelago/src/container/supervised_runtime.rs
@@ -1,7 +1,7 @@
//! Production systemd/Podman adapter. Application write admission/drain is an
//! explicit dependency: neither process pause nor a filesystem receipt is drain.
use super::{
- supervised_update::{self, PreparedTarget, RecoveryImage, Supervisor, Unit},
+ supervised_update::{self, Completion, PreparedTarget, RecoveryImage, Supervisor, Unit},
update_transaction::{Observed, Podman, Runtime, Target},
};
use anyhow::{Context, Result};
@@ -10,24 +10,155 @@ use std::{
collections::HashMap,
future::Future,
io::Write,
- os::unix::fs::{MetadataExt, OpenOptionsExt},
+ os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt},
path::{Path, PathBuf},
time::Duration,
};
pub(crate) trait DrainBarrier: Sync {
- /// Persist ownership before blocking admissions. Return only after API,
- /// direct uploads and workers have drained and the coherent backup finished.
+ /// Called after the node persisted original writable recovery images.
+ /// Persist ownership before blocking admissions or stopping writers. Return
+ /// only after API/direct uploads/workers drained and coherent backup finished.
fn acquire(
&self,
operation: &str,
originals: &[Unit],
+ recovery: bool,
) -> impl Future