diff --git a/clients/openframe-client/src/lib.rs b/clients/openframe-client/src/lib.rs index f815586c40..e5162ece98 100644 --- a/clients/openframe-client/src/lib.rs +++ b/clients/openframe-client/src/lib.rs @@ -761,8 +761,9 @@ impl Client { self.script_schedule_execution_listener.start().await?; info!("Script schedule execution listener started"); - // Start tool run manager - self.tool_run_manager.run().await?; + if let Err(e) = self.tool_run_manager.run().await { + error!("Failed to start tool run manager: {:#}", e); + } // Start mesh self-heal watcher (re-fetch .msh + bounce agent if held on a stale MeshID). self.mesh_self_heal_service.run().await?; diff --git a/clients/openframe-client/src/platform/preferences_writer.rs b/clients/openframe-client/src/platform/preferences_writer.rs index e35d221274..f9a3126783 100644 --- a/clients/openframe-client/src/platform/preferences_writer.rs +++ b/clients/openframe-client/src/platform/preferences_writer.rs @@ -19,8 +19,11 @@ pub fn write<'a>( } for (key, value) in &prefs { - let status = Command::new("sudo") + let status = Command::new("launchctl") .args([ + "asuser", + &user.uid.to_string(), + "sudo", "-u", &user.username, "defaults", @@ -32,8 +35,25 @@ pub fn write<'a>( .status() .with_context(|| format!("Failed to write preference '{}'", key))?; - if !status.success() { - anyhow::bail!("defaults write failed for '{}': exit {}", key, status); + if status.success() { + continue; + } + + let fallback = Command::new("sudo") + .args([ + "-u", + &user.username, + "defaults", + "write", + bundle_id, + key, + value, + ]) + .status() + .with_context(|| format!("Failed to write preference '{}'", key))?; + + if !fallback.success() { + anyhow::bail!("defaults write failed for '{}': {}", key, fallback); } } diff --git a/clients/openframe-client/src/platform/tool_updater/gui_app.rs b/clients/openframe-client/src/platform/tool_updater/gui_app.rs index 1d44999ee4..655c309ca4 100644 --- a/clients/openframe-client/src/platform/tool_updater/gui_app.rs +++ b/clients/openframe-client/src/platform/tool_updater/gui_app.rs @@ -1,13 +1,13 @@ use anyhow::{Context, Result}; use async_trait::async_trait; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use tracing::{error, info, warn}; use super::{ToolUpdater, ToolUpdaterDeps, UpdateContext}; use crate::models::{DownloadConfiguration, Installation, InstalledTool}; use crate::platform::preferences_writer::{args_to_pairs, write as write_preferences}; -use crate::platform::remove_app_bundle; use crate::platform::user_session::{get_console_user, launch_as_user}; +use crate::platform::DirectoryManager; pub struct GuiAppToolUpdater { deps: ToolUpdaterDeps, @@ -17,6 +17,46 @@ impl GuiAppToolUpdater { pub fn new(deps: ToolUpdaterDeps) -> Self { Self { deps } } + + fn backup_path_for(bundle: &Path) -> PathBuf { + let name = bundle + .file_name() + .map(|n| n.to_string_lossy().to_string()) + .unwrap_or_else(|| "bundle".to_string()); + bundle.with_file_name(format!(".{name}.update-backup")) + } + + async fn move_bundle_aside( + executable_path: &str, + tool_agent_id: &str, + ) -> Result> { + let Some(bundle) = DirectoryManager::find_app_bundle_path(Path::new(executable_path)) + else { + warn!(tool_id = %tool_agent_id, "No .app bundle in {} — updating without a backup", executable_path); + return Ok(None); + }; + let backup = Self::backup_path_for(&bundle); + + if !bundle.exists() { + if backup.exists() { + info!(tool_id = %tool_agent_id, "Adopting the backup left by an interrupted update"); + return Ok(Some(backup)); + } + return Ok(None); + } + if backup.exists() { + tokio::fs::remove_dir_all(&backup).await.ok(); + } + tokio::fs::rename(&bundle, &backup).await.with_context(|| { + format!( + "Failed to move {} aside to {}", + bundle.display(), + backup.display() + ) + })?; + info!(tool_id = %tool_agent_id, "Old app bundle kept at {}", backup.display()); + Ok(Some(backup)) + } } #[async_trait] @@ -34,8 +74,15 @@ impl ToolUpdater for GuiAppToolUpdater { tokio::time::sleep(tokio::time::Duration::from_secs(2)).await; + let backup_path = match &tool.installation { + Installation::GuiApp { + executable_path, .. + } => Self::move_bundle_aside(executable_path, tool_agent_id).await?, + _ => None, + }; + Ok(UpdateContext { - backup_path: None, + backup_path, needs_restart: true, }) } @@ -49,20 +96,13 @@ impl ToolUpdater for GuiAppToolUpdater { let tool_agent_id = &tool.tool_agent_id; info!(tool_id = %tool_agent_id, "Applying GuiApp update"); - let Installation::GuiApp { - executable_path, - bundle_id, - } = &tool.installation - else { + let Installation::GuiApp { bundle_id, .. } = &tool.installation else { anyhow::bail!( "Expected GuiApp installation type for tool: {}", tool_agent_id ); }; - info!(tool_id = %tool_agent_id, "Removing old app bundle"); - remove_app_bundle(executable_path).await?; - let applications_dir = PathBuf::from("/Applications"); info!(tool_id = %tool_agent_id, "Downloading and installing new version from: {}", config.link); @@ -93,6 +133,14 @@ impl ToolUpdater for GuiAppToolUpdater { let tool_agent_id = &tool.tool_agent_id; info!(tool_id = %tool_agent_id, "Finalizing GuiApp update"); + if let Some(backup) = &ctx.backup_path { + if backup.exists() { + if let Err(e) = tokio::fs::remove_dir_all(backup).await { + warn!(tool_id = %tool_agent_id, "Failed to remove the old app bundle {}: {:#}", backup.display(), e); + } + } + } + if !ctx.needs_restart { info!(tool_id = %tool_agent_id, "Restart not requested, skipping"); return Ok(()); @@ -154,10 +202,37 @@ impl ToolUpdater for GuiAppToolUpdater { Ok(()) } - async fn rollback(&self, tool: &InstalledTool, _ctx: &UpdateContext) -> Result<()> { + async fn rollback(&self, tool: &InstalledTool, ctx: &UpdateContext) -> Result<()> { let tool_agent_id = &tool.tool_agent_id; - warn!(tool_id = %tool_agent_id, - "Rollback requested for GuiApp but no backup available. User should reinstall from server."); + + let Some(backup) = &ctx.backup_path else { + warn!(tool_id = %tool_agent_id, + "Rollback requested for GuiApp but no backup was taken. User should reinstall from server."); + return Ok(()); + }; + + let Installation::GuiApp { + executable_path, .. + } = &tool.installation + else { + anyhow::bail!("Expected GuiApp installation type for tool: {tool_agent_id}"); + }; + let Some(bundle) = DirectoryManager::find_app_bundle_path(Path::new(executable_path)) + else { + anyhow::bail!("Could not resolve the .app bundle path for {tool_agent_id}"); + }; + + if bundle.exists() { + tokio::fs::remove_dir_all(&bundle).await.ok(); + } + tokio::fs::rename(backup, &bundle).await.with_context(|| { + format!( + "Failed to restore the previous app bundle from {}", + backup.display() + ) + })?; + + info!(tool_id = %tool_agent_id, "Restored the previous app bundle: {}", bundle.display()); Ok(()) } } diff --git a/clients/openframe-client/src/platform/tool_updater/standard.rs b/clients/openframe-client/src/platform/tool_updater/standard.rs index ae48522047..8b49ebccd8 100644 --- a/clients/openframe-client/src/platform/tool_updater/standard.rs +++ b/clients/openframe-client/src/platform/tool_updater/standard.rs @@ -31,7 +31,10 @@ impl ToolUpdater for StandardToolUpdater { .await .with_context(|| format!("Failed to stop tool: {}", tool_agent_id))?; - let agent_path = self.deps.directory_manager.get_agent_path(tool_agent_id); + let agent_path = self + .deps + .directory_manager + .get_tool_executable_path(&tool.tool_agent_id, tool.installation.executable_path()); clear_aside_binary(&agent_path, tool_agent_id).await; log_update_survivors(&self.deps, tool).await; @@ -52,7 +55,10 @@ impl ToolUpdater for StandardToolUpdater { let tool_agent_id = &tool.tool_agent_id; info!(tool_id = %tool_agent_id, "Applying Standard tool update"); - let agent_path = self.deps.directory_manager.get_agent_path(tool_agent_id); + let agent_path = self + .deps + .directory_manager + .get_tool_executable_path(&tool.tool_agent_id, tool.installation.executable_path()); download_and_write_binary(&self.deps, config, &agent_path, tool_agent_id).await?; Ok(None) } @@ -74,7 +80,10 @@ impl ToolUpdater for StandardToolUpdater { let tool_agent_id = &tool.tool_agent_id; info!(tool_id = %tool_agent_id, "Rolling back Standard tool update"); - let agent_path = self.deps.directory_manager.get_agent_path(tool_agent_id); + let agent_path = self + .deps + .directory_manager + .get_tool_executable_path(&tool.tool_agent_id, tool.installation.executable_path()); restore_from_backup(ctx.backup_path.as_ref(), &agent_path, tool_agent_id).await } } diff --git a/clients/openframe-client/src/service.rs b/clients/openframe-client/src/service.rs index cfd5f8b93c..e28987d772 100644 --- a/clients/openframe-client/src/service.rs +++ b/clients/openframe-client/src/service.rs @@ -83,14 +83,22 @@ fn windows_service_main(_args: Vec) { }; // Report that the service is running - let _ = set_service_status(&status_handle, ServiceState::Running); + let _ = set_service_status( + &status_handle, + ServiceState::Running, + ServiceExitCode::Win32(0), + ); // Create a Tokio runtime and run the service core let rt = match Runtime::new() { Ok(runtime) => runtime, Err(e) => { eprintln!("Failed to create Tokio runtime: {:?}", e); - let _ = set_service_status(&status_handle, ServiceState::Stopped); + let _ = set_service_status( + &status_handle, + ServiceState::Stopped, + ServiceExitCode::ServiceSpecific(1), + ); return; } }; @@ -114,17 +122,29 @@ fn windows_service_main(_args: Vec) { }); if let Err(e) = result { - eprintln!("Service core failed: {:?}", e); - let _ = set_service_status(&status_handle, ServiceState::Stopped); + error!("Service core failed: {:#}", e); + let _ = set_service_status( + &status_handle, + ServiceState::Stopped, + ServiceExitCode::ServiceSpecific(1), + ); } else { info!("Service stopped gracefully"); - let _ = set_service_status(&status_handle, ServiceState::Stopped); + let _ = set_service_status( + &status_handle, + ServiceState::Stopped, + ServiceExitCode::Win32(0), + ); } } /// Helper function to set service status #[cfg(windows)] -fn set_service_status(status_handle: &ServiceStatusHandle, state: ServiceState) -> Result<()> { +fn set_service_status( + status_handle: &ServiceStatusHandle, + state: ServiceState, + exit_code: ServiceExitCode, +) -> Result<()> { let status = ServiceStatus { service_type: ServiceType::OWN_PROCESS, current_state: state, @@ -133,7 +153,7 @@ fn set_service_status(status_handle: &ServiceStatusHandle, state: ServiceState) } else { ServiceControlAccept::empty() }, - exit_code: ServiceExitCode::Win32(0), + exit_code, checkpoint: 0, wait_hint: std::time::Duration::from_secs(5), process_id: None, diff --git a/clients/openframe-client/src/services/github_download_service.rs b/clients/openframe-client/src/services/github_download_service.rs index 7b8b09881b..5286ab3bb6 100644 --- a/clients/openframe-client/src/services/github_download_service.rs +++ b/clients/openframe-client/src/services/github_download_service.rs @@ -8,9 +8,32 @@ use bytes::Bytes; use reqwest::Client; use std::io::Cursor; use std::path::Path; +#[cfg(target_os = "macos")] +use std::path::{Component, PathBuf}; use tokio::time::Duration; use tracing::{info, warn}; +/// `Path::join` with an absolute entry replaces the base, and `..` is resolved by the OS. +#[cfg(target_os = "macos")] +fn safe_join(target_dir: &Path, entry_path: &Path) -> Result { + for component in entry_path.components() { + match component { + Component::Normal(_) | Component::CurDir => {} + _ => { + return Err(anyhow!( + "Refusing archive entry with an unsafe path: {}", + entry_path.display() + )) + } + } + } + Ok(target_dir.join(entry_path)) +} + +#[cfg(all(test, target_os = "macos"))] +#[path = "github_download_service_tests.rs"] +mod tests; + #[derive(Clone)] pub struct GithubDownloadService { http_client: Client, @@ -409,7 +432,7 @@ impl GithubDownloadService { for entry_result in archive.entries().context("Failed to read tar entries")? { let mut entry = entry_result.context("Failed to read tar entry")?; let path = entry.path().context("Failed to get entry path")?; - let dest_path = target_dir.join(&path); + let dest_path = safe_join(target_dir, &path)?; if entry.header().entry_type().is_dir() { fs::create_dir_all(&dest_path).with_context(|| { diff --git a/clients/openframe-client/src/services/github_download_service_tests.rs b/clients/openframe-client/src/services/github_download_service_tests.rs new file mode 100644 index 0000000000..b2d3f59510 --- /dev/null +++ b/clients/openframe-client/src/services/github_download_service_tests.rs @@ -0,0 +1,60 @@ +use super::safe_join; +use std::path::{Path, PathBuf}; + +fn target() -> PathBuf { + PathBuf::from("/Applications") +} + +fn join(entry: &str) -> Option { + safe_join(&target(), Path::new(entry)).ok() +} + +#[test] +fn accepts_leading_current_dir() { + assert_eq!( + join("./OpenFrame.app/Contents/MacOS/openframe-chat"), + Some(target().join("./OpenFrame.app/Contents/MacOS/openframe-chat")) + ); +} + +#[test] +fn accepts_a_plain_nested_entry() { + assert_eq!( + join("OpenFrame.app/Contents/Info.plist"), + Some(target().join("OpenFrame.app/Contents/Info.plist")) + ); +} + +#[test] +fn accepts_a_bare_filename() { + assert_eq!( + join("pax_global_header"), + Some(target().join("pax_global_header")) + ); +} + +#[test] +fn rejects_an_absolute_entry() { + assert!(join("/Library/LaunchDaemons/evil.plist").is_none()); +} + +#[test] +fn rejects_parent_traversal() { + assert!(join("../../../../etc/cron.d/evil").is_none()); +} +#[test] +fn rejects_interior_parent_traversal() { + assert!(join("OpenFrame.app/../../../etc/passwd").is_none()); +} + +#[test] +fn rejects_a_root_relative_entry() { + assert!(join("//srv/evil").is_none()); +} +#[test] +fn rejection_names_the_offending_entry() { + let err = safe_join(&target(), Path::new("../escape")) + .expect_err("traversal must be refused") + .to_string(); + assert!(err.contains("../escape"), "unhelpful error: {err}"); +} diff --git a/clients/openframe-client/src/services/initial_configuration_service.rs b/clients/openframe-client/src/services/initial_configuration_service.rs index 77ec07d0a9..07d8ce86d0 100644 --- a/clients/openframe-client/src/services/initial_configuration_service.rs +++ b/clients/openframe-client/src/services/initial_configuration_service.rs @@ -107,7 +107,7 @@ impl InitialConfigurationService { pub fn save(&self, config: &InitialConfiguration) -> Result<()> { let config_json = serde_json::to_string_pretty(config) .context("Failed to serialize initial configuration to JSON")?; - fs::write(&self.config_file_path, config_json).with_context(|| { + crate::utils::fs::atomic_write(&self.config_file_path, config_json).with_context(|| { format!( "Failed to write initial configuration file: {:?}", self.config_file_path diff --git a/clients/openframe-client/src/services/last_known_good_service.rs b/clients/openframe-client/src/services/last_known_good_service.rs index c5b409a93b..5687fc51d5 100644 --- a/clients/openframe-client/src/services/last_known_good_service.rs +++ b/clients/openframe-client/src/services/last_known_good_service.rs @@ -135,10 +135,10 @@ impl LastKnownGoodService { } Some(anchor_version) => { warn!( - "Rollback protection degraded: reserve missing, running {} below anchor {} — rebuilding reserve from running binary, anchor unchanged", + "Rollback protection degraded: reserve missing and running {} does not match anchor {} — leaving the reserve unset for a verified update to rebuild", running_version, anchor_version ); - self.copy_running_to_reserve() + Ok(()) } None => { info!( diff --git a/clients/openframe-client/src/services/tool_agent_update_service.rs b/clients/openframe-client/src/services/tool_agent_update_service.rs index 3e31287f01..b89160ee42 100644 --- a/clients/openframe-client/src/services/tool_agent_update_service.rs +++ b/clients/openframe-client/src/services/tool_agent_update_service.rs @@ -4,10 +4,10 @@ use crate::models::tool_agent_update_message::{AssetUpdate, ToolAgentUpdateMessa use crate::models::{Installation, InstalledAsset, ToolRecordState}; use crate::platform::{ binary_writer, clear_aside_binary, detect_actual_installation, needs_migration, run_migration, - run_update, DirectoryManager, ToolUpdaterDeps, + run_update, system_service, DirectoryManager, ToolUpdaterDeps, }; use crate::services::agent_configuration_service::AgentConfigurationService; -use crate::services::tool_run_manager::ToolRunManager; +use crate::services::tool_run_manager::{ToolRunManager, UpdatingGuard}; use crate::services::GithubDownloadService; use crate::services::InstalledAgentMessagePublisher; use crate::services::InstalledToolsService; @@ -72,6 +72,21 @@ impl ToolAgentUpdateService { tool_agent_id, new_version ); + let tool_lock = self.tool_run_manager.tool_lock(tool_agent_id).await; + let _lock_guard = match tool_lock.try_lock_owned() { + Ok(guard) => guard, + Err(_) => { + info!( + "Tool {} busy with another operation, deferring update", + tool_agent_id + ); + anyhow::bail!( + "tool {} busy, deferring update for redelivery", + tool_agent_id + ); + } + }; + // Check if tool is installed let mut installed_tool = match self .installed_tools_service @@ -190,8 +205,8 @@ impl ToolAgentUpdateService { return Ok(()); } - // Mark as updating once for all updates - self.tool_run_manager.mark_updating(tool_agent_id).await; + let _updating = + UpdatingGuard::acquire(&self.tool_run_manager, tool_agent_id, Some(_lock_guard)).await; // A Standard->GuiApp migration self-relaunches, so only relaunch here if it was already a GUI app. let was_gui_before_update = @@ -206,9 +221,6 @@ impl ToolAgentUpdateService { ) .await; - // Clear updating flag - for Standard tools the run manager relaunches them via this flag. - self.tool_run_manager.clear_updating(tool_agent_id).await; - if result.is_ok() && needs_repair { if let Err(e) = self .installed_tools_service @@ -248,24 +260,65 @@ impl ToolAgentUpdateService { ); self.do_tool_update(new_version, message, installed_tool) .await?; - } else if !assets_to_update.is_empty() { - // Only assets to update - stop tool once before all asset updates + } + + let mut stopped_for_assets = false; + if !needs_tool_update && !assets_to_update.is_empty() { + // Only assets to update - stop tool once before all asset updates. info!(tool_id = %tool_agent_id, "Stopping tool for asset updates"); + stopped_for_assets = true; self.tool_kill_service - .stop_tool(tool_agent_id) + .stop_installed_tool(installed_tool, false) .await .with_context(|| { format!("Failed to stop tool {} for asset updates", tool_agent_id) })?; } - // 2. Asset updates (tool already stopped by tool_update or above) + // 2. Asset updates (tool already stopped by tool_update or above). + let mut asset_result = Ok(()); for asset in assets_to_update { - self.do_asset_update(tool_agent_id, asset, installed_tool) - .await?; + asset_result = self + .do_asset_update(tool_agent_id, asset, installed_tool) + .await; + if asset_result.is_err() { + break; + } } - Ok(()) + if stopped_for_assets { + match &installed_tool.installation { + Installation::Service { service_name, .. } => { + self.restart_service_after_assets(tool_agent_id, service_name) + .await; + } + #[cfg(target_os = "macos")] + Installation::GuiApp { .. } => { + if let Err(e) = self + .tool_run_manager + .run_new_tool(installed_tool.clone()) + .await + { + warn!(tool_id = %tool_agent_id, "Failed to relaunch GuiApp after asset updates: {:#}", e); + } + } + _ => {} + } + } + + asset_result + } + + async fn restart_service_after_assets(&self, tool_agent_id: &str, service_name: &str) { + info!(tool_id = %tool_agent_id, service = %service_name, "Restarting service tool after asset updates"); + + if let Err(e) = system_service::start_service(service_name).await { + error!( + tool_id = %tool_agent_id, service = %service_name, + "Service tool could not be restarted after asset updates - it will stay down until a reinstall or reboot: {:#}", + e + ); + } } async fn do_tool_update( diff --git a/clients/openframe-client/src/services/tool_restart_service.rs b/clients/openframe-client/src/services/tool_restart_service.rs index c3f8809377..2dafa7791d 100644 --- a/clients/openframe-client/src/services/tool_restart_service.rs +++ b/clients/openframe-client/src/services/tool_restart_service.rs @@ -1,7 +1,7 @@ use crate::config::service_stop::TOOL_RESTART_TIMEOUT_SECS; use crate::models::{Installation, InstalledTool}; use crate::platform::system_service; -use crate::services::tool_run_manager::ToolRunManager; +use crate::services::tool_run_manager::{ToolRunManager, UpdatingGuard}; use crate::services::InstalledToolsService; use crate::services::ToolKillService; use anyhow::{Context, Result}; @@ -16,27 +16,6 @@ pub enum RestartOutcome { Busy, } -/// Clears the updating flag on drop (surviving cancellation and panic), releasing the tool lock only after the flag clears. -struct UpdatingGuard { - tool_run_manager: ToolRunManager, - tool_agent_id: String, - lock_guard: Option>, -} - -impl Drop for UpdatingGuard { - fn drop(&mut self) { - if let Ok(handle) = tokio::runtime::Handle::try_current() { - let manager = self.tool_run_manager.clone(); - let tool_agent_id = self.tool_agent_id.clone(); - let lock_guard = self.lock_guard.take(); - handle.spawn(async move { - manager.clear_updating(&tool_agent_id).await; - drop(lock_guard); - }); - } - } -} - #[derive(Clone)] pub struct ToolRestartService { installed_tools_service: InstalledToolsService, @@ -64,12 +43,8 @@ impl ToolRestartService { Ok(guard) => guard, Err(_) => return Ok(RestartOutcome::Busy), }; - self.tool_run_manager.mark_updating(tool_agent_id).await; - let _updating = UpdatingGuard { - tool_run_manager: self.tool_run_manager.clone(), - tool_agent_id: tool_agent_id.to_string(), - lock_guard: Some(lock_guard), - }; + let _updating = + UpdatingGuard::acquire(&self.tool_run_manager, tool_agent_id, Some(lock_guard)).await; // Hard cap so a wedged OS call can't hold the flag/lock forever and freeze callers (e.g. mesh self-heal). let outcome = tokio::time::timeout( Duration::from_secs(TOOL_RESTART_TIMEOUT_SECS), diff --git a/clients/openframe-client/src/services/tool_run_manager.rs b/clients/openframe-client/src/services/tool_run_manager.rs index db4ad4f398..00538029bc 100644 --- a/clients/openframe-client/src/services/tool_run_manager.rs +++ b/clients/openframe-client/src/services/tool_run_manager.rs @@ -29,6 +29,8 @@ use windows::{ }; const RETRY_DELAY_SECONDS: u64 = 5; +/// Ceiling for the escalating retry delay after a leftover process refuses to die. +const KILL_RETRY_MAX_DELAY_SECONDS: u64 = 300; /// Under the supervision lock, so a concurrent resume either finds the loop alive or relaunches its tool. async fn shutdown_break( @@ -554,10 +556,7 @@ impl ToolRunManager { } let tool_id = tool.tool_agent_id.clone(); info!(tool_id = %tool_id, "Relaunching tool supervisor after the aborted update"); - if let Err(e) = self.run_tool(tool, false).await { - warn!(tool_id = %tool_id, "Failed to relaunch tool after the aborted update: {:#}", e); - self.clear_running_tool(&tool_id).await; - } + self.run_tool(tool, false).await; } info!("Tool run manager: supervision resumed after the aborted update"); Ok(()) @@ -664,11 +663,12 @@ impl ToolRunManager { } for tool in tools { - if self.try_mark_running(&tool.tool_agent_id).await { - info!("Running tool {}", tool.tool_agent_id); - self.run_tool(tool, false).await?; + let tool_id = tool.tool_agent_id.clone(); + if self.try_mark_running(&tool_id).await { + info!("Running tool {}", tool_id); + self.run_tool(tool, false).await; } else { - warn!("Tool {} is already running - skipping", tool.tool_agent_id); + warn!("Tool {} is already running - skipping", tool_id); } } @@ -685,7 +685,8 @@ impl ToolRunManager { } info!("Running new single tool {}", installed_tool.tool_agent_id); - self.run_tool(installed_tool, true).await + self.run_tool(installed_tool, true).await; + Ok(()) } async fn try_mark_running(&self, tool_id: &str) -> bool { @@ -704,27 +705,14 @@ impl ToolRunManager { } #[allow(unused_variables)] - async fn run_tool(&self, tool: InstalledTool, new_tool: bool) -> Result<()> { + async fn run_tool(&self, tool: InstalledTool, new_tool: bool) { if tool.installation.is_service() { info!( "Installation::Service for {} - self-managed, skipping launch", tool.tool_agent_id ); self.clear_running_tool(&tool.tool_agent_id).await; - return Ok(()); - } - - #[cfg(not(target_os = "windows"))] - self.tool_kill_service - .stop_tool(&tool.tool_agent_id) - .await?; - - // Windows GUI apps are owned by the HKLM Run autorun, not us — never kill them. - #[cfg(target_os = "windows")] - if !tool.installation.is_gui_app() { - self.tool_kill_service - .stop_tool(&tool.tool_agent_id) - .await?; + return; } let updating_tools = self.updating_tools.clone(); @@ -732,11 +720,13 @@ impl ToolRunManager { let params_processor = self.params_processor.clone(); let running_tools = self.running_tools.clone(); let installed_tools_service = self.installed_tools_service.clone(); + let tool_kill_service = self.tool_kill_service.clone(); let mut installation = tool.installation.clone(); let mut run_command_args = tool.run_command_args.clone(); tokio::spawn(async move { let mut launch_backoff = FailureLogBackoff::new(); + let mut kill_leftovers_now = true; loop { // Self-update in progress — stop the loop entirely if shutdown_break(&shutting_down, &running_tools, &tool.tool_agent_id).await { @@ -768,6 +758,28 @@ impl ToolRunManager { let log_attempt = launch_backoff.should_log(); + // Windows GUI apps belong to the HKLM Run autorun — never kill them. + #[cfg(target_os = "windows")] + let kill_leftovers = !installation.is_gui_app(); + #[cfg(not(target_os = "windows"))] + let kill_leftovers = true; + + if kill_leftovers && kill_leftovers_now { + if let Err(e) = tool_kill_service.stop_tool(&tool.tool_agent_id).await { + let failures = launch_backoff.record_failure(log_attempt); + let delay = + (RETRY_DELAY_SECONDS * failures).min(KILL_RETRY_MAX_DELAY_SECONDS); + if log_attempt { + error!(tool_id = %tool.tool_agent_id, failed_attempts = failures, + "Failed to stop leftover tool processes - not launching, retrying in {} seconds: {:#}", + delay, e); + } + sleep(Duration::from_secs(delay)).await; + continue; + } + kill_leftovers_now = false; + } + let processed_args = match params_processor.process(&tool.tool_agent_id, run_command_args.clone()) { Ok(args) => args, @@ -880,7 +892,14 @@ impl ToolRunManager { if let Err(e) = crate::platform::preferences_writer::write(bid, prefs) { - error!(tool_id = %tool.tool_agent_id, "Failed to write preferences: {:#}", e); + let failures = launch_backoff.record_failure(log_attempt); + if log_attempt { + error!(tool_id = %tool.tool_agent_id, failed_attempts = failures, + "Failed to write GuiApp preferences - not launching, retrying in {} seconds: {:#}", + RETRY_DELAY_SECONDS, e); + } + sleep(Duration::from_secs(RETRY_DELAY_SECONDS)).await; + continue; } if tool.tool_agent_id == "openframe-chat" { @@ -1033,11 +1052,47 @@ impl ToolRunManager { "Failed to wait for tool process - restarting in {} seconds: {:#}", RETRY_DELAY_SECONDS, e); } } + kill_leftovers_now = true; sleep(Duration::from_secs(RETRY_DELAY_SECONDS)).await; } }); + } +} - Ok(()) +/// Clears the updating flag on drop, releasing the tool lock only once the flag is clear. +pub struct UpdatingGuard { + tool_run_manager: ToolRunManager, + tool_agent_id: String, + lock_guard: Option>, +} + +impl UpdatingGuard { + /// Marks the tool as updating and holds `lock_guard` (if any) until the flag clears again. + pub async fn acquire( + tool_run_manager: &ToolRunManager, + tool_agent_id: &str, + lock_guard: Option>, + ) -> Self { + tool_run_manager.mark_updating(tool_agent_id).await; + Self { + tool_run_manager: tool_run_manager.clone(), + tool_agent_id: tool_agent_id.to_string(), + lock_guard, + } + } +} + +impl Drop for UpdatingGuard { + fn drop(&mut self) { + if let Ok(handle) = tokio::runtime::Handle::try_current() { + let manager = self.tool_run_manager.clone(); + let tool_agent_id = self.tool_agent_id.clone(); + let lock_guard = self.lock_guard.take(); + handle.spawn(async move { + manager.clear_updating(&tool_agent_id).await; + drop(lock_guard); + }); + } } }