From 410a78983e48c15d670174214cb2813375057bb7 Mon Sep 17 00:00:00 2001 From: Abhijay Jain Date: Thu, 2 Jul 2026 13:36:41 +0530 Subject: [PATCH] fix(acp): propagate agent cleanup errors on session delete (#10112) Signed-off-by: Abhijay Jain --- crates/goose-server/src/routes/agent.rs | 10 +++++-- crates/goose/src/acp/server.rs | 22 ++++++++++++-- .../goose/src/acp/server/manage_sessions.rs | 10 +++++-- crates/goose/src/acp/server/new_session.rs | 21 ++++++++++++-- crates/goose/src/execution/manager.rs | 29 +++++++++++++++++++ crates/goose/src/gateway/handler.rs | 4 ++- 6 files changed, 86 insertions(+), 10 deletions(-) diff --git a/crates/goose-server/src/routes/agent.rs b/crates/goose-server/src/routes/agent.rs index 20c91d0db..155c621f1 100644 --- a/crates/goose-server/src/routes/agent.rs +++ b/crates/goose-server/src/routes/agent.rs @@ -799,8 +799,14 @@ async fn restart_agent_internal( session_id: &str, session: &Session, ) -> Result, ErrorResponse> { - // Remove existing agent (ignore error if not found) - let _ = state.agent_manager.remove_session(session_id).await; + state + .agent_manager + .remove_session_if_loaded(session_id) + .await + .map_err(|e| ErrorResponse { + message: format!("Failed to remove in-memory agent for session {session_id}: {e}"), + status: StatusCode::INTERNAL_SERVER_ERROR, + })?; let agent = state .get_agent_for_route(session_id.to_string()) diff --git a/crates/goose/src/acp/server.rs b/crates/goose/src/acp/server.rs index 14c5b446b..5365f1b9b 100644 --- a/crates/goose/src/acp/server.rs +++ b/crates/goose/src/acp/server.rs @@ -1178,7 +1178,10 @@ impl GooseAcpAgent { .await .internal_err_ctx("Failed to update session")?; - let _ = self.agent_manager.remove_session(&session_id).await; + self.agent_manager + .remove_session_if_loaded(&session_id) + .await + .internal_err_ctx("Failed to remove in-memory agent")?; session = self .session_manager @@ -2337,7 +2340,17 @@ impl GooseAcpAgent { if self.closed_session_ids.lock().await.contains(session_id) { self.sessions.lock().await.remove(session_id); - let _ = self.agent_manager.remove_session(session_id).await; + if let Err(error) = self + .agent_manager + .remove_session_if_loaded(session_id) + .await + { + tracing::warn!( + session_id, + %error, + "Failed to remove in-memory agent for closed session" + ); + } } } @@ -2969,7 +2982,10 @@ impl GooseAcpAgent { sessions.remove(session_id); drop(sessions); - let _ = self.agent_manager.remove_session(session_id).await; + self.agent_manager + .remove_session_if_loaded(session_id) + .await + .internal_err_ctx("Failed to remove in-memory agent")?; info!(session_id = %session_id, "ACP session closed"); Ok(CloseSessionResponse::new()) diff --git a/crates/goose/src/acp/server/manage_sessions.rs b/crates/goose/src/acp/server/manage_sessions.rs index d4d34c65f..3c79da1e2 100644 --- a/crates/goose/src/acp/server/manage_sessions.rs +++ b/crates/goose/src/acp/server/manage_sessions.rs @@ -104,7 +104,10 @@ impl GooseAcpAgent { .await .internal_err()?; self.sessions.lock().await.remove(&req.session_id); - let _ = self.agent_manager.remove_session(&req.session_id).await; + self.agent_manager + .remove_session_if_loaded(&req.session_id) + .await + .internal_err_ctx("Failed to remove in-memory agent")?; Ok(EmptyResponse {}) } @@ -254,7 +257,10 @@ impl GooseAcpAgent { .await .internal_err()?; self.sessions.lock().await.remove(&req.session_id); - let _ = self.agent_manager.remove_session(&req.session_id).await; + self.agent_manager + .remove_session_if_loaded(&req.session_id) + .await + .internal_err_ctx("Failed to remove in-memory agent")?; Ok(EmptyResponse {}) } diff --git a/crates/goose/src/acp/server/new_session.rs b/crates/goose/src/acp/server/new_session.rs index 74ad554e5..68934047e 100644 --- a/crates/goose/src/acp/server/new_session.rs +++ b/crates/goose/src/acp/server/new_session.rs @@ -11,6 +11,7 @@ use agent_client_protocol::{Client, ConnectionTo}; use goose_providers::model::ModelConfig; use std::collections::HashMap; use std::path::PathBuf; +use tracing::warn; struct InitialSessionConfig { provider: String, @@ -92,9 +93,25 @@ impl GooseAcpAgent { } async fn cleanup_failed_new_session(&self, session_id: &str) { - let _ = self.session_manager.delete_session(session_id).await; + if let Err(error) = self.session_manager.delete_session(session_id).await { + warn!( + session_id, + %error, + "Failed to delete session during new-session cleanup" + ); + } self.sessions.lock().await.remove(session_id); - let _ = self.agent_manager.remove_session(session_id).await; + if let Err(error) = self + .agent_manager + .remove_session_if_loaded(session_id) + .await + { + warn!( + session_id, + %error, + "Failed to remove in-memory agent during new-session cleanup" + ); + } } async fn configure_new_session( diff --git a/crates/goose/src/execution/manager.rs b/crates/goose/src/execution/manager.rs index 4eea38862..c21752e7f 100644 --- a/crates/goose/src/execution/manager.rs +++ b/crates/goose/src/execution/manager.rs @@ -328,6 +328,21 @@ impl AgentManager { Ok(()) } + /// Drops an in-memory agent when one is loaded for `session_id`. + pub async fn remove_session_if_loaded(&self, session_id: &str) -> Result<()> { + if let Some(token) = self.cancel_tokens.write().await.remove(session_id) { + token.cancel(); + } + let mut sessions = self.sessions.write().await; + if sessions.pop(session_id).is_none() { + return Ok(()); + } + drop(sessions); + self.prune_creation_lock(session_id).await; + info!("Removed session {}", session_id); + Ok(()) + } + pub async fn has_session(&self, session_id: &str) -> bool { self.sessions.read().await.contains(session_id) } @@ -484,6 +499,20 @@ mod tests { assert!(manager.remove_session(&session).await.is_err()); } + #[tokio::test] + async fn test_remove_session_if_loaded() { + let temp_dir = TempDir::new().unwrap(); + let manager = create_test_manager(&temp_dir).await; + let session = String::from("remove-if-loaded-test"); + + manager.remove_session_if_loaded(&session).await.unwrap(); + + manager.get_or_create_agent(session.clone()).await.unwrap(); + manager.remove_session_if_loaded(&session).await.unwrap(); + assert!(!manager.has_session(&session).await); + manager.remove_session_if_loaded(&session).await.unwrap(); + } + #[tokio::test] async fn test_concurrent_access() { let temp_dir = TempDir::new().unwrap(); diff --git a/crates/goose/src/gateway/handler.rs b/crates/goose/src/gateway/handler.rs index bbe1ecf3c..6ccf3d34a 100644 --- a/crates/goose/src/gateway/handler.rs +++ b/crates/goose/src/gateway/handler.rs @@ -308,7 +308,9 @@ impl GatewayHandler { // extension processes don't linger. let extensions_changed = self.sync_session_config(&session).await?; if extensions_changed { - let _ = self.agent_manager.remove_session(session_id).await; + self.agent_manager + .remove_session_if_loaded(session_id) + .await?; } let agent = self