Skip to main content

agent_client_protocol_polyfill/mcp_over_acp/
mod.rs

1//! MCP-over-ACP polyfill proxy.
2//!
3//! This proxy bridges MCP-over-ACP transport for agents that don't support
4//! `mcpCapabilities.acp` natively. It sits in the proxy chain and:
5//!
6//! - Intercepts `NewSessionRequest` to transform `McpServer::Http` entries with `acp:` URLs
7//!   into localhost TCP bridges
8//! - Handles `_mcp/connect`, `_mcp/message`, `_mcp/disconnect` by routing through those bridges
9//!
10//! # Usage
11//!
12//! ```rust,ignore
13//! use agent_client_protocol_polyfill::mcp_over_acp::McpOverAcpPolyfill;
14//!
15//! // Add to a conductor proxy chain
16//! let conductor = ConductorImpl::new_agent(
17//!     "conductor",
18//!     ProxiesAndAgent::new(my_agent).proxy(McpOverAcpPolyfill::http()),
19//!     McpBridgeMode::default(),
20//! );
21//! ```
22
23mod actor;
24pub(crate) mod http;
25pub(crate) mod stdio;
26
27use std::collections::HashMap;
28use std::path::PathBuf;
29
30use agent_client_protocol::schema::v1::{
31    McpServer, McpServerHttp, McpServerStdio, NewSessionRequest,
32};
33use agent_client_protocol::schema::{
34    InitializeProxyRequest, McpConnectRequest, McpConnectResponse, McpDisconnectNotification,
35    McpOverAcpMessage,
36};
37use agent_client_protocol::{
38    Agent, Client, Conductor, ConnectTo, ConnectionTo, Dispatch, Proxy, Role,
39};
40use futures::{SinkExt, channel::mpsc};
41use tokio::net::TcpListener;
42use tracing::info;
43
44use self::actor::BridgeConnectionActor;
45
46/// Internal messages for the polyfill's bridge management.
47#[derive(Debug)]
48pub(crate) enum BridgeMessage {
49    /// A new TCP connection was accepted and needs an ACP connection ID.
50    ConnectionReceived {
51        acp_id: String,
52        actor: BridgeConnectionActor,
53        connection: BridgeConnection,
54    },
55
56    /// ACP connection ID received — spawn the actor and store the connection.
57    ConnectionEstablished {
58        response: McpConnectResponse,
59        actor: BridgeConnectionActor,
60        connection: BridgeConnection,
61    },
62
63    /// MCP message from a bridge client that needs to be forwarded over ACP.
64    ClientToServer {
65        connection_id: String,
66        message: Dispatch,
67    },
68
69    /// Bridge client disconnected.
70    Disconnected {
71        notification: McpDisconnectNotification,
72    },
73}
74
75/// Connection handle for sending messages to an MCP client via a bridge.
76#[derive(Clone, Debug)]
77#[allow(dead_code)]
78pub(crate) struct BridgeConnection {
79    to_mcp_client_tx: mpsc::Sender<Dispatch>,
80}
81
82impl BridgeConnection {
83    pub fn new(to_mcp_client_tx: mpsc::Sender<Dispatch>) -> Self {
84        Self { to_mcp_client_tx }
85    }
86
87    #[allow(dead_code)]
88    pub async fn send(&mut self, message: Dispatch) -> Result<(), agent_client_protocol::Error> {
89        self.to_mcp_client_tx
90            .send(message)
91            .await
92            .map_err(|_| agent_client_protocol::Error::internal_error())
93    }
94}
95
96/// Mode for the MCP bridge transport.
97#[derive(Debug, Clone, Default)]
98pub enum BridgeMode {
99    /// Use stdio-based MCP bridge with a subprocess.
100    Stdio {
101        /// Command and args to spawn bridge processes.
102        conductor_command: Vec<String>,
103    },
104
105    /// Use HTTP-based MCP bridge (default).
106    #[default]
107    Http,
108}
109
110/// MCP-over-ACP polyfill proxy.
111///
112/// Bridges MCP-over-ACP transport for agents that don't support `mcpCapabilities.acp`.
113#[derive(Debug)]
114pub struct McpOverAcpPolyfill {
115    mode: BridgeMode,
116}
117
118impl McpOverAcpPolyfill {
119    /// Create a polyfill using HTTP bridge mode.
120    #[must_use]
121    pub fn http() -> Self {
122        Self {
123            mode: BridgeMode::Http,
124        }
125    }
126
127    /// Create a polyfill using stdio bridge mode.
128    #[must_use]
129    pub fn stdio(conductor_command: Vec<String>) -> Self {
130        Self {
131            mode: BridgeMode::Stdio { conductor_command },
132        }
133    }
134}
135
136impl ConnectTo<Conductor> for McpOverAcpPolyfill {
137    async fn connect_to(
138        self,
139        client: impl ConnectTo<Proxy>,
140    ) -> Result<(), agent_client_protocol::Error> {
141        let (bridge_tx, bridge_rx) = mpsc::channel(128);
142        let mode = self.mode;
143
144        Proxy
145            .builder()
146            .name("mcp-over-acp-polyfill")
147            .with_responder(BridgeResponder {
148                bridge_tx: bridge_tx.clone(),
149                bridge_rx,
150                bridge_connections: HashMap::new(),
151            })
152            .on_receive_request_from(
153                Client,
154                async move |request: InitializeProxyRequest,
155                            responder,
156                            cx: ConnectionTo<Conductor>| {
157                    // Forward initialize to successor, then set mcpCapabilities.acp = true
158                    // in the response to advertise that we handle MCP-over-ACP.
159                    cx.send_request_to(Agent, request.initialize)
160                        .on_receiving_result(async move |result| {
161                            responder.respond_with_result(result.map(|mut response| {
162                                response.agent_capabilities.mcp_capabilities.acp = true;
163                                response
164                            }))
165                        })
166                },
167                agent_client_protocol::on_receive_request!(),
168            )
169            .on_receive_request_from(
170                Client,
171                {
172                    let bridge_tx = bridge_tx.clone();
173                    async move |mut request: NewSessionRequest,
174                                responder,
175                                cx: ConnectionTo<Conductor>| {
176                        // Transform acp: URLs in MCP servers
177                        let mut listeners = BridgeListeners::default();
178                        for mcp_server in &mut request.mcp_servers {
179                            listeners
180                                .transform_mcp_server(cx.clone(), mcp_server, &bridge_tx, &mode)
181                                .await?;
182                        }
183                        // Forward modified request to successor
184                        cx.send_request_to(Agent, request)
185                            .forward_response_to(responder)
186                    }
187                },
188                agent_client_protocol::on_receive_request!(),
189            )
190            .connect_to(client)
191            .await
192    }
193}
194
195/// Manages active bridge listeners (TCP listeners for acp: URLs).
196#[derive(Default, Debug)]
197struct BridgeListeners {
198    listeners: HashMap<String, BridgeListener>,
199}
200
201#[derive(Clone, Debug)]
202struct BridgeListener {
203    server: McpServer,
204}
205
206impl BridgeListeners {
207    /// Transform an MCP server with `acp:` URL into a bridged localhost server.
208    async fn transform_mcp_server(
209        &mut self,
210        connection: ConnectionTo<impl Role>,
211        mcp_server: &mut McpServer,
212        bridge_tx: &mpsc::Sender<BridgeMessage>,
213        mode: &BridgeMode,
214    ) -> Result<(), agent_client_protocol::Error> {
215        let McpServer::Http(http) = mcp_server else {
216            return Ok(());
217        };
218
219        if !http.url.starts_with("acp:") {
220            return Ok(());
221        }
222
223        if !http.headers.is_empty() {
224            return Err(agent_client_protocol::Error::internal_error());
225        }
226
227        let name = http.name.clone();
228        let url = http.url.clone();
229
230        info!(
231            server_name = %name,
232            acp_id = %url,
233            "Detected MCP server with ACP transport, spawning TCP bridge"
234        );
235
236        let transformed = self
237            .spawn_bridge(connection, &name, &url, bridge_tx, mode)
238            .await?;
239        *mcp_server = transformed;
240        Ok(())
241    }
242
243    async fn spawn_bridge(
244        &mut self,
245        connection: ConnectionTo<impl Role>,
246        server_name: &str,
247        acp_id: &str,
248        bridge_tx: &mpsc::Sender<BridgeMessage>,
249        mode: &BridgeMode,
250    ) -> anyhow::Result<McpServer> {
251        if let Some(listener) = self.listeners.get(acp_id) {
252            return Ok(listener.server.clone());
253        }
254
255        let tcp_listener = TcpListener::bind("127.0.0.1:0").await?;
256        let tcp_port = tcp_listener.local_addr()?.port();
257
258        info!(acp_id = acp_id, tcp_port, "Bound listener for MCP bridge");
259
260        let new_server = match mode {
261            BridgeMode::Stdio { conductor_command } => McpServer::Stdio(
262                McpServerStdio::new(
263                    server_name.to_string(),
264                    PathBuf::from(&conductor_command[0]),
265                )
266                .args(
267                    conductor_command[1..]
268                        .iter()
269                        .cloned()
270                        .chain(vec!["mcp".to_string(), format!("{tcp_port}")])
271                        .collect::<Vec<_>>(),
272                ),
273            ),
274
275            BridgeMode::Http => McpServer::Http(McpServerHttp::new(
276                server_name.to_string(),
277                format!("http://localhost:{tcp_port}"),
278            )),
279        };
280
281        self.listeners.insert(
282            acp_id.to_string(),
283            BridgeListener {
284                server: new_server.clone(),
285            },
286        );
287
288        connection.spawn({
289            let acp_id = acp_id.to_string();
290            let bridge_tx = bridge_tx.clone();
291            let mode = mode.clone();
292            async move {
293                info!(
294                    acp_id = acp_id,
295                    tcp_port, "now accepting bridge connections"
296                );
297                match mode {
298                    BridgeMode::Stdio {
299                        conductor_command: _,
300                    } => stdio::run_tcp_listener(tcp_listener, acp_id, bridge_tx).await,
301                    BridgeMode::Http => {
302                        http::run_http_listener(tcp_listener, acp_id, bridge_tx).await
303                    }
304                }
305            }
306        })?;
307
308        Ok(new_server)
309    }
310}
311
312/// Responder that runs alongside the proxy, managing bridge state.
313struct BridgeResponder {
314    bridge_tx: mpsc::Sender<BridgeMessage>,
315    bridge_rx: mpsc::Receiver<BridgeMessage>,
316    bridge_connections: HashMap<String, BridgeConnection>,
317}
318
319impl std::fmt::Debug for BridgeResponder {
320    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
321        f.debug_struct("BridgeResponder")
322            .field("bridge_connections", &self.bridge_connections.len())
323            .finish_non_exhaustive()
324    }
325}
326
327impl agent_client_protocol::RunWithConnectionTo<Conductor> for BridgeResponder {
328    async fn run_with_connection_to(
329        mut self,
330        connection: ConnectionTo<Conductor>,
331    ) -> Result<(), agent_client_protocol::Error> {
332        use futures::StreamExt;
333
334        while let Some(message) = self.bridge_rx.next().await {
335            match message {
336                BridgeMessage::ConnectionReceived {
337                    acp_id,
338                    actor,
339                    connection: bridge_conn,
340                } => {
341                    // Send _mcp/connect request back through the chain.
342                    // When the response arrives, send ConnectionEstablished back to ourselves.
343                    connection
344                        .send_request_to(Client, McpConnectRequest { acp_id, meta: None })
345                        .on_receiving_result({
346                            let mut bridge_tx = self.bridge_tx.clone();
347                            async move |result| match result {
348                                Ok(response) => bridge_tx
349                                    .send(BridgeMessage::ConnectionEstablished {
350                                        response,
351                                        actor,
352                                        connection: bridge_conn,
353                                    })
354                                    .await
355                                    .map_err(|_| agent_client_protocol::Error::internal_error()),
356                                Err(_) => Ok(()),
357                            }
358                        })?;
359                }
360
361                BridgeMessage::ConnectionEstablished {
362                    response: McpConnectResponse { connection_id, .. },
363                    actor,
364                    connection: bridge_conn,
365                } => {
366                    self.bridge_connections
367                        .insert(connection_id.clone(), bridge_conn);
368                    connection.spawn(actor.run(connection_id))?;
369                }
370
371                BridgeMessage::ClientToServer {
372                    connection_id,
373                    message,
374                } => {
375                    let wrapped = message.map(
376                        |request, responder| {
377                            (
378                                McpOverAcpMessage {
379                                    connection_id: connection_id.clone(),
380                                    message: request,
381                                    meta: None,
382                                },
383                                responder,
384                            )
385                        },
386                        |notification| McpOverAcpMessage {
387                            connection_id: connection_id.clone(),
388                            message: notification,
389                            meta: None,
390                        },
391                    );
392                    connection.send_proxied_message_to(Client, wrapped)?;
393                }
394
395                BridgeMessage::Disconnected { notification } => {
396                    self.bridge_connections.remove(&notification.connection_id);
397                    connection.send_notification_to(Client, notification)?;
398                }
399            }
400        }
401        Ok(())
402    }
403}