diff --git a/crates/goose/src/acp/server.rs b/crates/goose/src/acp/server.rs index ac908d18..cb171431 100644 --- a/crates/goose/src/acp/server.rs +++ b/crates/goose/src/acp/server.rs @@ -30,7 +30,7 @@ use crate::session::session_manager::SessionType; use crate::session::{EnabledExtensionsState, Session, SessionManager}; use anyhow::Result; use fs_err as fs; -use futures::future::BoxFuture; +use futures::future::{BoxFuture, Either}; use goose_acp_macros::custom_methods; use rmcp::model::{CallToolResult, RawContent, ResourceContents, Role}; use sacp::schema::{ @@ -112,12 +112,26 @@ struct GooseAcpSession { pending_working_dir: Option, } +/// Progress stages signalled by the background agent setup task via the watch +/// channel. `ProviderReady` fires as soon as the provider (and goose-mode) +/// are initialized — before extensions finish loading. `FullyReady` fires +/// once every extension has been loaded (or failed). +#[derive(Clone)] +enum AgentSetupProgress { + /// Provider is initialized; extensions are still loading in the background. + ProviderReady(Arc), + /// Provider *and* all extensions are initialized. + FullyReady(Arc), +} + +type AgentSetupSignal = Option>; + /// The agent may still be initializing in the background (extension loading, /// provider setup). Callers that need the live agent (e.g. `on_prompt`) await /// the handle; callers that only need the session metadata can proceed without it. enum AgentHandle { Ready(Arc), - Loading(tokio::sync::watch::Receiver, String>>>), + Loading(tokio::sync::watch::Receiver), } struct AgentSetupRequest { @@ -878,7 +892,7 @@ impl GooseAcpAgent { fn spawn_agent_setup( &self, cx: &ConnectionTo, - agent_tx: tokio::sync::watch::Sender, String>>>, + agent_tx: tokio::sync::watch::Sender, req: AgentSetupRequest, ) { let AgentSetupRequest { @@ -912,7 +926,20 @@ impl GooseAcpAgent { tokio::spawn(async move { let t_setup = std::time::Instant::now(); debug!(target: "perf", sid = %sid, "perf: agent_setup start (background)"); - let result: Result<(), String> = async { + + // Shared config — read once, used by both phases. + let config = match Config::new(config_dir.join(CONFIG_YAML_NAME), "goose") { + Ok(c) => c, + Err(e) => { + let msg = e.to_string(); + error!(error = %msg, "Background agent setup failed (config)"); + let _ = agent_tx.send(Some(Err(msg))); + return; + } + }; + + // ── Phase 1: create agent + init provider (fast, ~55ms) ────── + let phase1: Result, String> = async { let agent = Arc::new(Agent::with_config(AgentConfig::new( session_manager, permission_manager, @@ -922,11 +949,59 @@ impl GooseAcpAgent { GoosePlatform::GooseCli, ))); - let config_path = config_dir.join(CONFIG_YAML_NAME); - let mut extensions = Config::new(&config_path, "goose") - .ok() - .map(|c| get_enabled_extensions_with_config(&c)) - .unwrap_or_default(); + // Init provider — reuse the pre-resolved name + model when + // available (already computed in on_new_session), otherwise + // fall back to reading config (e.g. load_session path). + let t_prov = std::time::Instant::now(); + let (provider_name, model_config) = match resolved_provider { + Some(resolved) => resolved, + None => resolve_provider_and_model_from_config(&config, &goose_session).await?, + }; + let ext_state = EnabledExtensionsState::extensions_or_default( + Some(&goose_session.extension_data), + &config, + ); + let provider = match prebuilt_provider { + Some(provider) => provider, + None => provider_factory(provider_name.to_string(), model_config, ext_state) + .await + .map_err(|e| e.to_string())?, + }; + agent + .update_provider(provider.clone(), &goose_session.id) + .await + .map_err(|e| e.to_string())?; + + agent + .update_goose_mode(goose_mode, &internal_session_id) + .await + .map_err(|e| e.to_string())?; + debug!(target: "perf", sid = %sid, ms = t_prov.elapsed().as_millis() as u64, "perf: agent_setup provider_init"); + + Ok(agent) + } + .await; + + let agent = match phase1 { + Ok(agent) => { + // Signal ProviderReady — unblocks setProvider / update_provider + // while extensions continue loading below. + let _ = + agent_tx.send(Some(Ok(AgentSetupProgress::ProviderReady(agent.clone())))); + debug!(target: "perf", sid = %sid, ms = t_setup.elapsed().as_millis() as u64, "perf: agent_setup provider_ready (signalled)"); + agent + } + Err(e) => { + error!(error = %e, "Background agent setup failed (provider init)"); + debug!(target: "perf", sid = %sid, ms = t_setup.elapsed().as_millis() as u64, "perf: agent_setup failed (provider)"); + let _ = agent_tx.send(Some(Err(e))); + return; + } + }; + + // ── Phase 2: load extensions (slow, may take seconds) ──────── + let phase2: Result<(), String> = async { + let mut extensions = get_enabled_extensions_with_config(&config); extensions.extend(builtins.iter().map(|b| builtin_to_extension_config(b))); let acp_developer = if (client_fs_capabilities.read_text_file @@ -1027,37 +1102,6 @@ impl GooseAcpAgent { .await; } - // Init provider — reuse the pre-resolved name + model when - // available (already computed in on_new_session), otherwise - // fall back to reading config (e.g. load_session path). - let t_prov = std::time::Instant::now(); - let config = Config::new(config_dir.join(CONFIG_YAML_NAME), "goose") - .map_err(|e| e.to_string())?; - let (provider_name, model_config) = match resolved_provider { - Some(resolved) => resolved, - None => resolve_provider_and_model_from_config(&config, &goose_session).await?, - }; - let ext_state = EnabledExtensionsState::extensions_or_default( - Some(&goose_session.extension_data), - &config, - ); - let provider = match prebuilt_provider { - Some(provider) => provider, - None => provider_factory(provider_name.to_string(), model_config, ext_state) - .await - .map_err(|e| e.to_string())?, - }; - agent - .update_provider(provider.clone(), &goose_session.id) - .await - .map_err(|e| e.to_string())?; - - agent - .update_goose_mode(goose_mode, &internal_session_id) - .await - .map_err(|e| e.to_string())?; - debug!(target: "perf", sid = %sid, ms = t_prov.elapsed().as_millis() as u64, "perf: agent_setup provider_init"); - let t_mcp = std::time::Instant::now(); let mcp_count = mcp_servers.len(); GooseAcpAgent::add_mcp_extensions(&agent, mcp_servers, &internal_session_id) @@ -1071,41 +1115,39 @@ impl GooseAcpAgent { "perf: agent_setup mcp_extensions" ); - // Apply any working directory that was set while we were loading. - { - let mut locked = sessions.lock().await; - if let Some(session) = locked.get_mut(session_id.0.as_ref()) { - if let Some(dir) = session.pending_working_dir.take() { - agent.extension_manager.update_working_dir(&dir).await; - } - session.agent = AgentHandle::Ready(agent.clone()); - } - } - - let _ = agent_tx.send(Some(Ok(agent))); - Ok(()) } .await; - match &result { - Ok(()) => debug!( - target: "perf", - sid = %sid, - ms = t_setup.elapsed().as_millis() as u64, - "perf: agent_setup done" - ), - Err(e) => { - error!(error = %e, "Background agent setup failed"); - debug!( - target: "perf", - sid = %sid, - ms = t_setup.elapsed().as_millis() as u64, - "perf: agent_setup failed" - ); - let _ = agent_tx.send(Some(Err(e.clone()))); + if let Err(e) = &phase2 { + // Extension failures are non-fatal — individual failures are + // already logged as warnings. Log the top-level error but + // don't block the session: the provider is ready and the agent + // is usable. + error!(error = %e, "Background agent setup: extension phase had errors"); + } + + // Promote the handle to Ready and apply any working directory that + // was set while we were loading — regardless of phase-2 outcome, + // since the agent (with its provider) is fully usable. + { + let mut locked = sessions.lock().await; + if let Some(session) = locked.get_mut(session_id.0.as_ref()) { + if let Some(dir) = session.pending_working_dir.take() { + agent.extension_manager.update_working_dir(&dir).await; + } + session.agent = AgentHandle::Ready(agent.clone()); } } + + let _ = agent_tx.send(Some(Ok(AgentSetupProgress::FullyReady(agent)))); + debug!( + target: "perf", + sid = %sid, + ms = t_setup.elapsed().as_millis() as u64, + "perf: agent_setup done{}", + if phase2.is_err() { " (with extension errors)" } else { "" } + ); }); } @@ -1601,8 +1643,7 @@ impl GooseAcpAgent { let internal_session_id = goose_session.id.clone(); - let (agent_tx, agent_rx) = - tokio::sync::watch::channel::, String>>>(None); + let (agent_tx, agent_rx) = tokio::sync::watch::channel::(None); let session = GooseAcpSession { agent: AgentHandle::Loading(agent_rx), @@ -1708,39 +1749,82 @@ impl GooseAcpAgent { }) } + /// Look up the session and return the agent if already ready, or the watch + /// receiver if still loading. Optionally sets a cancellation token on the + /// session (needed by `on_prompt`). + async fn get_agent_or_receiver( + &self, + thread_id: &str, + cancel_token: Option, + ) -> Result, tokio::sync::watch::Receiver>, sacp::Error> + { + let mut sessions = self.sessions.lock().await; + let session = sessions.get_mut(thread_id).ok_or_else(|| { + sacp::Error::resource_not_found(Some(thread_id.to_string())) + .data(format!("Session not found: {}", thread_id)) + })?; + if let Some(token) = cancel_token { + session.cancel_token = Some(token); + } + match &session.agent { + AgentHandle::Ready(agent) => Ok(Either::Left(agent.clone())), + AgentHandle::Loading(rx) => Ok(Either::Right(rx.clone())), + } + } + + /// Wait until the agent is **fully ready** (provider + all extensions). + /// Most callers (e.g. `on_prompt`, `on_get_tools`) should use this. async fn get_session_agent( &self, thread_id: &str, cancel_token: Option, ) -> Result, sacp::Error> { - let mut rx = { - let mut sessions = self.sessions.lock().await; - let session = sessions.get_mut(thread_id).ok_or_else(|| { - sacp::Error::resource_not_found(Some(thread_id.to_string())) - .data(format!("Session not found: {}", thread_id)) - })?; - if let Some(token) = cancel_token { - session.cancel_token = Some(token); - } - match &session.agent { - AgentHandle::Ready(agent) => return Ok(agent.clone()), - AgentHandle::Loading(rx) => rx.clone(), - } + let mut rx = match self.get_agent_or_receiver(thread_id, cancel_token).await? { + Either::Left(agent) => return Ok(agent), + Either::Right(rx) => rx, }; - // Drop the lock while we wait for the background setup to finish. - // spawn_agent_setup promotes the handle to Ready before signalling. - let agent = { - let guard = rx.wait_for(|v| v.is_some()).await.map_err(|_| { + // Wait specifically for FullyReady (not just ProviderReady). + let guard = rx + .wait_for(|v| { + matches!( + v, + Some(Ok(AgentSetupProgress::FullyReady(_))) | Some(Err(_)) + ) + }) + .await + .map_err(|_| { sacp::Error::internal_error().data("Agent setup task was dropped".to_string()) })?; - guard - .as_ref() - .unwrap() - .as_ref() - .map_err(|e| sacp::Error::internal_error().data(e.clone()))? - .clone() + match guard.as_ref().unwrap() { + Ok(AgentSetupProgress::FullyReady(agent)) => Ok(agent.clone()), + Err(e) => Err(sacp::Error::internal_error().data(e.clone())), + // wait_for predicate excludes ProviderReady + _ => unreachable!(), + } + } + + /// Wait only until the **provider** is initialized. Extensions may still + /// be loading in the background. Use this for operations that only touch + /// the provider (e.g. `update_provider`, `set_model`, `build_config_update`). + async fn get_session_agent_provider_ready( + &self, + thread_id: &str, + ) -> Result, sacp::Error> { + let mut rx = match self.get_agent_or_receiver(thread_id, None).await? { + Either::Left(agent) => return Ok(agent), + Either::Right(rx) => rx, }; - Ok(agent) + // Any signal (ProviderReady, FullyReady, or Err) unblocks us. + let guard = rx.wait_for(|v| v.is_some()).await.map_err(|_| { + sacp::Error::internal_error().data("Agent setup task was dropped".to_string()) + })?; + match guard.as_ref().unwrap() { + Ok(progress) => match progress { + AgentSetupProgress::ProviderReady(agent) + | AgentSetupProgress::FullyReady(agent) => Ok(agent.clone()), + }, + Err(e) => Err(sacp::Error::internal_error().data(e.clone())), + } } async fn add_mcp_extensions( @@ -1981,8 +2065,7 @@ impl GooseAcpAgent { debug!(target: "perf", sid = %sid, ms = t_db.elapsed().as_millis() as u64, "perf: load_session db_updates"); // ── Register the session immediately with a Loading handle ── - let (agent_tx, agent_rx) = - tokio::sync::watch::channel::, String>>>(None); + let (agent_tx, agent_rx) = tokio::sync::watch::channel::(None); let session = GooseAcpSession { agent: AgentHandle::Loading(agent_rx), @@ -2234,8 +2317,8 @@ impl GooseAcpAgent { debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: set_model load_config"); let t_step = std::time::Instant::now(); - let agent = self.get_session_agent(thread_id, None).await?; - debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: set_model get_session_agent (waits for agent setup)"); + let agent = self.get_session_agent_provider_ready(thread_id).await?; + debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: set_model get_session_agent_provider_ready"); let t_step = std::time::Instant::now(); let current_provider = agent.provider().await.map_err(|e| { @@ -2325,7 +2408,7 @@ impl GooseAcpAgent { .get_session(&internal_id, false) .await .map_err(|e| sacp::Error::internal_error().data(e.to_string()))?; - let agent = self.get_session_agent(&thread_id.0, None).await?; + let agent = self.get_session_agent_provider_ready(&thread_id.0).await?; let provider = agent.provider().await.map_err(|e| { sacp::Error::internal_error().data(format!("Failed to get provider: {}", e)) })?; @@ -2367,7 +2450,7 @@ impl GooseAcpAgent { sacp::Error::invalid_params().data(format!("Invalid mode: {}", mode_id)) })?; - let agent = self.get_session_agent(thread_id, None).await?; + let agent = self.get_session_agent_provider_ready(thread_id).await?; agent .update_goose_mode(mode, &internal_id) .await @@ -2407,8 +2490,8 @@ impl GooseAcpAgent { debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: update_provider load_config"); let t_step = std::time::Instant::now(); - let agent = self.get_session_agent(thread_id, None).await?; - debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: update_provider get_session_agent (waits for agent setup)"); + let agent = self.get_session_agent_provider_ready(thread_id).await?; + debug!(target: "perf", sid = %sid, ms = t_step.elapsed().as_millis() as u64, "perf: update_provider get_session_agent_provider_ready"); let t_step = std::time::Instant::now(); let current_provider = agent.provider().await.map_err(|e| { @@ -2602,8 +2685,7 @@ impl GooseAcpAgent { let internal_session_id = goose_session.id.clone(); - let (agent_tx, agent_rx) = - tokio::sync::watch::channel::, String>>>(None); + let (agent_tx, agent_rx) = tokio::sync::watch::channel::(None); let session = GooseAcpSession { agent: AgentHandle::Loading(agent_rx),