agent_client_protocol_polyfill/mcp_over_acp/
mod.rs1mod 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#[derive(Debug)]
48pub(crate) enum BridgeMessage {
49 ConnectionReceived {
51 acp_id: String,
52 actor: BridgeConnectionActor,
53 connection: BridgeConnection,
54 },
55
56 ConnectionEstablished {
58 response: McpConnectResponse,
59 actor: BridgeConnectionActor,
60 connection: BridgeConnection,
61 },
62
63 ClientToServer {
65 connection_id: String,
66 message: Dispatch,
67 },
68
69 Disconnected {
71 notification: McpDisconnectNotification,
72 },
73}
74
75#[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#[derive(Debug, Clone, Default)]
98pub enum BridgeMode {
99 Stdio {
101 conductor_command: Vec<String>,
103 },
104
105 #[default]
107 Http,
108}
109
110#[derive(Debug)]
114pub struct McpOverAcpPolyfill {
115 mode: BridgeMode,
116}
117
118impl McpOverAcpPolyfill {
119 #[must_use]
121 pub fn http() -> Self {
122 Self {
123 mode: BridgeMode::Http,
124 }
125 }
126
127 #[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 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 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 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#[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 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
312struct 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 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(¬ification.connection_id);
397 connection.send_notification_to(Client, notification)?;
398 }
399 }
400 }
401 Ok(())
402 }
403}