fix(indeehub): preserve reviewed runtimes across reconciliation and lifecycle

This commit is contained in:
archipelago
2026-10-07 18:51:32 -04:00
parent d19124f5dc
commit b0b95810e0
5 changed files with 312 additions and 15 deletions
+29 -12
View File
@@ -111,7 +111,7 @@ impl RpcHandler {
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
do_orchestrator_package_start(orchestrator.as_ref(), &to_start, &data_dir).await
} else {
do_package_start(&to_start).await
};
@@ -239,11 +239,12 @@ impl RpcHandler {
.await;
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_stop(orchestrator.as_ref(), &to_stop).await
do_orchestrator_package_stop(orchestrator.as_ref(), &to_stop, &data_dir).await
} else {
do_package_stop(&containers).await
};
@@ -348,7 +349,7 @@ impl RpcHandler {
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
do_orchestrator_package_restart(orchestrator.as_ref(), &to_restart, &data_dir).await
} else {
do_package_restart(&containers).await
};
@@ -872,17 +873,24 @@ async fn do_package_start(to_start: &[String]) -> Result<()> {
async fn do_orchestrator_package_start(
orchestrator: &dyn crate::container::traits::ContainerOrchestrator,
to_start: &[String],
data_dir: &std::path::Path,
) -> Result<()> {
for name in to_start {
orchestrator.validate_start(name).await?;
}
let mut errors = Vec::new();
for (i, name) in to_start.iter().enumerate() {
if i > 0 {
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
}
repair_before_package_start(name).await;
wait_before_package_start(name).await;
let managed = crate::container::supervised_update::installed_unit(data_dir, name)?.is_some();
if !managed {
repair_before_package_start(name).await;
wait_before_package_start(name).await;
}
match orchestrator.start(name).await {
Ok(()) => wait_after_orchestrator_start(name).await,
Err(e) if is_unknown_app_id_error(&e) => {
Err(e) if !managed && is_unknown_app_id_error(&e) => {
// Collect instead of `?`: aborting here skipped every later
// stack member (mempool gate flakes #73/#74 — see
// order_present_containers).
@@ -902,7 +910,9 @@ async fn do_orchestrator_package_start(
Err(anyhow::anyhow!("Start failed: {}", errors.join("; ")))
} else {
for name in to_start {
ensure_runtime_host_port_listener(name).await?;
if crate::container::supervised_update::installed_unit(data_dir, name)?.is_none() {
ensure_runtime_host_port_listener(name).await?;
}
}
Ok(())
}
@@ -1131,12 +1141,14 @@ async fn podman_start_container(container_name: &str) -> Result<Output> {
async fn do_orchestrator_package_stop(
orchestrator: &dyn crate::container::traits::ContainerOrchestrator,
containers: &[String],
data_dir: &std::path::Path,
) -> Result<()> {
let mut errors = Vec::new();
for name in containers {
let managed = crate::container::supervised_update::installed_unit(data_dir, name)?.is_some();
match orchestrator.stop(name).await {
Ok(()) => {}
Err(e) if is_unknown_app_id_error(&e) => {
Err(e) if !managed && is_unknown_app_id_error(&e) => {
if let Err(e) = do_package_stop(&[name.clone()]).await {
errors.push(format!("{}: {:#}", name, e));
}
@@ -1145,7 +1157,7 @@ async fn do_orchestrator_package_stop(
// stop removes the --rm container before the orchestrator's
// defensive `podman stop` fallback probes it, and a stack member
// can be absent outright (stopped earlier, or never revived).
Err(e) if is_missing_container_error(&format!("{:#}", e)) => {
Err(e) if !managed && is_missing_container_error(&format!("{:#}", e)) => {
tracing::debug!(container = %name, error = %e, "stop: container already absent — treating as stopped");
}
Err(e) => {
@@ -1231,7 +1243,7 @@ async fn cascade_restart_address_caching_dependents(
let prev = flip_package_state(state_manager, dep, PackageState::Restarting).await;
let target = vec![dep.to_string()];
let result = if let Some(orch) = orchestrator {
do_orchestrator_package_restart(orch.as_ref(), &target).await
do_orchestrator_package_restart(orch.as_ref(), &target, data_dir).await
} else {
do_package_restart(&target).await
};
@@ -1299,9 +1311,14 @@ fn uses_single_orchestrator_app(package_id: &str) -> bool {
async fn do_orchestrator_package_restart(
orchestrator: &dyn crate::container::traits::ContainerOrchestrator,
to_restart: &[String],
data_dir: &std::path::Path,
) -> Result<()> {
do_orchestrator_package_stop(orchestrator, to_restart).await?;
do_orchestrator_package_start(orchestrator, to_restart).await
for name in to_restart {
orchestrator.validate_start(name).await?;
}
let stop_order: Vec<String> = to_restart.iter().rev().cloned().collect();
do_orchestrator_package_stop(orchestrator, &stop_order, data_dir).await?;
do_orchestrator_package_start(orchestrator, to_restart, data_dir).await
}
/// Stop all containers with their per-container graceful-shutdown timeout.
@@ -2248,6 +2248,17 @@ impl ProdContainerOrchestrator {
if crate::app_ops::lifecycle_op_in_flight(&c.name) {
continue;
}
// Reviewed managed data and units belong to the supervised
// transaction, including ownership changes and restarts. A bad
// record must fail closed rather than permit generic repair.
match super::supervised_update::installed_unit(&self.data_dir, &c.name) {
Ok(None) if !held_names.contains(&c.name) => {}
Ok(_) => continue,
Err(error) => {
report.failures.push((c.name.clone(), error.to_string()));
continue;
}
}
// Throttled: first pass after the container appears, then
// hourly — not on every 30s tick (see ownership_sweep_due).
if !ownership_sweep_due(&c.name) {
@@ -2377,6 +2388,27 @@ impl ProdContainerOrchestrator {
let lock = self.app_lock(&app_id).await;
let _guard = lock.lock().await;
let managed_name = compute_container_name(&lm.manifest);
if super::supervised_update::installed_unit(&self.data_dir, &managed_name)?.is_some() {
// A reviewed recipe owns this runtime, including its original writable
// layer. Guard BEFORE staged stops, secret/config hooks and every drift
// repair: install_fresh's later refusal cannot undo a prior removal.
let stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await;
let removed = crate::crash_recovery::load_user_uninstalled(&self.data_dir).await;
if stopped.contains(&app_id) || stopped.contains(&managed_name) {
return Ok(ReconcileAction::Left("user-stopped".into()));
}
if removed.contains(&app_id) || removed.contains(&managed_name) {
return Ok(ReconcileAction::Left("user-uninstalled".into()));
}
self.sync_quadlet_unit(lm, &managed_name).await?;
let status = self.runtime.get_container_status(&managed_name).await
.context("Reviewed managed runtime is missing; explicit recovery required")?;
anyhow::ensure!(matches!(status.state, ContainerState::Running),
"Reviewed managed runtime is not running; recover its saved systemd unit explicitly instead of recreating from the catalog");
return Ok(ReconcileAction::NoOp);
}
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?;
@@ -2935,6 +2967,10 @@ impl ProdContainerOrchestrator {
}
async fn stop_staged_runtime(&self, name: &str) -> Result<()> {
anyhow::ensure!(
super::supervised_update::installed_unit(&self.data_dir, name)?.is_none(),
"Reviewed managed runtime cannot be stopped by generic staged cleanup"
);
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:#}"));
@@ -5178,7 +5214,46 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
result
}
async fn validate_start(&self, app_id: &str) -> Result<()> {
let lm = match self.loaded(app_id).await {
Ok(lm) => lm,
Err(error) => {
anyhow::ensure!(
super::supervised_update::installed_unit(&self.data_dir, app_id)?.is_none(),
"Reviewed managed manifest is unavailable: {error}"
);
return Ok(());
}
};
let name = compute_container_name(&lm.manifest);
if super::supervised_update::installed_unit(&self.data_dir, &name)?.is_some() {
anyhow::ensure!(!super::update_transaction::is_held(&self.data_dir, &name)?,
"Reviewed managed runtime is held for update recovery");
self.sync_quadlet_unit(&lm, &name).await?;
self.ensure_resolved_source_available(&lm).await?;
}
Ok(())
}
async fn start(&self, app_id: &str) -> Result<()> {
if let Ok(lm) = self.loaded(app_id).await {
let name = compute_container_name(&lm.manifest);
if super::supervised_update::installed_unit(&self.data_dir, &name)?.is_some() {
let lock = self.app_lock(app_id).await;
let _guard = lock.lock().await;
self.validate_start(app_id).await?;
// Explicit owner start is allowed, but only systemd may use the
// exact retained recipe. Never render the current catalog here.
quadlet::enable_now(&format!("{name}.service")).await?;
wait_for_container_stable_running(self.runtime.as_ref(), &name, 3, 90).await?;
self.state.write().await.disabled.remove(app_id);
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_uninstalled(&self.data_dir, app_id).await;
crate::crash_recovery::mark_installed(&self.data_dir, app_id).await;
return Ok(());
}
}
if self.start_staged_update(app_id).await? {
return Ok(());
}
@@ -5252,6 +5327,19 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
let lock = self.app_lock(app_id).await;
let _guard = lock.lock().await;
let name = compute_container_name(&lm.manifest);
if super::supervised_update::installed_unit(&self.data_dir, &name)?.is_some() {
anyhow::ensure!(!super::update_transaction::is_held(&self.data_dir, &name)?,
"Reviewed managed runtime is held for update recovery");
self.sync_quadlet_unit(&lm, &name).await?;
quadlet::stop_service(&format!("{name}.service")).await?;
if let Ok(status) = self.runtime.get_container_status(&name).await {
anyhow::ensure!(
matches!(status.state, ContainerState::Stopped | ContainerState::Exited | ContainerState::Created),
"Reviewed managed runtime is still active after systemd stop"
);
}
return Ok(());
}
// Per-app graceful-stop grace: manifest `stop_grace_secs` if declared,
// else the historical per-app table. Slow-to-SIGTERM apps (bitcoin-core
// 600s, lnd 330s, electrumx 300s, fedimint 60s…) otherwise get a too-short
@@ -5302,6 +5390,13 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
}
async fn restart(&self, app_id: &str) -> Result<()> {
if let Ok(lm) = self.loaded(app_id).await {
if super::supervised_update::installed_unit(&self.data_dir, &compute_container_name(&lm.manifest))?.is_some() {
self.validate_start(app_id).await?;
self.stop(app_id).await?;
return self.start(app_id).await;
}
}
if let Some(members) = self.mempool_umbrella_members(app_id).await {
tracing::info!(
app_id,
@@ -7868,6 +7963,86 @@ app:
assert_eq!(config, "key: new\n");
}
#[tokio::test]
async fn reviewed_runtime_survives_catalog_drift_and_refuses_repairs_before_mutation() {
use std::os::unix::fs::PermissionsExt;
for case in ["running", "stopped", "missing", "changed-unit", "missing-unit", "user-stopped", "user-uninstalled"] {
let rt = Arc::new(MockRuntime::default());
let orch = orch_with(rt.clone()).await;
let name = format!("managed-{}", uuid::Uuid::new_v4().simple());
let mut manifest = pull_manifest(&name, "catalog:new");
manifest.app.environment = vec!["NEW_CATALOG_ENV=changed".into()];
orch.insert_manifest_for_test(manifest, PathBuf::from("/tmp/catalog-drift")).await;
let body = "[Container]\nImage=original:retained\nEnvironment=OLD_ENV=preserved\nPublishPort=127.0.0.1:1234:80\n";
let records = orch.data_dir.join("update-transactions/installed-units");
std::fs::create_dir_all(&records).unwrap();
std::fs::write(records.join(format!("{name}.json")), serde_json::to_vec(&serde_json::json!({
"schema": 1, "operation": uuid::Uuid::new_v4().to_string(),
"name": name, "body": body, "mode": 0o600
})).unwrap()).unwrap();
let unit = quadlet::unit_dir().await.unwrap().join(format!("{name}.container"));
if case != "missing-unit" {
std::fs::write(&unit, if case == "changed-unit" { "operator changed" } else { body }).unwrap();
std::fs::set_permissions(&unit, std::fs::Permissions::from_mode(0o600)).unwrap();
}
if case != "missing" {
rt.set_state(&name, if case == "stopped" { ContainerState::Stopped } else { ContainerState::Running });
}
if case == "user-stopped" {
crate::crash_recovery::mark_user_stopped(&orch.data_dir, &name).await;
}
if case == "user-uninstalled" {
crate::crash_recovery::mark_user_uninstalled(&orch.data_dir, &name).await;
}
// Even a conflicting generic staged record may not tear down the
// admitted runtime before ensure_running gets a chance to protect it.
assert!(orch.stop_staged_runtime(&name).await.is_err());
assert!(rt.calls().is_empty());
let before = rt.containers.lock().unwrap().clone();
let lm = orch.loaded(&name).await.unwrap();
let result = orch.ensure_running(&lm).await;
if case == "running" {
assert_eq!(result.unwrap(), ReconcileAction::NoOp);
} else if case.starts_with("user-") {
assert_eq!(result.unwrap(), ReconcileAction::Left(case.into()));
assert!(rt.calls().is_empty());
} else {
assert!(result.is_err(), "{case} must refuse before mutation");
}
assert_eq!(*rt.containers.lock().unwrap(), before, "{case}");
assert!(rt.calls().iter().all(|call| call.starts_with("get_container_status:")), "{case}: {:?}", rt.calls());
if case == "running" {
assert_eq!(std::fs::read_to_string(&unit).unwrap(), body);
}
if case == "running" || case.starts_with("user-") {
rt.mark_image_present("original:retained");
let holds = orch.data_dir.join("update-transactions/holds");
std::fs::create_dir_all(&holds).unwrap();
std::fs::write(holds.join(&name), uuid::Uuid::new_v4().to_string()).unwrap();
assert!(orch.validate_start(&name).await.is_err());
assert!(orch.start(&name).await.is_err());
assert!(orch.stop(&name).await.is_err());
std::fs::remove_file(holds.join(&name)).unwrap();
orch.validate_start(&name).await.unwrap();
orch.start(&name).await.unwrap();
assert_eq!(std::fs::read_to_string(&unit).unwrap(), body);
assert!(!crate::crash_recovery::load_user_stopped(&orch.data_dir).await.contains(&name));
assert!(!crate::crash_recovery::load_user_uninstalled(&orch.data_dir).await.contains(&name));
} else {
// Missing images/units and modified units refuse an explicit
// start before service mutation; no catalog pull is attempted.
assert!(orch.validate_start(&name).await.is_err());
assert!(orch.start(&name).await.is_err());
if case == "missing-unit" || case == "changed-unit" {
assert!(orch.stop(&name).await.is_err());
}
}
assert!(rt.calls().iter().all(|call| call.starts_with("get_container_status:") || call.starts_with("image_exists:")), "{case}: {:?}", rt.calls());
assert_eq!(*rt.containers.lock().unwrap(), before, "{case}");
let _ = std::fs::remove_file(unit);
}
}
#[tokio::test]
async fn reconcile_noop_when_already_running() {
let rt = Arc::new(MockRuntime::default());
+6
View File
@@ -45,6 +45,12 @@ pub trait ContainerOrchestrator: Send + Sync {
Ok(0)
}
/// Read-only admission for explicit start/restart. Managed implementations
/// validate saved units and locally available images before any stack stop.
async fn validate_start(&self, _app_id: &str) -> Result<()> {
Ok(())
}
/// Start an already-created container.
async fn start(&self, app_id: &str) -> Result<()>;
+62 -3
View File
@@ -728,7 +728,21 @@ fn parse_health_from_status(status: &str) -> Option<String> {
/// Try to recover a container. Running containers need a real restart so
/// rootless network helpers such as pasta are recreated; `podman start` is a
/// no-op for a running container with a missing host listener.
async fn restart_container(name: &str, state: &str) -> bool {
fn automatic_recovery_allowed(data_dir: &Path, name: &str) -> bool {
match (
crate::container::supervised_update::installed_unit(data_dir, name),
crate::container::update_transaction::is_held(data_dir, name),
) {
(Ok(None), Ok(false)) => true,
_ => false, // Saved, held, or unreadable recovery evidence fails closed.
}
}
async fn restart_container(name: &str, state: &str, data_dir: &Path) -> bool {
if !automatic_recovery_allowed(data_dir, name) {
warn!(container = %name, "Automatic restart refused: reviewed managed runtime needs explicit recovery");
return false;
}
let action = if state == "running" {
"restart"
} else {
@@ -927,7 +941,7 @@ pub fn spawn_health_monitor(state: Arc<StateManager>, data_dir: PathBuf) {
}
if matches!(
pkg.state,
PackageState::Starting | PackageState::Stopping | PackageState::Restarting
PackageState::Starting | PackageState::Stopping | PackageState::Restarting | PackageState::Updating
) {
debug!(
"Skipping container during package lifecycle transition: {} ({:?})",
@@ -1038,6 +1052,24 @@ pub fn spawn_health_monitor(state: Arc<StateManager>, data_dir: PathBuf) {
let mut prev_tier: Option<StartupTier> = None;
for container in &unhealthy {
if !automatic_recovery_allowed(&data_dir, &container.name) {
let id = format!("health-managed-{}", container.name);
if !data.notifications.iter().any(|n| n.id == id) {
data.notifications.push(Notification {
id,
level: NotificationLevel::Error,
title: format!("{} is unhealthy", container.app_id),
message: "Automatic restart is paused to preserve the installed configuration. Check the app before using its Start or Restart controls.".into(),
timestamp: chrono::Utc::now().to_rfc3339(),
app_id: Some(container.app_id.clone()),
});
if data.notifications.len() > 20 {
data.notifications = data.notifications.split_off(data.notifications.len() - 20);
}
state_changed = true;
}
continue;
}
let tier = container_tier(&container.name);
// Reset counter after 1 hour for permanently failed containers
@@ -1127,7 +1159,7 @@ pub fn spawn_health_monitor(state: Arc<StateManager>, data_dir: PathBuf) {
// the restart resyncs cleanly instead of crash-looping.
maybe_recover_corrupt_electrumx(&container.name, attempt).await;
let restarted = restart_container(&container.name, &container.state).await;
let restarted = restart_container(&container.name, &container.state, &data_dir).await;
if !restarted || attempt >= MAX_RESTART_ATTEMPTS {
let notification = Notification {
@@ -1197,6 +1229,33 @@ pub fn spawn_health_monitor(state: Arc<StateManager>, data_dir: PathBuf) {
mod tests {
use super::*;
#[tokio::test]
async fn automatic_recovery_never_restarts_saved_held_or_damaged_managed_runtime() {
let root = tempfile::tempdir().unwrap();
let name = "indeedhub-api";
assert!(automatic_recovery_allowed(root.path(), name));
let installed = root.path().join("update-transactions/installed-units");
std::fs::create_dir_all(&installed).unwrap();
let record = installed.join(format!("{name}.json"));
std::fs::write(&record, serde_json::to_vec(&serde_json::json!({
"schema": 1, "operation": uuid::Uuid::new_v4().to_string(),
"name": name, "body": "[Container]\nImage=original:retained\n", "mode": 0o600
})).unwrap()).unwrap();
assert!(!automatic_recovery_allowed(root.path(), name));
assert!(!restart_container(name, "running", root.path()).await);
std::fs::write(&record, b"damaged").unwrap();
assert!(!automatic_recovery_allowed(root.path(), name));
assert!(!restart_container(name, "stopped", root.path()).await);
std::fs::remove_file(record).unwrap();
let holds = root.path().join("update-transactions/holds");
std::fs::create_dir_all(&holds).unwrap();
std::fs::write(holds.join(name), uuid::Uuid::new_v4().to_string()).unwrap();
assert!(!automatic_recovery_allowed(root.path(), name));
assert!(!restart_container(name, "running", root.path()).await);
std::fs::write(holds.join(name), b"damaged").unwrap();
assert!(!automatic_recovery_allowed(root.path(), name));
}
#[test]
fn test_restart_tracker_new_is_empty() {
let tracker = RestartTracker::new();