diff --git a/core/archipelago/src/api/rpc/package/update.rs b/core/archipelago/src/api/rpc/package/update.rs index 037f0f93..61379e10 100644 --- a/core/archipelago/src/api/rpc/package/update.rs +++ b/core/archipelago/src/api/rpc/package/update.rs @@ -305,7 +305,22 @@ impl RpcHandler { guard: &crate::container::update_transaction::Guard, ) -> Result<()> { use crate::container::update_transaction::{self, Podman}; - let targets = Podman::targets(images_to_pull).await?; + 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() @@ -316,7 +331,22 @@ impl RpcHandler { ); self.set_install_phase(package_id, InstallPhase::Preparing) .await; - update_transaction::execute(guard, package_id, &targets, &Podman).await + 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?; + crate::container::supervised_update::execute(guard, package_id, &targets, &adapter) + .await + } else { + update_transaction::execute(guard, package_id, &targets, &Podman).await + } } async fn recreate_container_for_update( diff --git a/core/archipelago/src/appgate/mod.rs b/core/archipelago/src/appgate/mod.rs index d82fef6d..dbdd1eed 100644 --- a/core/archipelago/src/appgate/mod.rs +++ b/core/archipelago/src/appgate/mod.rs @@ -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//` before this gate @@ -838,6 +848,16 @@ dashboard and check {name} under My Apps.

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