mermaid_cli/mcp/client.rs
1//! MCP protocol client — higher-level API over a [`Transport`] (stdio child
2//! process or Streamable HTTP endpoint).
3//!
4//! Implements the three protocol methods we need:
5//! - `initialize` — handshake and capability negotiation
6//! - `tools/list` — discover available tools
7//! - `tools/call` — invoke a tool and get results
8
9use anyhow::{Result, anyhow};
10use serde::{Deserialize, Serialize};
11use serde_json::{Value, json};
12
13use super::transport::Transport;
14
15/// MCP protocol client for a single server connection.
16pub struct McpClient {
17 transport: Transport,
18 /// Server info from initialization
19 pub server_info: Option<ServerInfo>,
20 /// Set once [`Self::shutdown`] runs, so a later `call_tool` returns a clean
21 /// "stopped" error instead of a broken-pipe transport error — the manager
22 /// keeps the entry in its frozen map, so the client outlives its process.
23 shutdown: std::sync::atomic::AtomicBool,
24}
25
26/// Info returned by the server during initialization
27#[derive(Debug, Clone, Serialize, Deserialize)]
28pub struct ServerInfo {
29 pub name: String,
30 pub version: Option<String>,
31}
32
33/// A tool definition discovered from an MCP server
34#[derive(Debug, Clone)]
35pub struct McpToolDef {
36 pub name: String,
37 pub description: String,
38 pub input_schema: Value,
39 /// The server's `annotations.readOnlyHint` — an UNTRUSTED self-declaration
40 /// that the tool has no side effects. Absent ⇒ false, i.e. write-shaped
41 /// (fail closed). Feeds the external-writes policy floor.
42 pub read_only_hint: bool,
43}
44
45/// Result of calling an MCP tool
46#[derive(Debug, Clone)]
47pub struct McpToolResult {
48 pub content: Vec<ContentBlock>,
49 pub is_error: bool,
50}
51
52/// A content block in an MCP tool result.
53///
54/// Per the 2025-11-25 spec, servers may return text, image, audio,
55/// `resource_link` (URI reference), or embedded resource content. Older
56/// servers only emit text/image.
57#[derive(Debug, Clone)]
58pub enum ContentBlock {
59 Text(String),
60 Image {
61 data: String,
62 mime_type: String,
63 },
64 /// Audio content — base64-encoded data + mime type (e.g., `audio/wav`).
65 /// Routed to the model's image attachment channel for now; adapters
66 /// that don't support audio will silently drop the bytes but keep
67 /// the text hint from the tool output.
68 Audio {
69 data: String,
70 mime_type: String,
71 },
72 /// URI reference to an external resource. Rendered as text for the
73 /// model so it can follow up with another tool call if needed.
74 ResourceLink {
75 uri: String,
76 name: Option<String>,
77 description: Option<String>,
78 mime_type: Option<String>,
79 },
80 /// Embedded resource — same shape as a `read_resource` response.
81 /// Either `text` or `blob` (base64) is present depending on the
82 /// resource's kind. Rendered as text for the model.
83 Resource {
84 uri: String,
85 mime_type: Option<String>,
86 text: Option<String>,
87 blob: Option<String>,
88 },
89}
90
91/// Parse one `tools/list` entry. `None` for nameless entries (skipped, as
92/// before). Annotations are optional per the MCP spec; a missing or
93/// non-boolean `readOnlyHint` is treated as false — write-shaped, fail
94/// closed.
95fn tool_def_from_json(tool: &Value) -> Option<McpToolDef> {
96 let name = tool.get("name").and_then(|v| v.as_str())?;
97 if name.is_empty() {
98 return None;
99 }
100 Some(McpToolDef {
101 name: name.to_string(),
102 description: tool
103 .get("description")
104 .and_then(|v| v.as_str())
105 .unwrap_or("")
106 .to_string(),
107 input_schema: tool
108 .get("inputSchema")
109 .cloned()
110 .unwrap_or_else(|| json!({"type": "object", "properties": {}})),
111 read_only_hint: tool
112 .pointer("/annotations/readOnlyHint")
113 .and_then(|v| v.as_bool())
114 .unwrap_or(false),
115 })
116}
117
118impl McpClient {
119 /// Create a new MCP client wrapping a transport.
120 pub(super) fn new(transport: Transport) -> Self {
121 Self {
122 transport,
123 server_info: None,
124 shutdown: std::sync::atomic::AtomicBool::new(false),
125 }
126 }
127
128 /// Perform the MCP initialization handshake.
129 ///
130 /// Sends `initialize` request with our client info and protocol version,
131 /// then sends `notifications/initialized` to signal readiness.
132 ///
133 /// # Errors
134 ///
135 /// The `initialize` request failing — transport, timeout, or a JSON-RPC
136 /// error from the server — and the `notifications/initialized` send. A
137 /// server that omits `serverInfo` or negotiates down to an older
138 /// `protocolVersion` is not an error: the name falls back to `"unknown"`
139 /// and whatever version it names is what subsequent requests carry.
140 pub async fn initialize(&mut self) -> Result<ServerInfo> {
141 let result = self
142 .transport
143 .send_request(
144 "initialize",
145 json!({
146 // MCP spec version as of 2026-04. Servers negotiate
147 // down to older versions if they don't support this;
148 // spec requires them to respond with their latest
149 // supported version, which we currently accept
150 // silently. Bump when MCP ships a newer revision
151 // with features we depend on.
152 "protocolVersion": "2025-11-25",
153 "capabilities": {},
154 "clientInfo": {
155 "name": "mermaid",
156 "version": env!("CARGO_PKG_VERSION"),
157 }
158 }),
159 )
160 .await?;
161
162 // Parse server info
163 let server_info = ServerInfo {
164 name: result
165 .pointer("/serverInfo/name")
166 .and_then(|v| v.as_str())
167 .unwrap_or("unknown")
168 .to_string(),
169 version: result
170 .pointer("/serverInfo/version")
171 .and_then(|v| v.as_str())
172 .map(|s| s.to_string()),
173 };
174
175 // Record the negotiated protocol version BEFORE the initialized
176 // notification: over HTTP every request after initialize — including
177 // that notification — must carry the MCP-Protocol-Version header.
178 if let Some(version) = result.get("protocolVersion").and_then(|v| v.as_str()) {
179 self.transport.set_protocol_version(version);
180 }
181
182 // Send initialized notification
183 self.transport
184 .send_notification("notifications/initialized", json!({}))
185 .await?;
186
187 self.server_info = Some(server_info.clone());
188 Ok(server_info)
189 }
190
191 /// Discover all tools available from this server, following `nextCursor`
192 /// pagination so a server that pages its tool list isn't silently truncated
193 /// to page one. Bounded by a page cap so a server that echoes a stuck cursor
194 /// can't loop forever.
195 ///
196 /// # Errors
197 ///
198 /// Any page's `tools/list` request failing, and a response with no
199 /// `tools` array. A tool entry that does not parse is skipped rather than
200 /// failing the discovery, and hitting the page cap returns what was
201 /// collected — so a short list is not necessarily the server's whole
202 /// catalog.
203 pub async fn list_tools(&self) -> Result<Vec<McpToolDef>> {
204 const MAX_PAGES: usize = 100;
205 let mut tools = Vec::new();
206 let mut cursor: Option<String> = None;
207
208 for _ in 0..MAX_PAGES {
209 let params = match &cursor {
210 Some(c) => json!({ "cursor": c }),
211 None => json!({}),
212 };
213 let result = self.transport.send_request("tools/list", params).await?;
214
215 let tools_array = result
216 .get("tools")
217 .and_then(|v| v.as_array())
218 .ok_or_else(|| anyhow!("MCP tools/list response missing 'tools' array"))?;
219
220 for tool in tools_array {
221 if let Some(def) = tool_def_from_json(tool) {
222 tools.push(def);
223 }
224 }
225
226 match result.get("nextCursor").and_then(|v| v.as_str()) {
227 Some(next) if !next.is_empty() => cursor = Some(next.to_string()),
228 _ => break,
229 }
230 }
231
232 Ok(tools)
233 }
234
235 /// Call a tool on this server and return the result.
236 ///
237 /// # Errors
238 ///
239 /// The `tools/call` request failing: transport, the tool-call timeout, or
240 /// a JSON-RPC error. A tool that runs and reports failure is not among
241 /// them — that is `isError` on the returned [`McpToolResult`], which the
242 /// model is meant to see and react to.
243 pub async fn call_tool(&self, name: &str, arguments: &Value) -> Result<McpToolResult> {
244 let params = json!({
245 "name": name,
246 "arguments": arguments,
247 });
248
249 let result = self
250 .transport
251 .send_request_with_timeout("tools/call", params, Transport::tool_call_timeout_secs())
252 .await?;
253
254 let is_error = result
255 .get("isError")
256 .and_then(|v| v.as_bool())
257 .unwrap_or(false);
258
259 let content_array = result
260 .get("content")
261 .and_then(|v| v.as_array())
262 .cloned()
263 .unwrap_or_default();
264
265 let content = content_array
266 .iter()
267 .filter_map(parse_content_block)
268 .collect();
269
270 Ok(McpToolResult { content, is_error })
271 }
272
273 /// Shut down the transport (kills the server process).
274 pub async fn shutdown(&self) {
275 self.shutdown
276 .store(true, std::sync::atomic::Ordering::Release);
277 self.transport.shutdown().await;
278 }
279
280 /// `true` once [`Self::shutdown`] has run. The manager checks this so a
281 /// `call_tool` to a stopped-but-still-registered server returns a clean
282 /// error rather than a broken-pipe transport failure.
283 pub fn is_shutdown(&self) -> bool {
284 self.shutdown.load(std::sync::atomic::Ordering::Acquire)
285 }
286}
287
288/// Decode one `tools/call` content block. Unknown block types fall back to
289/// their `text` field; a block with nothing usable yields `None`.
290fn parse_content_block(block: &Value) -> Option<ContentBlock> {
291 let str_field = |v: &Value, key: &str| v.get(key).and_then(|v| v.as_str()).map(String::from);
292 let block_type = block.get("type").and_then(|v| v.as_str()).unwrap_or("");
293 match block_type {
294 "text" => str_field(block, "text").map(ContentBlock::Text),
295 "image" => Some(ContentBlock::Image {
296 data: str_field(block, "data").unwrap_or_default(),
297 mime_type: str_field(block, "mimeType").unwrap_or_else(|| "image/png".to_string()),
298 }),
299 "audio" => Some(ContentBlock::Audio {
300 data: str_field(block, "data").unwrap_or_default(),
301 mime_type: str_field(block, "mimeType").unwrap_or_else(|| "audio/wav".to_string()),
302 }),
303 "resource_link" => {
304 let uri = str_field(block, "uri").filter(|uri| !uri.is_empty())?;
305 Some(ContentBlock::ResourceLink {
306 uri,
307 name: str_field(block, "name"),
308 description: str_field(block, "description"),
309 mime_type: str_field(block, "mimeType"),
310 })
311 },
312 "resource" => {
313 // Embedded resource — nested under `resource`.
314 let res = block.get("resource")?;
315 let uri = str_field(res, "uri").filter(|uri| !uri.is_empty())?;
316 Some(ContentBlock::Resource {
317 uri,
318 mime_type: str_field(res, "mimeType"),
319 text: str_field(res, "text"),
320 blob: str_field(res, "blob"),
321 })
322 },
323 // Unknown content type — treat as text if it has a text field.
324 _ => str_field(block, "text").map(ContentBlock::Text),
325 }
326}
327
328#[cfg(test)]
329mod tests {
330 use super::*;
331
332 #[test]
333 fn tool_def_parses_read_only_hint() {
334 // Annotated read-only tool carries the hint through.
335 let def = tool_def_from_json(&json!({
336 "name": "get_thing",
337 "description": "d",
338 "inputSchema": {"type": "object"},
339 "annotations": {"readOnlyHint": true}
340 }))
341 .unwrap();
342 assert!(def.read_only_hint);
343
344 // Absent annotations (the common case) ⇒ write-shaped, fail closed.
345 let def = tool_def_from_json(&json!({"name": "send_thing"})).unwrap();
346 assert!(!def.read_only_hint);
347 assert_eq!(
348 def.input_schema,
349 json!({"type": "object", "properties": {}})
350 );
351
352 // Non-boolean hints and nameless entries are rejected safely.
353 let def = tool_def_from_json(&json!({
354 "name": "odd",
355 "annotations": {"readOnlyHint": "yes"}
356 }))
357 .unwrap();
358 assert!(!def.read_only_hint);
359 assert!(tool_def_from_json(&json!({"description": "nameless"})).is_none());
360 }
361}