diff --git a/core/archipelago/src/api/rpc/package/async_lifecycle.rs b/core/archipelago/src/api/rpc/package/async_lifecycle.rs index 65a80ce2..36172e87 100644 --- a/core/archipelago/src/api/rpc/package/async_lifecycle.rs +++ b/core/archipelago/src/api/rpc/package/async_lifecycle.rs @@ -89,6 +89,15 @@ impl RpcHandler { match handler.handle_package_install(params).await { Ok(_) => { info!("package.install {}: complete", package_id_spawn); + for id in [&package_id_spawn, &format!("archy-{}", package_id_spawn)] { + crate::crash_recovery::clear_user_uninstalled(&handler.config.data_dir, id) + .await; + } + crate::crash_recovery::mark_installed( + &handler.config.data_dir, + &package_id_spawn, + ) + .await; // The install pipeline has verified the container is up // and healthy (see install.rs post-start exit check). // Kick the scanner first so the fresh manifest (with @@ -184,17 +193,20 @@ impl RpcHandler { // phase is cleared (None) so no stale InstallPhase // lingers on the card. let err_msg = format!("Install failed: {:#}", e); - let (mut data, _) = handler.state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(&package_id_spawn) { - entry.state = PackageState::Stopped; - entry.install_progress = Some(crate::data_model::InstallProgress { - size: 0, - downloaded: 0, - phase: None, - message: Some(err_msg), - }); - handler.state_manager.update_data(data).await; - } + handler + .state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(&package_id_spawn) { + entry.state = PackageState::Stopped; + entry.install_progress = Some(crate::data_model::InstallProgress { + size: 0, + downloaded: 0, + phase: None, + message: Some(err_msg), + }); + } + }) + .await; } } }); @@ -252,6 +264,11 @@ impl RpcHandler { match handler.handle_package_uninstall(params).await { Ok(_) => { info!("package.uninstall {}: complete", package_id_spawn); + for id in [&package_id_spawn, &format!("archy-{}", package_id_spawn)] { + crate::crash_recovery::mark_user_uninstalled(&handler.config.data_dir, id) + .await; + crate::crash_recovery::clear_installed(&handler.config.data_dir, id).await; + } // Inner handler already removed the package entry on // success. Nothing more to do here. } @@ -382,52 +399,56 @@ impl RpcHandler { /// call, but fires before the spawn so the UI sees it immediately. async fn flip_to_installing(state_manager: &StateManager, package_id: &str) { use crate::data_model::{Description, Manifest, PackageDataEntry, StaticFiles}; - let (mut data, _) = state_manager.get_snapshot().await; - let entry = data - .package_data - .entry(package_id.to_string()) - .or_insert_with(|| PackageDataEntry { - state: PackageState::Installing, - health: None, - exit_code: None, - static_files: StaticFiles { - license: String::new(), - instructions: String::new(), - // Leave icon empty during the transient Installing window: - // hardcoding `.png` is wrong for ~half our apps (many use - // `.svg` / `.webp`), producing a broken-image flicker until - // the scanner refreshes the entry. The frontend's `icon` - // computed falls through to `curatedMap.get(id)?.icon` which - // has the correct extensions for known apps. - icon: String::new(), - }, - manifest: Manifest { - id: package_id.to_string(), - title: package_id.to_string(), - version: String::new(), - description: Description { - short: "Installing...".to_string(), - long: String::new(), - }, - release_notes: String::new(), - license: String::new(), - wrapper_repo: String::new(), - upstream_repo: String::new(), - support_site: String::new(), - marketing_site: String::new(), - donation_url: None, - author: None, - website: None, - interfaces: None, - tier: None, - }, - installed: None, - install_progress: None, - uninstall_stage: None, - available_update: None, - }); - entry.state = PackageState::Installing; - state_manager.update_data(data).await; + state_manager + .mutate_data(|data| { + let entry = data + .package_data + .entry(package_id.to_string()) + .or_insert_with(|| PackageDataEntry { + ui_ready: None, + state: PackageState::Installing, + health: None, + exit_code: None, + static_files: StaticFiles { + license: String::new(), + instructions: String::new(), + // Leave icon empty during the transient Installing window: + // hardcoding `.png` is wrong for ~half our apps (many use + // `.svg` / `.webp`), producing a broken-image flicker until + // the scanner refreshes the entry. The frontend's `icon` + // computed falls through to `curatedMap.get(id)?.icon` which + // has the correct extensions for known apps. + icon: String::new(), + }, + manifest: Manifest { + id: package_id.to_string(), + title: package_id.to_string(), + version: String::new(), + description: Description { + short: "Installing...".to_string(), + long: String::new(), + }, + release_notes: String::new(), + license: String::new(), + wrapper_repo: String::new(), + upstream_repo: String::new(), + support_site: String::new(), + marketing_site: String::new(), + donation_url: None, + author: None, + website: None, + interfaces: None, + tier: None, + }, + installed: None, + install_progress: None, + uninstall_stage: None, + available_update: None, + }); + entry.ui_ready = Some(false); + entry.state = PackageState::Installing; + }) + .await; } /// True when the failed install still has a real footprint: any container @@ -485,20 +506,23 @@ async fn remove_entry_with_notification( id_prefix: &str, message: &str, ) { - let (mut data, _) = handler.state_manager.get_snapshot().await; - data.package_data.remove(package_id); - data.notifications.push(crate::data_model::Notification { - id: format!("{id_prefix}-{package_id}"), - level: crate::data_model::NotificationLevel::Error, - title: format!("Could not install {package_id}"), - message: message.to_string(), - timestamp: chrono::Utc::now().to_rfc3339(), - app_id: Some(package_id.to_string()), - }); - while data.notifications.len() > 20 { - data.notifications.remove(0); - } - handler.state_manager.update_data(data).await; + handler + .state_manager + .mutate_data(|data| { + data.package_data.remove(package_id); + data.notifications.push(crate::data_model::Notification { + id: format!("{id_prefix}-{package_id}"), + level: crate::data_model::NotificationLevel::Error, + title: format!("Could not install {package_id}"), + message: message.to_string(), + timestamp: chrono::Utc::now().to_rfc3339(), + app_id: Some(package_id.to_string()), + }); + while data.notifications.len() > 20 { + data.notifications.remove(0); + } + }) + .await; } /// Flip an existing entry's state and return the pre-flip value (or None if @@ -508,18 +532,23 @@ async fn flip_package_state( package_id: &str, new_state: PackageState, ) -> Option { - let (mut data, _) = state_manager.get_snapshot().await; - let prev = data.package_data.get(package_id).map(|e| e.state.clone()); - if let Some(entry) = data.package_data.get_mut(package_id) { - entry.state = new_state; - state_manager.update_data(data).await; - } else { - warn!( - "flip_package_state: no entry for {} — cannot flip", - package_id - ); - } - prev + state_manager + .mutate_data(|data| { + let prev = data.package_data.get(package_id).map(|e| e.state.clone()); + if let Some(entry) = data.package_data.get_mut(package_id) { + if new_state != PackageState::Running { + entry.ui_ready = Some(false); + } + entry.state = new_state; + } else { + warn!( + "flip_package_state: no entry for {} — cannot flip", + package_id + ); + } + prev + }) + .await } /// Set state unconditionally (no-op if entry no longer exists). @@ -528,13 +557,18 @@ async fn set_package_state( package_id: &str, new_state: PackageState, ) { - let (mut data, _) = state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(package_id) { - if entry.state != new_state { - entry.state = new_state; - state_manager.update_data(data).await; - } - } + state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(package_id) { + if entry.state != new_state { + if new_state != PackageState::Running { + entry.ui_ready = Some(false); + } + entry.state = new_state; + } + } + }) + .await } /// Set state and clear the uninstall_stage label. Used when an uninstall @@ -545,12 +579,17 @@ async fn set_package_state_and_clear_uninstall_stage( package_id: &str, new_state: PackageState, ) { - let (mut data, _) = state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(package_id) { - entry.state = new_state; - entry.uninstall_stage = None; - state_manager.update_data(data).await; - } + state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(package_id) { + if new_state != PackageState::Running { + entry.ui_ready = Some(false); + } + entry.state = new_state; + entry.uninstall_stage = None; + } + }) + .await } /// Kick the container scanner to run immediately and wait for it to finish diff --git a/core/archipelago/src/api/rpc/package/install.rs b/core/archipelago/src/api/rpc/package/install.rs index 86259562..a81208d3 100644 --- a/core/archipelago/src/api/rpc/package/install.rs +++ b/core/archipelago/src/api/rpc/package/install.rs @@ -545,7 +545,7 @@ impl RpcHandler { // Keep legacy install flow as default while migration is in progress. if orchestrator_managed { let orchestrator_app_id = orchestrator_install_app_id(package_id); - self.set_install_phase(package_id, InstallPhase::CreatingContainer) + self.set_install_phase(package_id, InstallPhase::PreparingApp) .await; install_log(&format!( "INSTALL ORCH: {} — attempting orchestrator install as {}", @@ -2053,25 +2053,8 @@ fn parse_setup_token(lines: &[&str]) -> Option { } async fn cleanup_stale_package_ports(package_id: &str) { - match package_id { - "grafana" => cleanup_stale_pasta_port("3000").await, - "homeassistant" | "home-assistant" => cleanup_stale_pasta_port("8123").await, - "searxng" => cleanup_stale_pasta_port("8888").await, - "uptime-kuma" => cleanup_stale_pasta_port("3002").await, - "gitea" => { - cleanup_stale_pasta_port("3001").await; - cleanup_stale_pasta_port("2222").await; - cleanup_stale_pasta_port("3000").await; - } - "nginx-proxy-manager" => { - cleanup_stale_pasta_port("8081").await; - cleanup_stale_pasta_port("8084").await; - cleanup_stale_pasta_port("8444").await; - } - "nextcloud" => cleanup_stale_pasta_port("8085").await, - "portainer" => cleanup_stale_pasta_port("9000").await, - _ => {} - } + // Never kill by port: another app or the management gate may own it. + crate::container::ghost_reaper::reap_for_app(package_id).await; } fn install_command_tail( @@ -2196,93 +2179,11 @@ async fn cleanup_start_conflict(package_id: &str, stderr: &str) -> bool { return true; } - match package_id { - "grafana" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("3000").await; - true - } - "homeassistant" | "home-assistant" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("8123").await; - true - } - "searxng" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("8888").await; - true - } - "uptime-kuma" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("3002").await; - true - } - "gitea" if stderr.contains("pasta failed") || stderr.contains("address already in use") => { - cleanup_stale_pasta_port("3001").await; - cleanup_stale_pasta_port("2222").await; - cleanup_stale_pasta_port("3000").await; - true - } - "nginx-proxy-manager" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("8081").await; - cleanup_stale_pasta_port("8084").await; - cleanup_stale_pasta_port("8444").await; - true - } - "nextcloud" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("8085").await; - true - } - "portainer" - if stderr.contains("pasta failed") || stderr.contains("address already in use") => - { - cleanup_stale_pasta_port("9000").await; - true - } - _ => false, + if stderr.contains("pasta failed") || stderr.contains("address already in use") { + crate::container::ghost_reaper::reap_for_app(package_id).await; + return true; } -} - -async fn cleanup_stale_pasta_port(port: &str) { - // NEVER kill our own process. The daemon holds catalog app ports over - // IPv6 (the mesh app-port relay), so a blunt `fuser -k /tcp` would - // terminate archipelago itself mid-install — installs failed and apps - // vanished on a test node 2026-07-27. Kill every listener on the port - // EXCEPT our PID (and our process group), leaving the relay/daemon alive. - let self_pid = std::process::id(); - let kill_listener = format!( - "ss -ltnp 'sport = :{port}' 2>/dev/null | sed -n 's/.*pid=\\([0-9]*\\).*/\\1/p' | \ - while read p; do [ \"$p\" = \"{self_pid}\" ] || kill \"$p\" 2>/dev/null; done || true", - ); - let _ = tokio::process::Command::new("sh") - .args(["-c", &kill_listener]) - .output() - .await; - - // sudo fuser -k, but exclude our own PID: fuser prints the PIDs holding - // the port; kill each except self. (`fuser -k` has no exclusion flag.) - let fuser_kill = format!( - "for p in $(sudo fuser {port}/tcp 2>/dev/null); do [ \"$p\" = \"{self_pid}\" ] || sudo kill \"$p\" 2>/dev/null; done || true", - ); - let _ = tokio::process::Command::new("sh") - .args(["-c", &fuser_kill]) - .output() - .await; - - let pattern = format!("pasta.*{}", port); - let _ = tokio::process::Command::new("pkill") - .args(["-f", &pattern]) - .output() - .await; - tokio::time::sleep(std::time::Duration::from_secs(1)).await; + false } async fn repair_nextcloud_permissions() { diff --git a/core/archipelago/src/api/rpc/package/progress.rs b/core/archipelago/src/api/rpc/package/progress.rs index 47767e5a..248c0428 100644 --- a/core/archipelago/src/api/rpc/package/progress.rs +++ b/core/archipelago/src/api/rpc/package/progress.rs @@ -14,20 +14,23 @@ impl RpcHandler { /// the rare case where the pull stream actually parses, but podman /// almost never emits parseable progress on a piped stderr. pub(super) async fn set_install_progress(&self, package_id: &str, downloaded: u64, size: u64) { - let (mut data, _rev) = self.state_manager.get_snapshot().await; - let entry = data - .package_data - .entry(package_id.to_string()) - .or_insert_with(|| create_installing_entry(package_id)); - entry.state = PackageState::Installing; - let existing_phase = entry.install_progress.as_ref().and_then(|p| p.phase); - entry.install_progress = Some(InstallProgress { - size, - downloaded, - phase: existing_phase, - message: None, - }); - self.state_manager.update_data(data).await; + self.state_manager + .mutate_data(|data| { + let entry = data + .package_data + .entry(package_id.to_string()) + .or_insert_with(|| create_installing_entry(package_id)); + entry.ui_ready = Some(false); + entry.state = PackageState::Installing; + let existing_phase = entry.install_progress.as_ref().and_then(|p| p.phase); + entry.install_progress = Some(InstallProgress { + size, + downloaded, + phase: existing_phase, + message: None, + }); + }) + .await; } /// Set the install pipeline phase and broadcast. This is the @@ -35,76 +38,86 @@ impl RpcHandler { /// percentage and a user-facing label. Byte counters are retained /// for the rare case podman emits parseable progress. pub(super) async fn set_install_phase(&self, package_id: &str, phase: InstallPhase) { - let (mut data, _rev) = self.state_manager.get_snapshot().await; - let entry = data - .package_data - .entry(package_id.to_string()) - .or_insert_with(|| create_installing_entry(package_id)); - // Preparing / PullingImage / CreatingContainer / StartingContainer / - // WaitingHealthy / PostInstall all map to the Installing state. - // Updates use Updating state — the wrapper has already flipped - // state to Updating, so don't clobber it. - if entry.state != PackageState::Updating { - entry.state = PackageState::Installing; - } - let (size, downloaded) = entry - .install_progress - .as_ref() - .map(|p| (p.size, p.downloaded)) - .unwrap_or((0, 0)); - entry.install_progress = Some(InstallProgress { - size, - downloaded, - phase: Some(phase), - message: None, - }); - self.state_manager.update_data(data).await; + self.state_manager + .mutate_data(|data| { + let entry = data + .package_data + .entry(package_id.to_string()) + .or_insert_with(|| create_installing_entry(package_id)); + // Preparing / PullingImage / CreatingContainer / StartingContainer / + // WaitingHealthy / PostInstall all map to the Installing state. + // Updates use Updating state — the wrapper has already flipped + // state to Updating, so don't clobber it. + if entry.state != PackageState::Updating { + entry.ui_ready = Some(false); + entry.state = PackageState::Installing; + } + let (size, downloaded) = entry + .install_progress + .as_ref() + .map(|p| (p.size, p.downloaded)) + .unwrap_or((0, 0)); + entry.install_progress = Some(InstallProgress { + size, + downloaded, + phase: Some(phase), + message: None, + }); + }) + .await; } /// Set a user-facing install status message (e.g. "Waiting for Bitcoin /// to start…") without disturbing the current phase/byte counters. pub(super) async fn set_install_message(&self, package_id: &str, message: &str) { - let (mut data, _rev) = self.state_manager.get_snapshot().await; - let entry = data - .package_data - .entry(package_id.to_string()) - .or_insert_with(|| create_installing_entry(package_id)); - if entry.state != PackageState::Updating { - entry.state = PackageState::Installing; - } - let (size, downloaded, phase) = entry - .install_progress - .as_ref() - .map(|p| (p.size, p.downloaded, p.phase)) - .unwrap_or((0, 0, None)); - entry.install_progress = Some(InstallProgress { - size, - downloaded, - phase, - message: Some(message.to_string()), - }); - self.state_manager.update_data(data).await; + self.state_manager + .mutate_data(|data| { + let entry = data + .package_data + .entry(package_id.to_string()) + .or_insert_with(|| create_installing_entry(package_id)); + if entry.state != PackageState::Updating { + entry.ui_ready = Some(false); + entry.state = PackageState::Installing; + } + let (size, downloaded, phase) = entry + .install_progress + .as_ref() + .map(|p| (p.size, p.downloaded, p.phase)) + .unwrap_or((0, 0, None)); + entry.install_progress = Some(InstallProgress { + size, + downloaded, + phase, + message: Some(message.to_string()), + }); + }) + .await; } /// Clear install progress after pull completes or fails. pub(super) async fn clear_install_progress(&self, package_id: &str) { - let (mut data, _rev) = self.state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(package_id) { - entry.install_progress = None; - } - self.state_manager.update_data(data).await; + self.state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(package_id) { + entry.install_progress = None; + } + }) + .await; } /// Set the uninstall stage label so the UI can show what's happening /// instead of a generic spinner. Each call broadcasts a state change /// — call sparingly (one per pipeline phase, not per container). pub(super) async fn set_uninstall_stage(&self, package_id: &str, stage: &str) { - let (mut data, _rev) = self.state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(package_id) { - entry.uninstall_stage = Some(stage.to_string()); - entry.state = crate::data_model::PackageState::Removing; - } - self.state_manager.update_data(data).await; + self.state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(package_id) { + entry.uninstall_stage = Some(stage.to_string()); + entry.state = crate::data_model::PackageState::Removing; + } + }) + .await; } /// Update install progress (static method for use in async closures). @@ -114,25 +127,28 @@ impl RpcHandler { downloaded: u64, total: u64, ) { - let (mut data, _rev) = state_manager.get_snapshot().await; - let entry = data - .package_data - .entry(package_id.to_string()) - .or_insert_with(|| create_installing_entry(package_id)); - let existing_phase = entry.install_progress.as_ref().and_then(|p| p.phase); - entry.install_progress = Some(InstallProgress { - size: total, - downloaded, - phase: existing_phase, - message: None, - }); - state_manager.update_data(data).await; + state_manager + .mutate_data(|data| { + let entry = data + .package_data + .entry(package_id.to_string()) + .or_insert_with(|| create_installing_entry(package_id)); + let existing_phase = entry.install_progress.as_ref().and_then(|p| p.phase); + entry.install_progress = Some(InstallProgress { + size: total, + downloaded, + phase: existing_phase, + message: None, + }); + }) + .await; } } /// Create a minimal PackageDataEntry for a package being installed. fn create_installing_entry(package_id: &str) -> PackageDataEntry { PackageDataEntry { + ui_ready: None, state: PackageState::Installing, health: None, exit_code: None, diff --git a/core/archipelago/src/api/rpc/package/runtime.rs b/core/archipelago/src/api/rpc/package/runtime.rs index de627004..ac65bc94 100644 --- a/core/archipelago/src/api/rpc/package/runtime.rs +++ b/core/archipelago/src/api/rpc/package/runtime.rs @@ -1431,10 +1431,9 @@ async fn repair_before_package_start(container_name: &str) { // published port and the data-dir file locks, so the replacement either // fails to bind (`address already in use`) or starts and dies on the // lock — and `Restart=always` loops it there forever. Ordered before - // the port cleanup below: killing the owner is what actually frees the - // port, and the port sweep alone cannot tell a ghost from a live app. + // starting the replacement. A port sweep cannot distinguish a ghost + // from the dashboard gate or another live app and must never kill it. crate::container::ghost_reaper::reap_for_app(container_name).await; - cleanup_runtime_host_ports(container_name).await; } async fn wait_before_package_start(container_name: &str) { @@ -1579,7 +1578,6 @@ async fn repair_netbird_network() { async fn repair_nginx_proxy_manager_container() { repair_nginx_proxy_manager_dirs().await; if !nginx_proxy_manager_has_legacy_admin_port().await { - cleanup_nginx_proxy_manager_ports().await; return; } @@ -1588,7 +1586,7 @@ async fn repair_nginx_proxy_manager_container() { ) .await; let _ = podman_control(&["rm", "-f", "nginx-proxy-manager"]).await; - cleanup_nginx_proxy_manager_ports().await; + crate::container::ghost_reaper::reap_for_app("nginx-proxy-manager").await; if let Err(err) = recreate_nginx_proxy_manager_container().await { tracing::warn!(error = %err, "failed to recreate stale nginx-proxy-manager container"); } @@ -1812,6 +1810,9 @@ fn manifest_host_ports(container_name: &str) -> Vec { pub(super) fn manifest_apps_dirs() -> Vec { let mut dirs = Vec::new(); + if let Some(root) = std::env::var_os("ARCHIPELAGO_APPS_DIR") { + dirs.push(root.into()); + } if let Ok(manifest_dir) = std::env::var("CARGO_MANIFEST_DIR") { dirs.push(Path::new(&manifest_dir).join("../../apps")); } @@ -2032,51 +2033,10 @@ async fn cleanup_start_conflict(container_name: &str, stderr: &str) { return; } - let ports = runtime_host_ports(container_name); - if !ports.is_empty() { - cleanup_ports(&ports).await; - return; - } -} - -async fn cleanup_runtime_host_ports(container_name: &str) { - let ports = runtime_host_ports(container_name); - if !ports.is_empty() { - cleanup_ports(&ports).await; - } -} - -async fn cleanup_nginx_proxy_manager_ports() { - cleanup_ports(&[8081, 8084, 8444]).await; -} - -async fn cleanup_ports(ports: &[u16]) { - for port in ports { - cleanup_stale_pasta_port(&port.to_string()).await; - } -} - -async fn cleanup_stale_pasta_port(port: &str) { - let kill_listener = format!( - "ss -ltnp 'sport = :{}' 2>/dev/null | sed -n 's/.*pid=\\([0-9]*\\).*/\\1/p' | xargs -r kill 2>/dev/null || true", - port - ); - let _ = tokio::process::Command::new("sh") - .args(["-c", &kill_listener]) - .output() - .await; - - let pattern = format!("pasta.*{}", port); - let _ = tokio::process::Command::new("pkill") - .args(["-f", &pattern]) - .output() - .await; - let pattern = format!("rootlessport.*{}", port); - let _ = tokio::process::Command::new("pkill") - .args(["-f", &pattern]) - .output() - .await; - tokio::time::sleep(std::time::Duration::from_secs(1)).await; + // Only reap processes proven to belong to an absent container. The app + // gate shares the app's port on other addresses and lives in this daemon; + // killing port owners (or matching argv with pkill) kills the dashboard. + crate::container::ghost_reaper::reap_for_app(container_name).await; } pub(super) fn is_missing_companion_ok(name: &str, stderr: &str) -> bool { @@ -2095,13 +2055,16 @@ async fn flip_package_state( package_id: &str, transitional: PackageState, ) -> Option { - let (mut data, _) = state_manager.get_snapshot().await; - let prev = data.package_data.get(package_id).map(|e| e.state.clone()); - if let Some(entry) = data.package_data.get_mut(package_id) { - entry.state = transitional; - state_manager.update_data(data).await; - } - prev + state_manager + .mutate_data(|data| { + let prev = data.package_data.get(package_id).map(|e| e.state.clone()); + if let Some(entry) = data.package_data.get_mut(package_id) { + entry.ui_ready = Some(false); + entry.state = transitional; + } + prev + }) + .await } /// Write the package entry's final state. No-op if the entry has since @@ -2111,13 +2074,18 @@ async fn set_package_state( package_id: &str, new_state: PackageState, ) { - let (mut data, _) = state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(package_id) { - if entry.state != new_state { - entry.state = new_state; - state_manager.update_data(data).await; - } - } + state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(package_id) { + if entry.state != new_state { + if new_state != PackageState::Running { + entry.ui_ready = Some(false); + } + entry.state = new_state; + } + } + }) + .await } pub(super) async fn reconcile_companions_for(package_id: &str) { @@ -2185,6 +2153,20 @@ pub(super) fn orchestrator_uninstall_app_ids(package_id: &str) -> Vec { mod tests { use super::*; + #[tokio::test] + async fn port_conflict_cleanup_preserves_live_host_listener() { + // The previous ss|kill sweep terminated the daemon's app gate on a + // restart. Keep a real listening socket owned by this test process. + let listener = tokio::net::TcpListener::bind("127.0.0.2:2342") + .await + .unwrap(); + let addr = listener.local_addr().unwrap(); + cleanup_start_conflict("photoprism", "address already in use").await; + let client = tokio::net::TcpStream::connect(addr).await.unwrap(); + let _connection = listener.accept().await.unwrap(); + drop(client); + } + #[test] fn missing_container_classifier_covers_podman5_phrasings() { // Regression (.228 gate 2026-07-08): podman 5.x `inspect` on a missing diff --git a/core/archipelago/src/api/rpc/transitional.rs b/core/archipelago/src/api/rpc/transitional.rs index 69992e69..b03052fc 100644 --- a/core/archipelago/src/api/rpc/transitional.rs +++ b/core/archipelago/src/api/rpc/transitional.rs @@ -150,23 +150,31 @@ async fn flip_to_transitional( app_id: &str, transitional: PackageState, ) -> Option { - let (mut data, _) = state_manager.get_snapshot().await; - let prev = data.package_data.get(app_id).map(|e| e.state.clone()); - if let Some(entry) = data.package_data.get_mut(app_id) { - entry.state = transitional; - state_manager.update_data(data).await; - } - prev + state_manager + .mutate_data(|data| { + let prev = data.package_data.get(app_id).map(|e| e.state.clone()); + if let Some(entry) = data.package_data.get_mut(app_id) { + entry.ui_ready = Some(false); + entry.state = transitional; + } + prev + }) + .await } /// Set the entry's state to `new_state`. No-ops if the entry has since been /// removed (e.g. uninstall ran concurrently). async fn set_state(state_manager: &StateManager, app_id: &str, new_state: PackageState) { - let (mut data, _) = state_manager.get_snapshot().await; - if let Some(entry) = data.package_data.get_mut(app_id) { - if entry.state != new_state { - entry.state = new_state; - state_manager.update_data(data).await; - } - } + state_manager + .mutate_data(|data| { + if let Some(entry) = data.package_data.get_mut(app_id) { + if entry.state != new_state { + if new_state != PackageState::Running { + entry.ui_ready = Some(false); + } + entry.state = new_state; + } + } + }) + .await } diff --git a/core/archipelago/src/appgate/identity.rs b/core/archipelago/src/appgate/identity.rs index 716b28d3..269ef7a2 100644 --- a/core/archipelago/src/appgate/identity.rs +++ b/core/archipelago/src/appgate/identity.rs @@ -114,6 +114,9 @@ impl PortMap { /// there. fn apps_dirs() -> Vec { let mut dirs = Vec::new(); + if let Some(root) = std::env::var_os("ARCHIPELAGO_APPS_DIR") { + dirs.push(root.into()); + } if let Ok(manifest_dir) = std::env::var("CARGO_MANIFEST_DIR") { dirs.push(PathBuf::from(manifest_dir).join("../../apps")); } diff --git a/core/archipelago/src/appgate/listener.rs b/core/archipelago/src/appgate/listener.rs index 99127aa3..6689e63f 100644 --- a/core/archipelago/src/appgate/listener.rs +++ b/core/archipelago/src/appgate/listener.rs @@ -144,6 +144,34 @@ pub fn shared_status() -> Arc> { .clone() } +static REFRESH_KICK: std::sync::LazyLock = + std::sync::LazyLock::new(tokio::sync::Notify::new); +static REFRESH_REV: std::sync::LazyLock> = + std::sync::LazyLock::new(|| tokio::sync::watch::channel(0).0); + +/// Installation must not wait for the minute sweep before becoming reachable. +/// Wait for a completed sweep, bounded if shutdown/startup prevents one. +pub async fn refresh_now() { + let mut completed = REFRESH_REV.subscribe(); + REFRESH_KICK.notify_one(); + let _ = tokio::time::timeout(std::time::Duration::from_secs(3), completed.changed()).await; +} + +pub fn port_claimed(status: &GateStatus, port: u16) -> bool { + let mut external = false; + let mut tor = false; + for (claimed_port, address) in &status.claimed { + if *claimed_port != port { + continue; + } + if let Ok(ip) = address.parse::() { + tor |= ip == GATE_TOR_UPSTREAM; + external |= !ip.is_loopback(); + } + } + external && tor +} + /// Run the gate. Returns only on shutdown. pub async fn run( gate: Arc, @@ -162,11 +190,12 @@ pub async fn run( loop { tokio::select! { - _ = interval.tick() => { - sweep(&gate, &status, &mut held, &shutdown_rx).await; - } + _ = interval.tick() => {} + _ = REFRESH_KICK.notified() => {} _ = shutdown_rx.changed() => return, } + sweep(&gate, &status, &mut held, &shutdown_rx).await; + REFRESH_REV.send_modify(|revision| *revision = revision.wrapping_add(1)); } } @@ -461,3 +490,19 @@ mod tests { assert!(!status.is_fully_enforced()); } } + +#[cfg(test)] +mod readiness_tests { + use super::*; + #[test] + fn readiness_requires_external_and_tor_claims_for_the_same_port() { + let mut status = GateStatus::default(); + assert!(!port_claimed(&status, 3001)); + status.claimed.push((3001, "127.0.0.2".into())); + assert!(!port_claimed(&status, 3001)); + status.claimed.push((3002, "192.0.2.10".into())); + assert!(!port_claimed(&status, 3001)); + status.claimed.push((3001, "192.0.2.10".into())); + assert!(port_claimed(&status, 3001)); + } +} diff --git a/core/archipelago/src/assistant/evals.rs b/core/archipelago/src/assistant/evals.rs index f37c14e8..83858827 100644 --- a/core/archipelago/src/assistant/evals.rs +++ b/core/archipelago/src/assistant/evals.rs @@ -322,6 +322,7 @@ async fn eval_rpc_handler() -> (Arc, tempfile::TempDir) { fn installed_entry(app_id: &str) -> crate::data_model::PackageDataEntry { use crate::data_model::{Description, Manifest, PackageDataEntry, PackageState, StaticFiles}; PackageDataEntry { + ui_ready: None, state: PackageState::Running, health: None, exit_code: None, diff --git a/core/archipelago/src/assistant/mod.rs b/core/archipelago/src/assistant/mod.rs index 4b1ba050..800e9325 100644 --- a/core/archipelago/src/assistant/mod.rs +++ b/core/archipelago/src/assistant/mod.rs @@ -1069,6 +1069,7 @@ mod tests { Description, Manifest, PackageDataEntry, PackageState, StaticFiles, }; PackageDataEntry { + ui_ready: None, state: PackageState::Running, health: None, exit_code: None, diff --git a/core/archipelago/src/container/docker_packages.rs b/core/archipelago/src/container/docker_packages.rs index 82095c63..1a7dbd04 100644 --- a/core/archipelago/src/container/docker_packages.rs +++ b/core/archipelago/src/container/docker_packages.rs @@ -3,8 +3,9 @@ use anyhow::Result; use archipelago_container::{ - ContainerRuntime as ContainerRuntimeTrait, ContainerState, PodmanClient, + ContainerRuntime as ContainerRuntimeTrait, ContainerState, ContainerStatus, PodmanClient, }; +use futures_util::StreamExt; use std::collections::HashMap; use std::sync::Arc; use tracing::{debug, info}; @@ -25,8 +26,15 @@ impl DockerPackageScanner { } /// Scan Docker containers and convert to package data - pub async fn scan_containers(&self) -> Result> { - let containers = self.runtime.list_containers().await?; + pub async fn scan_containers( + &self, + data_dir: &std::path::Path, + cached: &HashMap, + ) -> Result> { + let mut containers = self.runtime.list_containers().await?; + let installed = crate::crash_recovery::load_installed_apps(data_dir).await; + let uninstalled = crate::crash_recovery::load_user_uninstalled(data_dir).await; + restore_absent_installed(&mut containers, &installed, &uninstalled); debug!("Found {} containers", containers.len()); @@ -139,6 +147,18 @@ impl DockerPackageScanner { continue; } + if container.id.is_empty() { + if let Some(previous) = cached.get(&app_id) { + let mut held = previous.clone(); + held.state = PackageState::Stopped; + held.ui_ready = Some(false); + held.health = None; + held.exit_code = None; + packages.insert(app_id.clone(), held); + continue; + } + } + // Get metadata for this app let metadata = get_app_metadata(&app_id); // Manifest-owned metadata (icon) wins over the static table: the @@ -179,14 +199,22 @@ impl DockerPackageScanner { let tor_address = read_tor_address(&app_id).await; // Extract actual version from container image tag - let running_version = image_versions::extract_version_from_image(&container.image); + let running_version = if container.id.is_empty() { + String::new() // Absence cannot establish the installed image version. + } else { + image_versions::extract_version_from_image(&container.image) + }; // Decoupled from the binary OTA: prefer the remote app catalog, // falling back to the image-versions.sh pin when uncovered/offline. - let available_update = - crate::container::app_catalog::available_update_for_app(&app_id, &container.image); + let available_update = if container.id.is_empty() { + None + } else { + crate::container::app_catalog::available_update_for_app(&app_id, &container.image) + }; let package = PackageDataEntry { + ui_ready: Some(false), state: package_state.clone(), health: container.health.clone(), exit_code: if package_state == PackageState::Exited { @@ -283,10 +311,215 @@ impl DockerPackageScanner { ); } + let probes: Vec<_> = packages + .iter() + .filter_map(|(id, pkg)| { + if pkg.state != PackageState::Running { + return None; + } + let url = pkg + .installed + .as_ref()? + .interface_addresses + .get("main")? + .lan_address + .clone()?; + Some((id.clone(), url)) + }) + .collect(); + let mut results = futures_util::stream::iter( + probes + .into_iter() + .map(|(id, url)| async move { (id, launch_http_ready(&url).await) }), + ) + .buffer_unordered(8); + while let Some((id, ready)) = results.next().await { + if let Some(pkg) = packages.get_mut(&id) { + pkg.ui_ready = Some(ready); + } + } + // HTTP on loopback can precede the LAN/Tor listener after install. + let port_map = crate::appgate::identity::build_port_map(); + let gated: Vec<_> = packages + .iter() + .filter_map(|(id, pkg)| { + if pkg.ui_ready != Some(true) { + return None; + } + let url = pkg + .installed + .as_ref()? + .interface_addresses + .get("main")? + .lan_address + .as_deref()?; + let port = launch_url_port(url)?; + port_map + .gated(port) + .filter(|gate| gate.declared) + .map(|_| (id.clone(), port)) + }) + .collect(); + if !gated.is_empty() { + use crate::appgate::listener::{port_claimed, refresh_now, shared_status}; + let status = shared_status(); + let needs_refresh = { + let current = status.read().await; + gated.iter().any(|(_, port)| !port_claimed(¤t, *port)) + }; + if needs_refresh { + refresh_now().await; + } + let current = status.read().await; + for (id, port) in gated { + if !port_claimed(¤t, port) { + packages.get_mut(&id).unwrap().ui_ready = Some(false); + } + } + } Ok(packages) } } +/// Quadlet removes containers during ordinary stops/restarts. Rebuild installed +/// entries even on the daemon's first scan; a runtime absence is not uninstall. +fn restore_absent_installed( + containers: &mut Vec, + installed: &std::collections::HashSet, + uninstalled: &std::collections::HashSet, +) { + fn canonical(name: &str) -> &str { + let name = name.strip_prefix("archy-").unwrap_or(name); + match name { + "immich_server" => "immich", + _ => name, + } + } + let mut present: std::collections::HashSet = containers + .iter() + .map(|c| canonical(&c.name).to_owned()) + .collect(); + let removed: std::collections::HashSet<_> = + uninstalled.iter().map(|id| canonical(id)).collect(); + for name in installed { + let id = canonical(name); + if removed.contains(id) || !present.insert(id.to_owned()) { + continue; + } + containers.push(ContainerStatus { + id: String::new(), + name: id.to_owned(), + state: ContainerState::Stopped, + health: None, + exit_code: None, + started_at: None, + image: String::new(), + created: String::new(), + ports: Vec::new(), + lan_address: None, + }); + } +} + +/// Probe the actual loopback upstream, not the app gate's login page. A bound +/// TCP socket alone can still reset requests or serve a startup 503. +async fn launch_http_ready(candidate: &str) -> bool { + let Ok(mut url) = reqwest::Url::parse(candidate) else { + return false; + }; + if !matches!(url.scheme(), "http" | "https") { + return false; + } + if url.set_host(Some("127.0.0.1")).is_err() { + return false; + } + static CLIENT: std::sync::OnceLock = std::sync::OnceLock::new(); + let client = CLIENT.get_or_init(|| { + reqwest::Client::builder() + .no_proxy() + .timeout(std::time::Duration::from_secs(2)) + .redirect(reqwest::redirect::Policy::none()) + // Self-signed local app certificates are normal. This client only + // contacts loopback and never sends credentials or follows redirects. + .danger_accept_invalid_certs(true) + .build() + .expect("local readiness client") + }); + match client.get(url).send().await { + Ok(response) => matches!(response.status().as_u16(), 200..=399 | 401 | 403), + Err(_) => false, + } +} + +#[cfg(test)] +mod lifecycle_regression_tests { + use super::*; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + #[test] + fn registry_survives_empty_runtime_and_deduplicates_aliases() { + let installed = ["archy-gitea", "gitea", "immich_server", "archy-removed"] + .into_iter() + .map(str::to_owned) + .collect(); + let removed = ["removed".to_owned()].into_iter().collect(); + let mut containers = Vec::new(); + restore_absent_installed(&mut containers, &installed, &removed); + assert_eq!(containers.len(), 2); + assert!(containers + .iter() + .all(|c| c.state == ContainerState::Stopped)); + containers[0].state = ContainerState::Running; + restore_absent_installed(&mut containers, &installed, &removed); + assert_eq!(containers.len(), 2); + assert_eq!(containers[0].state, ContainerState::Running); + } + + #[tokio::test] + async fn readiness_rejects_startup_errors_and_accepts_auth_and_redirects() { + for (status, expected) in [ + (200, true), + (302, true), + (401, true), + (403, true), + (404, false), + (500, false), + (502, false), + (503, false), + ] { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let task = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.unwrap(); + let mut buf = [0; 2048]; + let n = stream.read(&mut buf).await.unwrap(); + assert!(String::from_utf8_lossy(&buf[..n]).starts_with("GET /start HTTP/1.1")); + stream.write_all(format!("HTTP/1.1 {status} Test\r\nContent-Length: 0\r\nConnection: close\r\n\r\n").as_bytes()).await.unwrap(); + }); + assert_eq!( + launch_http_ready(&format!("http://localhost:{port}/start")).await, + expected, + "status {status}" + ); + task.await.unwrap(); + } + } + + #[tokio::test] + async fn readiness_rejects_tcp_accept_without_http() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let task = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + drop(stream); + }); + assert!(!launch_http_ready(&format!("http://localhost:{port}/")).await); + task.await.unwrap(); + assert!(!launch_http_ready(&format!("http://localhost:{port}/")).await); + assert!(!launch_http_ready("file:///tmp/test").await); + } +} + struct AppMetadata { title: String, description: String, diff --git a/core/archipelago/src/container/prod_orchestrator.rs b/core/archipelago/src/container/prod_orchestrator.rs index 801193ef..d764d47d 100644 --- a/core/archipelago/src/container/prod_orchestrator.rs +++ b/core/archipelago/src/container/prod_orchestrator.rs @@ -3483,11 +3483,9 @@ impl ProdContainerOrchestrator { } async fn cleanup_stale_grafana_port(&self) { - let _ = tokio::process::Command::new("pkill") - .args(["-f", "pasta.*3001"]) - .output() - .await; - tokio::time::sleep(std::time::Duration::from_secs(1)).await; + // Port 3001 can belong to Gitea or the daemon's gate. Reap only a + // Grafana container proven absent from Podman's inventory. + crate::container::ghost_reaper::reap_for_app("grafana").await; } async fn detect_host_facts(&self) -> HostFacts { diff --git a/core/archipelago/src/crash_recovery.rs b/core/archipelago/src/crash_recovery.rs index 428d90dc..b4a66e3b 100644 --- a/core/archipelago/src/crash_recovery.rs +++ b/core/archipelago/src/crash_recovery.rs @@ -194,6 +194,7 @@ pub async fn clear_user_stopped(data_dir: &Path, name: &str) { // Installation is a decision, not a runtime observation, so it gets a record // of its own that no amount of downtime erodes. const INSTALLED_APPS_FILE: &str = "installed-apps.json"; +static INSTALLED_APPS_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); /// Load the durable set of installed app ids / container names. pub async fn load_installed_apps(data_dir: &Path) -> std::collections::HashSet { @@ -220,12 +221,23 @@ pub async fn load_installed_apps_if_recorded( async fn save_installed_apps(data_dir: &Path, installed: &std::collections::HashSet) { let path = data_dir.join(INSTALLED_APPS_FILE); if let Ok(json) = serde_json::to_string_pretty(installed) { - let _ = fs::write(&path, json).await; + let tmp = path.with_extension("json.tmp"); + let result = async { + fs::write(&tmp, json).await?; + fs::File::open(&tmp).await?.sync_all().await?; + fs::rename(&tmp, &path).await?; + fs::File::open(data_dir).await?.sync_all().await + } + .await; + if let Err(error) = result { + warn!(%error, "could not persist installed apps"); + } } } /// Record that an app is installed. Called when an install succeeds. pub async fn mark_installed(data_dir: &Path, name: &str) { + let _guard = INSTALLED_APPS_LOCK.lock().await; let mut installed = load_installed_apps(data_dir).await; if installed.insert(name.to_string()) { save_installed_apps(data_dir, &installed).await; @@ -235,6 +247,7 @@ pub async fn mark_installed(data_dir: &Path, name: &str) { /// Forget an app. Called on uninstall, beside `mark_user_uninstalled` — the /// two must move together or a reinstall-after-uninstall leaves a stale claim. pub async fn clear_installed(data_dir: &Path, name: &str) { + let _guard = INSTALLED_APPS_LOCK.lock().await; let mut installed = load_installed_apps(data_dir).await; if installed.remove(name) { save_installed_apps(data_dir, &installed).await; @@ -252,6 +265,7 @@ pub async fn clear_installed(data_dir: &Path, name: &str) { /// need it. Runs on every boot, so an app installed before the upgrade is /// still picked up whenever it is next seen alive. pub async fn backfill_installed_apps(data_dir: &Path, present_container_names: &[String]) { + let _guard = INSTALLED_APPS_LOCK.lock().await; if present_container_names.is_empty() { return; } @@ -1497,3 +1511,26 @@ mod tests { ); } } + +#[cfg(test)] +mod installed_concurrency_tests { + use super::*; + #[tokio::test] + async fn concurrent_install_records_are_not_lost() { + let dir = tempfile::tempdir().unwrap(); + let mut tasks = Vec::new(); + for i in 0..24 { + let path = dir.path().to_owned(); + tasks.push(tokio::spawn(async move { + mark_installed(&path, &format!("app-{i}")).await; + })); + } + for task in tasks { + task.await.unwrap(); + } + assert_eq!(load_installed_apps(dir.path()).await.len(), 24); + clear_installed(dir.path(), "app-3").await; + assert_eq!(load_installed_apps(dir.path()).await.len(), 23); + assert!(!dir.path().join("installed-apps.json.tmp").exists()); + } +} diff --git a/core/archipelago/src/data_model.rs b/core/archipelago/src/data_model.rs index 45f2191f..77540ddd 100644 --- a/core/archipelago/src/data_model.rs +++ b/core/archipelago/src/data_model.rs @@ -146,6 +146,10 @@ pub enum PackageState { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct PackageDataEntry { + /// Whether the app's HTTP upstream answered this scan (independent of + /// container health and blockchain sync). Missing on older nodes. + #[serde(rename = "ui-ready", default, skip_serializing_if = "Option::is_none")] + pub ui_ready: Option, pub state: PackageState, /// Container health: "healthy", "unhealthy", "starting", or null #[serde(skip_serializing_if = "Option::is_none")] @@ -297,6 +301,8 @@ pub enum InstallPhase { /// `podman pull` in progress (the longest phase — up to several /// minutes for large images on slow networks). PullingImage, + /// Orchestrator owns download/build and startup as one operation. + PreparingApp, /// Creating data directories, writing app-specific configs /// (bitcoin.conf, lnd.conf, searxng settings.yml, chown). CreatingContainer, diff --git a/core/archipelago/src/server.rs b/core/archipelago/src/server.rs index 43652abc..66d27b59 100644 --- a/core/archipelago/src/server.rs +++ b/core/archipelago/src/server.rs @@ -1765,12 +1765,17 @@ fn merge_preserving_transitional( }; crate::data_model::PackageDataEntry { - state, + state: state.clone(), // install_progress and uninstall_stage are also owned by the // initiating op (same reason as state) — keep them. install_progress: existing.install_progress.clone(), uninstall_stage: existing.uninstall_stage.clone(), // Everything else comes from the fresh scan. + ui_ready: if state == crate::data_model::PackageState::Running { + fresh.ui_ready + } else { + Some(false) + }, health: fresh.health.clone(), exit_code: fresh.exit_code, static_files: fresh.static_files.clone(), @@ -1809,7 +1814,10 @@ async fn scan_and_update_packages( absence_tracker: &mut HashMap, transitional_since: &mut HashMap, ) -> Result<()> { - let mut packages = scanner.scan_containers().await?; + let (before_scan, _) = state.get_snapshot().await; + let mut packages = scanner + .scan_containers(data_dir, &before_scan.package_data) + .await?; let user_stopped = crate::crash_recovery::load_user_stopped(data_dir).await; for (id, pkg) in packages.iter_mut() { if pkg.state == crate::data_model::PackageState::Exited && user_stopped.contains(id) { @@ -1870,11 +1878,14 @@ async fn scan_and_update_packages( // once at load ~2). Better to keep saying "scanning…" than to say "empty". if packages.is_empty() && (!first_scan || !installed_registry.is_empty()) { if tor_changed || update_changed { - let mut data = current_data; - data.server_info.tor_address = tor_addr.clone(); - data.server_info.node_address = tor_addr.as_ref().map(|t| identity.node_address(t)); - data.server_info.status_info.updated = update_available; - state.update_data(data).await; + state + .mutate_data(|data| { + data.server_info.tor_address = tor_addr.clone(); + data.server_info.node_address = + tor_addr.as_ref().map(|t| identity.node_address(t)); + data.server_info.status_info.updated = update_available; + }) + .await; } return Ok(()); } @@ -1899,6 +1910,13 @@ async fn scan_and_update_packages( // died without cleanup and let the scan override it. let now = Instant::now(); for (id, pkg) in &packages { + if user_uninstalled.contains(id) + || user_uninstalled.contains(&format!("archy-{id}")) + || (before_scan.package_data.contains_key(id) + && !current_data.package_data.contains_key(id)) + { + continue; + } absence_tracker.remove(id); let existing = merged.get(id); let overwrite = match existing { @@ -2054,22 +2072,40 @@ async fn scan_and_update_packages( } if changed || tor_changed || first_scan || update_changed { - let mut data = current_data; - data.package_data = merged; - data.server_info.tor_address = tor_addr.clone(); - data.server_info.node_address = tor_addr.as_ref().map(|t| identity.node_address(t)); - data.server_info.status_info.containers_scanned = true; - data.server_info.status_info.updated = update_available; - state.update_data(data).await; - debug!( - "📦 State changed (packages={}, tor={}, first_scan={}, update={}), broadcasting update", - changed, tor_changed, first_scan, update_changed - ); + state + .mutate_data(|data| { + // A lifecycle operation may have started/finished while this scan + // awaited probes or disk I/O. Never overwrite that newer entry or + // resurrect one that an uninstall removed in the meantime. + apply_scanned_packages(&mut data.package_data, ¤t_data.package_data, &merged); + data.server_info.tor_address = tor_addr.clone(); + data.server_info.node_address = tor_addr.as_ref().map(|t| identity.node_address(t)); + data.server_info.status_info.containers_scanned = true; + data.server_info.status_info.updated = update_available; + }) + .await; } Ok(()) } +fn apply_scanned_packages( + latest: &mut HashMap, + base: &HashMap, + scanned: &HashMap, +) { + for (id, fresh) in scanned { + if latest.get(id) == base.get(id) { + latest.insert(id.clone(), fresh.clone()); + } + } + for id in base.keys() { + if !scanned.contains_key(id) && latest.get(id) == base.get(id) { + latest.remove(id); + } + } +} + async fn normalize_reachable_package_health( packages: &mut HashMap, ) { @@ -2268,6 +2304,7 @@ mod merge_tests { fn make_entry(state: PackageState, health: Option<&str>) -> PackageDataEntry { PackageDataEntry { + ui_ready: None, state, health: health.map(|s| s.to_string()), exit_code: None, @@ -2280,6 +2317,37 @@ mod merge_tests { } } + #[test] + fn stale_scan_cannot_remove_new_installs_or_overwrite_lifecycle_changes() { + let running = make_entry(PackageState::Running, Some("healthy")); + let restarting = make_entry(PackageState::Restarting, None); + let base = [ + ("restart".into(), running.clone()), + ("uninstalled".into(), running.clone()), + ] + .into_iter() + .collect(); + let mut latest = [ + ("restart".into(), restarting.clone()), + ("new".into(), running.clone()), + ] + .into_iter() + .collect(); + let scanned = [ + ("restart".into(), running.clone()), + ("uninstalled".into(), running.clone()), + ] + .into_iter() + .collect(); + apply_scanned_packages(&mut latest, &base, &scanned); + assert_eq!(latest.get("restart"), Some(&restarting)); + assert_eq!(latest.get("new"), Some(&running)); + assert!(!latest.contains_key("uninstalled")); + apply_scanned_packages(&mut latest, &base, &HashMap::new()); + assert_eq!(latest.get("restart"), Some(&restarting)); + assert!(latest.contains_key("new")); + } + #[test] fn peer_path_filter_allows_content_catalog_and_items() { // Regression: the content *catalog* is exactly "/content" (no trailing diff --git a/core/archipelago/src/state.rs b/core/archipelago/src/state.rs index d2b319d4..b835290c 100644 --- a/core/archipelago/src/state.rs +++ b/core/archipelago/src/state.rs @@ -54,6 +54,21 @@ impl StateManager { let _ = self.broadcast_tx.send(message); } + /// Apply a small state change while holding the write lock. A lifecycle + /// task must not replace the entire model from an earlier snapshot. + pub async fn mutate_data(&self, change: impl FnOnce(&mut DataModel) -> T) -> T { + let mut data = self.data.write().await; + let result = change(&mut data); + let mut rev = self.revision.write().await; + *rev += 1; + let _ = self.broadcast_tx.send(WebSocketMessage { + rev: *rev, + data: Some(data.clone()), + patch: None, + }); + result + } + /// Get a WebSocket message with the current state pub async fn get_initial_message(&self) -> WebSocketMessage { let (data, rev) = self.get_snapshot().await; @@ -190,3 +205,29 @@ mod tests { assert_eq!(rev, 1); } } + +#[cfg(test)] +mod atomic_mutation_tests { + use super::*; + #[tokio::test] + async fn concurrent_updates_preserve_independent_entries() { + let state = Arc::new(StateManager::new()); + let mut tasks = Vec::new(); + for i in 0..24 { + let state = state.clone(); + tasks.push(tokio::spawn(async move { + state + .mutate_data(|data| { + data.peer_health.insert(format!("peer-{i}"), true); + }) + .await; + })); + } + for task in tasks { + task.await.unwrap(); + } + let (data, revision) = state.get_snapshot().await; + assert_eq!(data.peer_health.len(), 24); + assert_eq!(revision, 24); + } +} diff --git a/neode-ui/src/stores/__tests__/appLauncher.test.ts b/neode-ui/src/stores/__tests__/appLauncher.test.ts index db871c13..68bf8ee7 100644 --- a/neode-ui/src/stores/__tests__/appLauncher.test.ts +++ b/neode-ui/src/stores/__tests__/appLauncher.test.ts @@ -34,6 +34,7 @@ vi.mock('@/api/rpc-client', () => ({ vi.stubGlobal('open', mockWindowOpen) import { useAppLauncherStore, senderMatchesApp } from '../appLauncher' +import { useAppStore } from '../app' describe('useAppLauncherStore', () => { beforeEach(() => { @@ -54,6 +55,25 @@ describe('useAppLauncherStore', () => { }) }) + it('blocks both browser and embedded launch while HTTP is unready', () => { + const app = useAppStore() + app.data = { 'package-data': { gitea: { state: 'running', 'ui-ready': false, health: 'healthy', manifest: { id: 'gitea', title: 'Gitea' } } } } as never + const launcher = useAppLauncherStore() + launcher.openSession('gitea') + expect(launcher.panelAppId).toBeNull() + launcher.open({ url: 'http://192.0.2.10:3001/', title: 'Gitea', openInNewTab: true }) + expect(mockWindowOpen).not.toHaveBeenCalled() + expect(launcher.isOpen).toBe(false) + }) + + it('also gates a dynamic app resolved through its runtime URL', () => { + useAppStore().data = { 'package-data': { custom: { state: 'running', 'ui-ready': false, manifest: { id: 'custom', title: 'Custom' }, installed: { 'interface-addresses': { main: { 'lan-address': 'http://localhost:18993/' } } } } } } as never + const launcher = useAppLauncherStore() + launcher.open({ url: 'http://192.0.2.10:18993/', title: 'Custom', openInNewTab: true }) + expect(mockWindowOpen).not.toHaveBeenCalled() + expect(launcher.isOpen).toBe(false) + }) + it('starts closed with empty state', () => { const store = useAppLauncherStore() expect(store.isOpen).toBe(false) diff --git a/neode-ui/src/stores/__tests__/server.test.ts b/neode-ui/src/stores/__tests__/server.test.ts new file mode 100644 index 00000000..96ce5b0d --- /dev/null +++ b/neode-ui/src/stores/__tests__/server.test.ts @@ -0,0 +1,42 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { createPinia, setActivePinia } from 'pinia' +import { reactive, nextTick } from 'vue' + +const fake = reactive<{ packages: Record }>({ packages: {} }) +vi.mock('../sync', () => ({ useSyncStore: () => fake })) +vi.mock('../../api/rpc-client', () => ({ rpcClient: {} })) +import { useServerStore } from '../server' + +function installing(phase = 'preparing-app') { + return { state: 'installing', manifest: { title: 'Git Workshop' }, 'install-progress': { phase, size: 0, downloaded: 0 } } +} + +describe('installation state after hard refresh', () => { + beforeEach(() => { + setActivePinia(createPinia()) + fake.packages = {} + }) + + it('restores an in-flight install from an already-loaded server snapshot', () => { + fake.packages = { 'archipelago-source': installing() } + const store = useServerStore() + expect(store.isInstalling('archipelago-source')).toBe(true) + expect(store.installingApps.get('archipelago-source')).toMatchObject({ + progress: 20, + message: 'Downloading, building and starting app…', + }) + }) + + it('keeps a long download visible and clears it on terminal success', async () => { + const store = useServerStore() + fake.packages = { 'nginx-proxy-manager': installing() } + await nextTick() + expect(store.isInstalling('nginx-proxy-manager')).toBe(true) + fake.packages = { 'nginx-proxy-manager': installing() } + await nextTick() + expect(store.installingApps.get('nginx-proxy-manager')?.progress).toBe(20) + fake.packages = { 'nginx-proxy-manager': { state: 'running' } } + await nextTick() + expect(store.isInstalling('nginx-proxy-manager')).toBe(false) + }) +}) diff --git a/neode-ui/src/stores/appLauncher.ts b/neode-ui/src/stores/appLauncher.ts index 50e5818c..aca7a029 100644 --- a/neode-ui/src/stores/appLauncher.ts +++ b/neode-ui/src/stores/appLauncher.ts @@ -239,6 +239,11 @@ export const useAppLauncherStore = defineStore('appLauncher', () => { const panelPath = ref(null) function openSessionNow(appId: string, opts: LaunchOptions = {}) { + const pkg = useAppStore().data?.['package-data']?.[appId] + if (pkg?.['ui-ready'] === false) { + useToast().info(`${pkg.manifest?.title || appId} is not ready to open yet`) + return + } recordAppLaunch(appId) const mobile = isMobileViewport() @@ -295,7 +300,7 @@ export const useAppLauncherStore = defineStore('appLauncher', () => { // Apply the same readiness gate here so a container that has just entered // `running` cannot race nginx and show a transient 502 to the user. const pkg = useAppStore().data?.['package-data']?.[appId] - if (pkg && pkg.state === 'running' && !isAppReadyForLaunch(pkg)) { + if (pkg && (pkg['ui-ready'] === false || (pkg.state === 'running' && !isAppReadyForLaunch(pkg)))) { useToast().info(`${pkg.manifest?.title || appId} is still starting — try again in a moment`) return } @@ -394,6 +399,12 @@ export const useAppLauncherStore = defineStore('appLauncher', () => { let launchUrl = normalizeLaunchUrl(payload.url, titleHintId) const resolvedId = resolveAppIdFromUrl(launchUrl) || titleHintId + const pkg = resolvedId ? useAppStore().data?.['package-data']?.[resolvedId] : undefined + if (pkg?.['ui-ready'] === false) { + useToast().info(`${pkg.manifest?.title || resolvedId} is not ready to open yet`) + return + } + // Scheme discipline for everything launched on this host. Ports fronted // by the node's app gate (manifest auth gated/open) serve TLS on the same // port — on an HTTPS connection those must open over https. Ports that @@ -472,6 +483,17 @@ export const useAppLauncherStore = defineStore('appLauncher', () => { // Check /app/{id}/ path-style routes first (HTTPS proxy mode) const m = u.pathname.match(/^\/app\/([a-z0-9._-]+)(?:\/|$)/i) if (m?.[1]) return m[1].toLowerCase() + // Dynamic/sideloaded apps have no entry in the static port map. + if (u.hostname === window.location.hostname && u.port) { + for (const [id, pkg] of Object.entries(useAppStore().data?.['package-data'] || {})) { + const address = pkg.installed?.['interface-addresses']?.main?.['lan-address'] + if (!address) continue + try { + const runtime = new URL(address) + if (runtime.port === u.port && ['localhost', '127.0.0.1', window.location.hostname].includes(runtime.hostname)) return id + } catch { /* malformed runtime address is not a launch target */ } + } + } // Check port-based apps const appId = PORT_TO_APP_ID[u.port] if (appId) return appId diff --git a/neode-ui/src/stores/server.ts b/neode-ui/src/stores/server.ts index 9b88d72f..fd213d12 100644 --- a/neode-ui/src/stores/server.ts +++ b/neode-ui/src/stores/server.ts @@ -25,6 +25,7 @@ import type { InstallPhase } from '../types/api' const PHASE_INFO: Record = { 'preparing': { progress: 5, message: 'Preparing…', status: 'downloading' }, 'pulling-image': { progress: 20, message: 'Downloading image…', status: 'downloading' }, + 'preparing-app': { progress: 20, message: 'Downloading, building and starting app…', status: 'downloading' }, 'creating-container': { progress: 70, message: 'Creating container…', status: 'installing' }, 'starting-container': { progress: 80, message: 'Starting container…', status: 'starting' }, 'waiting-healthy': { progress: 88, message: 'Finalizing first start…', status: 'starting' }, @@ -149,7 +150,7 @@ export const useServerStore = defineStore('server', () => { uninstallingApps.value.delete(appId) } } - }, { deep: true }) + }, { deep: true, immediate: true }) function setInstallProgress(appId: string, progress: Partial & { id: string; title: string }) { const existing = installingApps.value.get(appId) diff --git a/neode-ui/src/types/api.ts b/neode-ui/src/types/api.ts index 57c0ff3f..ca1f7f32 100644 --- a/neode-ui/src/types/api.ts +++ b/neode-ui/src/types/api.ts @@ -88,6 +88,7 @@ export const PackageState = { export type PackageState = typeof PackageState[keyof typeof PackageState] export interface PackageDataEntry { + 'ui-ready'?: boolean // HTTP upstream readiness, separate from container health state: PackageState health?: string | null // "healthy", "unhealthy", "starting", or null 'exit-code'?: number | null // container exit code: 0 = clean stop, non-zero = crash @@ -180,6 +181,7 @@ export type ServiceStatus = typeof ServiceStatus[keyof typeof ServiceStatus] export type InstallPhase = | 'preparing' | 'pulling-image' + | 'preparing-app' | 'creating-container' | 'starting-container' | 'waiting-healthy' diff --git a/neode-ui/src/views/AppDetails.vue b/neode-ui/src/views/AppDetails.vue index 367d9853..2d30a802 100644 --- a/neode-ui/src/views/AppDetails.vue +++ b/neode-ui/src/views/AppDetails.vue @@ -259,7 +259,7 @@ const canLaunch = computed(() => { const hasRuntimeAddress = !!pkg.value.installed?.['interface-addresses']?.main?.['lan-address'] const hasKnownLaunchUrl = typeof window !== 'undefined' && !!resolveAppUrl(pkg.value.manifest.id) const hasUI = !!(pkg.value.manifest.interfaces?.main?.ui || hasRuntimeAddress || hasKnownLaunchUrl) - return hasUI && pkg.value.state === 'running' && pkg.value.health !== 'starting' && pkg.value.health !== 'unhealthy' + return hasUI && pkg.value['ui-ready'] !== false && pkg.value.state === 'running' && pkg.value.health !== 'starting' && pkg.value.health !== 'unhealthy' }) const features = computed(() => { diff --git a/neode-ui/src/views/AppSession.vue b/neode-ui/src/views/AppSession.vue index d752e4c7..4e9fc1d6 100644 --- a/neode-ui/src/views/AppSession.vue +++ b/neode-ui/src/views/AppSession.vue @@ -39,6 +39,7 @@ :must-open-new-tab="mustOpenNewTab" :auto-retry-count="autoRetryCount" :refresh-key="refreshKey" + :ui-ready-blocked="packageEntry?.['ui-ready'] === false" :blocked-reason="blockedReason" :blocked-title="blockedTitle" :warming-up="warmingUp" @@ -375,6 +376,20 @@ const panelClasses = computed(() => { return `${base} app-session-overlay` }) +// A cold/restarting upstream is held outside the iframe. Start one fresh +// load when the scanner observes HTTP readiness; no manual refresh required. +watch(() => packageEntry.value?.['ui-ready'], (ready, previous) => { + if (ready === false) { + if (loadTimeoutId) clearTimeout(loadTimeoutId) + if (autoRetryId) clearTimeout(autoRetryId) + if (iframeCheckId) clearTimeout(iframeCheckId) + loading.value = false + } else if (previous === false && ready === true) { + autoRetryCount.value = 0 + refresh() + } +}) + // --- Lifecycle handlers --- function onLoad() { @@ -433,6 +448,7 @@ function refresh() { function startLoadTimeout() { if (loadTimeoutId) clearTimeout(loadTimeoutId) + if (packageEntry.value?.['ui-ready'] === false) return loadTimeoutId = setTimeout(() => { if (loading.value) { loading.value = false @@ -442,11 +458,13 @@ function startLoadTimeout() { } function openNewTabAndBack() { + if (packageEntry.value?.['ui-ready'] === false) return if (appUrl.value) openExternalUrl(appUrl.value) closeSession() } function openNewTab() { + if (packageEntry.value?.['ui-ready'] === false) return if (appUrl.value) openExternalUrl(appUrl.value) } diff --git a/neode-ui/src/views/appSession/AppSessionFrame.vue b/neode-ui/src/views/appSession/AppSessionFrame.vue index 4cb98dcf..5df52fe7 100644 --- a/neode-ui/src/views/appSession/AppSessionFrame.vue +++ b/neode-ui/src/views/appSession/AppSessionFrame.vue @@ -6,7 +6,7 @@ first, then sync status arrives), and the sync screen is strictly more informative, so it takes precedence instead of the two rendering on top of each other. --> - + -
+
@@ -78,7 +78,8 @@

{{ warmingUp ? `${appTitle} is starting…` : blockedReason ? blockedTitle : (mustOpenNewTab ? 'This app opens in a new tab' : 'App not reachable') }}

- + + @@ -95,6 +96,7 @@ Retry now

-
+

App not configured

No URL found for {{ appId }}

@@ -133,6 +135,7 @@ const props = defineProps<{ mustOpenNewTab: boolean autoRetryCount: number refreshKey: number + uiReadyBlocked?: boolean blockedReason?: string blockedTitle?: string // True while the container is up but its probe hasn't answered yet and the diff --git a/neode-ui/src/views/appSession/__tests__/AppSessionFrame.test.ts b/neode-ui/src/views/appSession/__tests__/AppSessionFrame.test.ts index b616198a..a9fc14da 100644 --- a/neode-ui/src/views/appSession/__tests__/AppSessionFrame.test.ts +++ b/neode-ui/src/views/appSession/__tests__/AppSessionFrame.test.ts @@ -66,3 +66,19 @@ describe('AppSessionFrame warm-up state', () => { expect(text).toContain('This app opens in a new tab') }) }) + +describe('HTTP readiness gate', () => { + it('does not show a missing-configuration error during initial installation', () => { + const frame = mountFrame({ appUrl: '', uiReadyBlocked: true, blockedReason: 'Waiting for the app to be ready…' }) + expect(frame.text()).not.toContain('App not configured') + expect(frame.find('iframe').exists()).toBe(false) + }) + it('does not mount an iframe before readiness, then opens automatically', async () => { + const frame = mountFrame({ iframeBlocked: false, uiReadyBlocked: true, blockedReason: 'Waiting for the app to be ready…', blockedTitle: 'App not ready' }) + expect(frame.find('iframe').exists()).toBe(false) + expect(frame.text()).toContain('opens automatically') + expect(frame.text()).not.toContain('Open in new tab') + await frame.setProps({ uiReadyBlocked: false, blockedReason: '' }) + expect(frame.find('iframe').exists()).toBe(true) + }) +}) diff --git a/neode-ui/src/views/apps/__tests__/appsConfig.test.ts b/neode-ui/src/views/apps/__tests__/appsConfig.test.ts index c363b68d..9766736b 100644 --- a/neode-ui/src/views/apps/__tests__/appsConfig.test.ts +++ b/neode-ui/src/views/apps/__tests__/appsConfig.test.ts @@ -184,3 +184,24 @@ describe('appsConfig service filtering', () => { expect(canLaunch(pkg)).toBe(true) }) }) + +describe('HTTP readiness independent of container health', () => { + it('blocks fixed launch URLs while the HTTP upstream is unavailable', () => { + for (const id of ['gitea', 'filebrowser', 'fedimint', 'lnd']) { + const pkg = makePkg(id, id, 'other') + pkg['ui-ready'] = false + pkg.health = 'healthy' + expect(canLaunch(pkg)).toBe(false) + expect(isAppReadyForLaunch(pkg)).toBe(false) + expect(launchBlockedReason(id, pkg)).toContain('Waiting') + pkg['ui-ready'] = true + expect(isAppReadyForLaunch(pkg)).toBe(true) + } + }) + it('allows a ready companion while its backend is syncing', () => { + const pkg = makePkg('lnd', 'Lightning', 'bitcoin') + pkg.health = 'starting' + pkg['ui-ready'] = true + expect(isAppReadyForLaunch(pkg)).toBe(true) + }) +}) diff --git a/neode-ui/src/views/apps/appsConfig.ts b/neode-ui/src/views/apps/appsConfig.ts index 1d2952b0..95a1b68e 100644 --- a/neode-ui/src/views/apps/appsConfig.ts +++ b/neode-ui/src/views/apps/appsConfig.ts @@ -244,6 +244,7 @@ export function resolveAppIcon(id: string, pkg: PackageDataEntry, curatedIcon?: export function canLaunch(pkg: PackageDataEntry): boolean { if (isWebOnlyApp(pkg.manifest.id)) return true + if (pkg['ui-ready'] === false) return false // Headless backends never get a Launch button, even with a published port. if (isServicePackage(pkg.manifest.id, pkg)) return false const hasRuntimeAddress = !!pkg.installed?.['interface-addresses']?.main?.['lan-address'] @@ -277,6 +278,7 @@ export function canLaunch(pkg: PackageDataEntry): boolean { * health check retain the legacy state/port behaviour. */ export function isAppReadyForLaunch(pkg: PackageDataEntry): boolean { + if (pkg['ui-ready'] !== undefined) return pkg['ui-ready'] const manifest = pkg.manifest as unknown as Record const hasHealthCheck = Boolean(manifest.health_check || manifest['health-check']) if (!hasHealthCheck) return pkg.health !== 'unhealthy' @@ -285,6 +287,10 @@ export function isAppReadyForLaunch(pkg: PackageDataEntry): boolean { export function launchBlockedReason(id: string, pkg?: PackageDataEntry | null): string { const appId = pkg?.manifest?.id || id + if (pkg?.['ui-ready'] === false && !isServicePackage(appId, pkg)) { + if (pkg.state === PackageState.Stopped || pkg.state === PackageState.Exited) return 'App is stopped. Start it to open it.' + return 'Waiting for the app to be ready…' + } if ( (appId === 'fedimint' || appId === 'fedimintd') && (pkg?.state === PackageState.Starting || (pkg?.state === PackageState.Running && pkg?.health === 'starting'))