perf(acp): parallelize extension loading in ACP server (#8098)
Signed-off-by: Douwe Osinga <douwe@squareup.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Co-authored-by: Douwe Osinga <douwe@squareup.com>
This commit is contained in:
@@ -479,21 +479,30 @@ impl GooseAcpAgent {
|
|||||||
let skip_developer = acp_developer.is_some();
|
let skip_developer = acp_developer.is_some();
|
||||||
let sid_str = session_id.map(|s| s.0.to_string());
|
let sid_str = session_id.map(|s| s.0.to_string());
|
||||||
|
|
||||||
for ext in extensions {
|
if skip_developer {
|
||||||
if skip_developer && ext.name() == "developer" {
|
extensions.retain(|ext| ext.name() != "developer");
|
||||||
continue;
|
|
||||||
}
|
|
||||||
let name = ext.name().to_string();
|
|
||||||
match agent
|
|
||||||
.extension_manager
|
|
||||||
.add_extension(ext, None, None, sid_str.as_deref())
|
|
||||||
.await
|
|
||||||
{
|
|
||||||
Ok(_) => info!(extension = %name, "extension loaded"),
|
|
||||||
Err(e) => warn!(extension = %name, error = %e, "extension load failed"),
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let ext_manager = &agent.extension_manager;
|
||||||
|
let extension_futures = extensions
|
||||||
|
.into_iter()
|
||||||
|
.map(|ext| {
|
||||||
|
let ext_manager = Arc::clone(ext_manager);
|
||||||
|
let sid = sid_str.clone();
|
||||||
|
async move {
|
||||||
|
let name = ext.name().to_string();
|
||||||
|
match ext_manager
|
||||||
|
.add_extension(ext, None, None, sid.as_deref())
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(_) => info!(extension = %name, "extension loaded"),
|
||||||
|
Err(e) => warn!(extension = %name, error = %e, "extension load failed"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
futures::future::join_all(extension_futures).await;
|
||||||
|
|
||||||
if let Some((client, config)) = acp_developer {
|
if let Some((client, config)) = acp_developer {
|
||||||
let info = client.get_info().cloned();
|
let info = client.get_info().cloned();
|
||||||
agent
|
agent
|
||||||
@@ -932,10 +941,11 @@ impl GooseAcpAgent {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn add_mcp_extensions(
|
async fn add_mcp_extensions(
|
||||||
agent: &Agent,
|
agent: &Arc<Agent>,
|
||||||
mcp_servers: Vec<McpServer>,
|
mcp_servers: Vec<McpServer>,
|
||||||
session_id: &str,
|
session_id: &str,
|
||||||
) -> Result<(), sacp::Error> {
|
) -> Result<(), sacp::Error> {
|
||||||
|
let mut configs = Vec::with_capacity(mcp_servers.len());
|
||||||
for mcp_server in mcp_servers {
|
for mcp_server in mcp_servers {
|
||||||
let config = match mcp_server_to_extension_config(mcp_server) {
|
let config = match mcp_server_to_extension_config(mcp_server) {
|
||||||
Ok(c) => c,
|
Ok(c) => c,
|
||||||
@@ -943,10 +953,24 @@ impl GooseAcpAgent {
|
|||||||
return Err(sacp::Error::invalid_params().data(msg));
|
return Err(sacp::Error::invalid_params().data(msg));
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let name = config.name().to_string();
|
configs.push(config);
|
||||||
if let Err(e) = agent.add_extension(config, session_id).await {
|
}
|
||||||
return Err(sacp::Error::internal_error()
|
|
||||||
.data(format!("Failed to add MCP server '{}': {}", name, e)));
|
if configs.is_empty() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
let results = agent
|
||||||
|
.add_extensions_bulk(configs, session_id)
|
||||||
|
.await
|
||||||
|
.map_err(|e| sacp::Error::internal_error().data(e.to_string()))?;
|
||||||
|
for result in &results {
|
||||||
|
if !result.success {
|
||||||
|
let error_msg = result.error.as_deref().unwrap_or("unknown error");
|
||||||
|
return Err(sacp::Error::internal_error().data(format!(
|
||||||
|
"Failed to add MCP server '{}': {}",
|
||||||
|
result.name, error_msg
|
||||||
|
)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|||||||
@@ -761,6 +761,71 @@ impl Agent {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Load multiple extensions in parallel, persisting state once at the end.
|
||||||
|
///
|
||||||
|
/// Unlike `add_extension`, this avoids per-extension persistence and acquires
|
||||||
|
/// the container lock once upfront to prevent serialisation of the parallel futures.
|
||||||
|
pub async fn add_extensions_bulk(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
extensions: Vec<ExtensionConfig>,
|
||||||
|
session_id: &str,
|
||||||
|
) -> anyhow::Result<Vec<ExtensionLoadResult>> {
|
||||||
|
let working_dir = match self
|
||||||
|
.config
|
||||||
|
.session_manager
|
||||||
|
.get_session(session_id, false)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(session) => Some(session.working_dir),
|
||||||
|
Err(e) => {
|
||||||
|
warn!("Failed to get session for bulk load: {}", e);
|
||||||
|
None
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let container = self.container.lock().await.clone();
|
||||||
|
|
||||||
|
let extension_futures = extensions
|
||||||
|
.into_iter()
|
||||||
|
.map(|config| {
|
||||||
|
let ext_manager = Arc::clone(&self.extension_manager);
|
||||||
|
let working_dir = working_dir.clone();
|
||||||
|
let container = container.clone();
|
||||||
|
let sid = session_id.to_string();
|
||||||
|
|
||||||
|
async move {
|
||||||
|
let name = config.name().to_string();
|
||||||
|
match ext_manager
|
||||||
|
.add_extension(config, working_dir, container.as_ref(), Some(&sid))
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(_) => ExtensionLoadResult {
|
||||||
|
name,
|
||||||
|
success: true,
|
||||||
|
error: None,
|
||||||
|
},
|
||||||
|
Err(e) => {
|
||||||
|
let error_msg = e.to_string();
|
||||||
|
warn!("Failed to load extension {}: {}", name, error_msg);
|
||||||
|
ExtensionLoadResult {
|
||||||
|
name,
|
||||||
|
success: false,
|
||||||
|
error: Some(error_msg),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
|
let results = futures::future::join_all(extension_futures).await;
|
||||||
|
|
||||||
|
if results.iter().any(|r| r.success) {
|
||||||
|
self.persist_extension_state(session_id).await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(results)
|
||||||
|
}
|
||||||
|
|
||||||
async fn add_extension_inner(
|
async fn add_extension_inner(
|
||||||
&self,
|
&self,
|
||||||
extension: ExtensionConfig,
|
extension: ExtensionConfig,
|
||||||
|
|||||||
Reference in New Issue
Block a user