agent_client_protocol_rmcp/
lib.rs1use agent_client_protocol::mcp_server::{McpConnectionTo, McpServer, McpServerConnect};
43#[cfg(feature = "unstable_mcp_over_acp")]
44use agent_client_protocol::mcp_server::{McpOutcome, McpRequest, McpRequestContext, McpService};
45use agent_client_protocol::role;
46use agent_client_protocol::{ByteStreams, ConnectTo, DynConnectTo, NullRun, Role};
47#[cfg(feature = "unstable_mcp_over_acp")]
48use futures::future::BoxFuture;
49use futures_concurrency::future::TryJoin as _;
50use rmcp::ServiceExt;
51#[cfg(feature = "unstable_mcp_over_acp")]
52use std::sync::{Arc, OnceLock};
53use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
54
55mod builder;
56#[cfg(feature = "unstable_mcp_over_acp")]
57mod native;
58
59pub use agent_client_protocol::mcp_server::{EnabledTools, McpTool};
60pub use agent_client_protocol::{tool_fn, tool_fn_mut};
61pub use builder::McpServerBuilder;
62
63pub trait McpServerExt<Counterpart: Role> {
65 fn builder(name: impl ToString) -> McpServerBuilder<Counterpart, NullRun> {
67 McpServerBuilder::new(name.to_string())
68 }
69
70 fn from_rmcp<S>(
80 name: impl ToString,
81 new_fn: impl Fn() -> S + Send + Sync + 'static,
82 ) -> McpServer<Counterpart, NullRun>
83 where
84 S: rmcp::Service<rmcp::RoleServer>,
85 {
86 struct RmcpServer<F> {
87 name: String,
88 new_fn: F,
89 }
90
91 #[cfg(feature = "unstable_mcp_over_acp")]
92 struct SharedRmcp<F, S> {
93 new_fn: Arc<F>,
94 service: OnceLock<Arc<S>>,
95 }
96
97 #[cfg(feature = "unstable_mcp_over_acp")]
98 impl<Counterpart, F, S> McpService<Counterpart> for SharedRmcp<F, S>
99 where
100 Counterpart: Role,
101 F: Fn() -> S + Send + Sync + 'static,
102 S: rmcp::Service<rmcp::RoleServer>,
103 {
104 fn execute(
105 &self,
106 request: McpRequest,
107 context: McpRequestContext<Counterpart>,
108 ) -> BoxFuture<'static, Result<McpOutcome, agent_client_protocol::Error>> {
109 let service = self.service.get_or_init(|| Arc::new((self.new_fn)()));
110 native::execute(service.clone(), request, context)
111 }
112 }
113
114 impl<Counterpart, F, S> McpServerConnect<Counterpart> for RmcpServer<F>
115 where
116 Counterpart: Role,
117 F: Fn() -> S + Send + Sync + 'static,
118 S: rmcp::Service<rmcp::RoleServer>,
119 {
120 fn name(&self) -> String {
121 self.name.clone()
122 }
123
124 fn connect(
125 &self,
126 _cx: McpConnectionTo<Counterpart>,
127 ) -> DynConnectTo<role::mcp::Client> {
128 let service = (self.new_fn)();
129 DynConnectTo::new(RmcpServerComponent { service })
130 }
131 }
132
133 #[cfg(feature = "unstable_mcp_over_acp")]
134 {
135 let new_fn = Arc::new(new_fn);
138 McpServer::new_service_with_standalone(
139 SharedRmcp {
140 new_fn: new_fn.clone(),
141 service: OnceLock::new(),
142 },
143 RmcpServer {
144 name: name.to_string(),
145 new_fn: move || new_fn(),
146 },
147 NullRun,
148 )
149 }
150 #[cfg(not(feature = "unstable_mcp_over_acp"))]
151 {
152 McpServer::new(
153 RmcpServer {
154 name: name.to_string(),
155 new_fn,
156 },
157 NullRun,
158 )
159 }
160 }
161}
162
163impl<Counterpart: Role> McpServerExt<Counterpart> for McpServer<Counterpart> {}
164
165struct RmcpServerComponent<S> {
167 service: S,
168}
169
170impl<S> ConnectTo<role::mcp::Client> for RmcpServerComponent<S>
171where
172 S: rmcp::Service<rmcp::RoleServer>,
173{
174 async fn connect_to(
175 self,
176 client: impl ConnectTo<role::mcp::Server>,
177 ) -> Result<(), agent_client_protocol::Error> {
178 let (mcp_server_stream, mcp_client_stream) = tokio::io::duplex(8192);
180 let (mcp_server_read, mcp_server_write) = tokio::io::split(mcp_server_stream);
181 let (mcp_client_read, mcp_client_write) = tokio::io::split(mcp_client_stream);
182
183 let bytes_to_acp = async {
184 let byte_streams =
186 ByteStreams::new(mcp_client_write.compat_write(), mcp_client_read.compat());
187
188 ConnectTo::<role::mcp::Client>::connect_to(byte_streams, client).await
189 };
190
191 let bytes_to_rmcp = async {
192 let running_server = self
194 .service
195 .serve((mcp_server_read, mcp_server_write))
196 .await
197 .map_err(agent_client_protocol::Error::into_internal_error)?;
198
199 running_server
201 .waiting()
202 .await
203 .map(|_quit_reason| ())
204 .map_err(agent_client_protocol::Error::into_internal_error)
205 };
206
207 (bytes_to_acp, bytes_to_rmcp).try_join().await?;
208 Ok(())
209 }
210}
211
212#[cfg(test)]
213mod tests {
214 use agent_client_protocol::{mcp_server::McpServer, role};
215
216 use crate::McpServerExt as _;
217
218 #[test]
219 fn builds_standalone_server_without_acp_transport_feature() {
220 let _server: McpServer<role::mcp::Client, _> = McpServer::builder("standalone").build();
221 }
222}