Recover stopped update staging draft on current private-image preflight
This commit is contained in:
@@ -362,7 +362,11 @@ impl RpcHandler {
|
|||||||
set_package_state(
|
set_package_state(
|
||||||
&handler.state_manager,
|
&handler.state_manager,
|
||||||
&package_id_spawn,
|
&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)
|
pre_state.clone().unwrap_or(PackageState::Running)
|
||||||
} else {
|
} else {
|
||||||
PackageState::Running
|
PackageState::Running
|
||||||
|
|||||||
@@ -526,7 +526,18 @@ pub(in crate::api::rpc) async fn get_containers_for_app(package_id: &str) -> Res
|
|||||||
.await
|
.await
|
||||||
.context("podman ps timed out while listing containers")?
|
.context("podman ps timed out while listing containers")?
|
||||||
.context("Failed to list 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 all: Vec<&str> = stdout.lines().filter(|s| !s.is_empty()).collect();
|
||||||
|
|
||||||
let patterns = all_container_names(package_id);
|
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 {
|
mod tests {
|
||||||
use super::{all_container_names, get_data_dirs_for_app, get_health_check_args};
|
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]
|
#[test]
|
||||||
fn bitcoin_variant_container_names_are_precise() {
|
fn bitcoin_variant_container_names_are_precise() {
|
||||||
let core = all_container_names("bitcoin-core");
|
let core = all_container_names("bitcoin-core");
|
||||||
|
|||||||
@@ -126,7 +126,19 @@ impl RpcHandler {
|
|||||||
Err(e) => {
|
Err(e) => {
|
||||||
tracing::error!("package.start {} failed: {:#}", package_id_owned, e);
|
tracing::error!("package.start {} failed: {:#}", package_id_owned, e);
|
||||||
install_log(&format!("START FAIL: {} — {:#}", package_id_owned, e)).await;
|
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;
|
set_package_state(&state_manager, &package_id_owned, prev).await;
|
||||||
} else {
|
} else {
|
||||||
warn!(
|
warn!(
|
||||||
|
|||||||
@@ -4,7 +4,7 @@
|
|||||||
//! remove old container(s) → recreate (orchestrator-first, legacy fallback) → verify running.
|
//! remove old container(s) → recreate (orchestrator-first, legacy fallback) → verify running.
|
||||||
//! Data volumes are preserved (bind mounts, not stored in container).
|
//! 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::install::install_log;
|
||||||
use super::progress::parse_pull_progress;
|
use super::progress::parse_pull_progress;
|
||||||
use super::runtime::stop_timeout_secs;
|
use super::runtime::stop_timeout_secs;
|
||||||
@@ -13,6 +13,7 @@ use crate::api::rpc::RpcHandler;
|
|||||||
use crate::container::image_versions;
|
use crate::container::image_versions;
|
||||||
use crate::data_model::{InstallPhase, PackageState};
|
use crate::data_model::{InstallPhase, PackageState};
|
||||||
use anyhow::{Context, Result};
|
use anyhow::{Context, Result};
|
||||||
|
use std::{collections::HashSet, path::Path};
|
||||||
use tokio::io::{AsyncBufReadExt, BufReader};
|
use tokio::io::{AsyncBufReadExt, BufReader};
|
||||||
use tracing::{error, info, warn};
|
use tracing::{error, info, warn};
|
||||||
|
|
||||||
@@ -60,9 +61,24 @@ impl RpcHandler {
|
|||||||
let targets = pinned
|
let targets = pinned
|
||||||
.as_ref()
|
.as_ref()
|
||||||
.map(|target| self.resolve_images_to_pull(package_id, target));
|
.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 {
|
if let Some(targets) = &targets {
|
||||||
let installed = inspect_update_images(package_id).await?;
|
if !update_targets_need_change(targets, &installed)? && !markers.stopped {
|
||||||
if !update_targets_need_change(targets, &installed)? {
|
|
||||||
install_log(&format!(
|
install_log(&format!(
|
||||||
"UPDATE SKIP: {} — target versions already installed",
|
"UPDATE SKIP: {} — target versions already installed",
|
||||||
package_id
|
package_id
|
||||||
@@ -111,6 +127,19 @@ impl RpcHandler {
|
|||||||
if let Some(orchestrator) = self.orchestrator.as_ref() {
|
if let Some(orchestrator) = self.orchestrator.as_ref() {
|
||||||
match orchestrator.upgrade(orchestrator_app_id).await {
|
match orchestrator.upgrade(orchestrator_app_id).await {
|
||||||
Ok(()) => {
|
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 {
|
if let Some(targets) = &targets {
|
||||||
verify_update_targets(
|
verify_update_targets(
|
||||||
targets,
|
targets,
|
||||||
@@ -574,13 +603,69 @@ impl RpcHandler {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn inspect_update_images(package_id: &str) -> Result<Vec<(String, String)>> {
|
#[derive(Default)]
|
||||||
let containers = get_containers_for_app(package_id).await?;
|
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!(
|
anyhow::ensure!(
|
||||||
!containers.is_empty(),
|
!markers.uninstalled,
|
||||||
"No containers found for {}",
|
"{} was uninstalled; use an explicit install instead of update",
|
||||||
package_id
|
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");
|
let mut command = tokio::process::Command::new("podman");
|
||||||
command.arg("inspect").args(&containers).kill_on_drop(true);
|
command.arg("inspect").args(&containers).kill_on_drop(true);
|
||||||
let output = tokio::time::timeout(std::time::Duration::from_secs(30), command.output())
|
let output = tokio::time::timeout(std::time::Duration::from_secs(30), command.output())
|
||||||
@@ -626,6 +711,7 @@ fn verify_update_targets(
|
|||||||
targets: &[(String, String)],
|
targets: &[(String, String)],
|
||||||
installed: &[(String, String)],
|
installed: &[(String, String)],
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
|
anyhow::ensure!(!targets.is_empty(), "No update targets resolved");
|
||||||
for (app_id, target) in targets {
|
for (app_id, target) in targets {
|
||||||
let running = installed_image_for_target(app_id, installed).ok_or_else(|| {
|
let running = installed_image_for_target(app_id, installed).ok_or_else(|| {
|
||||||
anyhow::anyhow!("Update {}: target container missing after recreate", app_id)
|
anyhow::anyhow!("Update {}: target container missing after recreate", app_id)
|
||||||
@@ -651,6 +737,7 @@ fn update_targets_need_change(
|
|||||||
installed: &[(String, String)],
|
installed: &[(String, String)],
|
||||||
) -> Result<bool> {
|
) -> Result<bool> {
|
||||||
use std::cmp::Ordering;
|
use std::cmp::Ordering;
|
||||||
|
anyhow::ensure!(!targets.is_empty(), "No update targets resolved");
|
||||||
let mut changed = false;
|
let mut changed = false;
|
||||||
for (app_id, target) in targets {
|
for (app_id, target) in targets {
|
||||||
let running = installed_image_for_target(app_id, installed);
|
let running = installed_image_for_target(app_id, installed);
|
||||||
@@ -771,9 +858,83 @@ mod tests {
|
|||||||
use super::{
|
use super::{
|
||||||
candidate_app_ids_for_container, immutable_update_image, orchestrator_update_app_id,
|
candidate_app_ids_for_container, immutable_update_image, orchestrator_update_app_id,
|
||||||
should_try_orchestrator_update, update_targets_need_change, uses_legacy_update_flow,
|
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]
|
#[tokio::test]
|
||||||
async fn stack_image_failure_precedes_every_lifecycle_action() {
|
async fn stack_image_failure_precedes_every_lifecycle_action() {
|
||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex};
|
||||||
|
|||||||
@@ -298,7 +298,22 @@ impl DockerPackageScanner {
|
|||||||
uninstall_stage: None,
|
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);
|
packages.insert(app_id.clone(), package);
|
||||||
info!(
|
info!(
|
||||||
"Detected container: {} ({})",
|
"Detected container: {} ({})",
|
||||||
|
|||||||
@@ -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<()> {
|
async fn run_step(step: &HookStep, container: &str, app_id: &str, data_dir: &Path) -> Result<()> {
|
||||||
match step {
|
match step {
|
||||||
HookStep::Exec { exec } => {
|
HookStep::Exec { exec } => {
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
pub mod app_catalog;
|
pub mod app_catalog;
|
||||||
pub mod node_catalog;
|
|
||||||
pub mod app_gate_config;
|
pub mod app_gate_config;
|
||||||
pub mod bitcoin_ui;
|
pub mod bitcoin_ui;
|
||||||
pub mod boot_reconciler;
|
pub mod boot_reconciler;
|
||||||
@@ -14,11 +13,12 @@ pub mod image_policy;
|
|||||||
pub mod image_versions;
|
pub mod image_versions;
|
||||||
pub mod lnd;
|
pub mod lnd;
|
||||||
pub mod migration_backup;
|
pub mod migration_backup;
|
||||||
|
pub mod node_catalog;
|
||||||
pub mod npm;
|
pub mod npm;
|
||||||
pub mod prod_orchestrator;
|
pub mod prod_orchestrator;
|
||||||
pub mod quadlet;
|
pub mod quadlet;
|
||||||
pub mod registry;
|
|
||||||
pub mod registration_pin;
|
pub mod registration_pin;
|
||||||
|
pub mod registry;
|
||||||
pub mod secrets;
|
pub mod secrets;
|
||||||
pub mod traits;
|
pub mod traits;
|
||||||
pub mod ui_detection;
|
pub mod ui_detection;
|
||||||
@@ -29,3 +29,5 @@ pub use dev_orchestrator::DevContainerOrchestrator;
|
|||||||
pub use docker_packages::DockerPackageScanner;
|
pub use docker_packages::DockerPackageScanner;
|
||||||
pub use prod_orchestrator::ProdContainerOrchestrator;
|
pub use prod_orchestrator::ProdContainerOrchestrator;
|
||||||
pub use traits::ContainerOrchestrator;
|
pub use traits::ContainerOrchestrator;
|
||||||
|
|
||||||
|
mod staged_update;
|
||||||
|
|||||||
@@ -1422,6 +1422,17 @@ fn host_port_collisions<'m>(
|
|||||||
out
|
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
|
/// Internal: track a manifest together with the absolute directory it was loaded
|
||||||
/// from, so Build sources can resolve relative `context:` paths.
|
/// from, so Build sources can resolve relative `context:` paths.
|
||||||
#[derive(Debug, Clone)]
|
#[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 disk_gb = self.disk_gb().await;
|
||||||
let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB
|
let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB
|
||||||
|| crate::settings::bitcoin_storage::load(&self.data_dir)
|
|| crate::settings::bitcoin_storage::load(&self.data_dir)
|
||||||
@@ -2336,6 +2349,12 @@ impl ProdContainerOrchestrator {
|
|||||||
let lock = self.app_lock(&app_id).await;
|
let lock = self.app_lock(&app_id).await;
|
||||||
let _guard = lock.lock().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?;
|
self.ensure_app_secrets(&app_id).await?;
|
||||||
|
|
||||||
// Don't fight the Bitcoin-implementation switch: bitcoin-core and
|
// 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.
|
/// 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<()> {
|
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?;
|
self.ensure_app_secrets(&lm.manifest.app.id).await?;
|
||||||
let mut resolved_manifest = lm.manifest.clone();
|
let mut resolved_manifest = lm.manifest.clone();
|
||||||
self.resolve_dynamic_env(&mut resolved_manifest).await?;
|
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(|| {
|
let resolved = resolved_manifest.app.container.resolve().ok_or_else(|| {
|
||||||
anyhow::anyhow!(
|
anyhow::anyhow!(
|
||||||
@@ -3008,7 +3187,17 @@ impl ProdContainerOrchestrator {
|
|||||||
// freshly created container — exactly when container mutations (e.g.
|
// freshly created container — exactly when container mutations (e.g.
|
||||||
// indeedhub's nginx X-Frame-Options strip + nostr-provider injection) must
|
// indeedhub's nginx X-Frame-Options strip + nostr-provider injection) must
|
||||||
// be re-applied. Best-effort + idempotent: never fails the install.
|
// 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 uses_pasta_network(&resolved_manifest) {
|
||||||
if let Err(err) = wait_for_manifest_host_ports(
|
if let Err(err) = wait_for_manifest_host_ports(
|
||||||
&resolved_manifest,
|
&resolved_manifest,
|
||||||
@@ -4825,6 +5014,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn install(&self, app_id: &str) -> Result<String> {
|
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?;
|
let lm = self.loaded(app_id).await?;
|
||||||
// Optional shared-service preconditions are checked before recording
|
// Optional shared-service preconditions are checked before recording
|
||||||
// installation or creating anything. A headless adapter must not claim
|
// installation or creating anything. A headless adapter must not claim
|
||||||
@@ -4896,7 +5087,7 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
|||||||
// point). Just delegate.
|
// point). Just delegate.
|
||||||
self.prepare_for_start(&lm.manifest).await?;
|
self.prepare_for_start(&lm.manifest).await?;
|
||||||
let action = self.ensure_running(&lm).await?;
|
let action = self.ensure_running(&lm).await?;
|
||||||
match action {
|
let result = match action {
|
||||||
ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed => {
|
ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed => {
|
||||||
Ok(name)
|
Ok(name)
|
||||||
}
|
}
|
||||||
@@ -4915,10 +5106,17 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
|||||||
self.install_fresh(&lm).await?;
|
self.install_fresh(&lm).await?;
|
||||||
Ok(name)
|
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<()> {
|
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 {
|
if let Some(members) = self.mempool_umbrella_members(app_id).await {
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
app_id,
|
app_id,
|
||||||
@@ -5136,6 +5334,8 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
|||||||
self.remove_quadlet_unit_if_present(&name).await?;
|
self.remove_quadlet_unit_if_present(&name).await?;
|
||||||
}
|
}
|
||||||
self.state.write().await.disabled.insert(app_id.to_string());
|
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(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
let lm = self.loaded(app_id).await?;
|
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
|
// stale claim behind would let desired-state recovery recreate the very
|
||||||
// app that was just uninstalled.
|
// app that was just uninstalled.
|
||||||
crate::crash_recovery::clear_installed(&self.data_dir, app_id).await;
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Upgrade: stop-remove-reinstall (re-pulls or rebuilds as required).
|
/// Upgrade: stop-remove-reinstall (re-pulls or rebuilds as required).
|
||||||
async fn upgrade(&self, app_id: &str) -> Result<()> {
|
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 lock = self.app_lock(app_id).await;
|
||||||
let _guard = lock.lock().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);
|
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();
|
let mut resolved = lm.manifest.clone();
|
||||||
resolve_catalog_image(&mut resolved);
|
resolve_catalog_image(&mut resolved);
|
||||||
if resolved.app.container.build.is_none() {
|
if resolved.app.container.build.is_none() {
|
||||||
@@ -5213,7 +5481,10 @@ impl ContainerOrchestrator for ProdContainerOrchestrator {
|
|||||||
running.image,
|
running.image,
|
||||||
target
|
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.stop_container(&name).await;
|
||||||
let _ = self.runtime.remove_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> {
|
async fn status(&self, app_id: &str) -> Result<ContainerStatus> {
|
||||||
@@ -5778,6 +6064,7 @@ mod tests {
|
|||||||
fail_image_exists: StdMutex<HashMap<String, String>>,
|
fail_image_exists: StdMutex<HashMap<String, String>>,
|
||||||
/// If set, `start_container` for this container fails with this message.
|
/// If set, `start_container` for this container fails with this message.
|
||||||
fail_start: StdMutex<HashMap<String, String>>,
|
fail_start: StdMutex<HashMap<String, String>>,
|
||||||
|
fail_stop: StdMutex<HashMap<String, String>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl MockRuntime {
|
impl MockRuntime {
|
||||||
@@ -5841,6 +6128,12 @@ mod tests {
|
|||||||
) -> Result<String> {
|
) -> Result<String> {
|
||||||
self.record(format!("create_container:{name}:offset={port_offset}"));
|
self.record(format!("create_container:{name}:offset={port_offset}"));
|
||||||
self.set_state(name, ContainerState::Created);
|
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
|
self.created_env
|
||||||
.lock()
|
.lock()
|
||||||
.unwrap()
|
.unwrap()
|
||||||
@@ -5862,6 +6155,9 @@ mod tests {
|
|||||||
}
|
}
|
||||||
async fn stop_container(&self, name: &str) -> Result<()> {
|
async fn stop_container(&self, name: &str) -> Result<()> {
|
||||||
self.record(format!("stop_container:{name}"));
|
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);
|
self.set_state(name, ContainerState::Stopped);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -6310,6 +6606,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]
|
#[tokio::test]
|
||||||
async fn install_fresh_pull() {
|
async fn install_fresh_pull() {
|
||||||
let rt = Arc::new(MockRuntime::default());
|
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)
|
||||||
|
}
|
||||||
@@ -62,6 +62,11 @@ pub trait ContainerOrchestrator: Send + Sync {
|
|||||||
/// Pull/rebuild the image and recreate the container from scratch.
|
/// Pull/rebuild the image and recreate the container from scratch.
|
||||||
async fn upgrade(&self, app_id: &str) -> Result<()>;
|
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.
|
/// Current state of a single container.
|
||||||
async fn status(&self, app_id: &str) -> Result<ContainerStatus>;
|
async fn status(&self, app_id: &str) -> Result<ContainerStatus>;
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user