Skip to main content

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}