| 1 | use std::collections::HashMap; |
| 2 | use std::collections::hash_map::DefaultHasher; |
| 3 | use std::hash::{Hash, Hasher}; |
| 4 | |
| 5 | use anyhow::{Context, Result, bail}; |
| 6 | use serde::de::DeserializeOwned; |
| 7 | use serde::{Deserialize, Serialize}; |
| 8 | use serde_json::{Value, json}; |
| 9 | |
| 10 | mod stdio_client; |
| 11 | // Unix-gated as well as test-gated: every helper in here builds and spawns a |
| 12 | // POSIX-sh script, so the tests that use it are `#[cfg(unix)]` and on Windows |
| 13 | // the whole module compiles to dead code, which `-D warnings` rejects. |
| 14 | #[cfg(all(test, unix))] |
| 15 | mod test_support; |
| 16 | |
| 17 | pub use stdio_client::ChildProcessMcpClient; |
| 18 | |
| 19 | /// Configuration for a single MCP server process. |
| 20 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 21 | pub struct McpServerConfig { |
| 22 | /// Unique server identifier used for tool name qualification. |
| 23 | pub name: String, |
| 24 | /// Path or name of the server executable. |
| 25 | pub command: String, |
| 26 | /// Command-line arguments passed to the server process. |
| 27 | #[serde(default)] |
| 28 | pub args: Vec<String>, |
| 29 | /// Environment variables set for the server process. |
| 30 | #[serde(default)] |
| 31 | pub env: HashMap<String, String>, |
| 32 | /// Whether this server should be started. Disabled servers are skipped. |
| 33 | #[serde(default = "default_true")] |
| 34 | pub enabled: bool, |
| 35 | } |
| 36 | |
| 37 | /// Filter controlling which tools from an MCP server are exposed. |
| 38 | /// |
| 39 | /// When `allow` is empty, all tools are permitted (unless denied). |
| 40 | /// `deny` takes precedence over `allow`. |
| 41 | #[derive(Debug, Clone, Serialize, Deserialize, Default)] |
| 42 | pub struct ToolFilter { |
| 43 | /// Tool names to expose. Empty means expose all. |
| 44 | #[serde(default)] |
| 45 | pub allow: Vec<String>, |
| 46 | /// Tool names to exclude. Takes precedence over `allow`. |
| 47 | #[serde(default)] |
| 48 | pub deny: Vec<String>, |
| 49 | } |
| 50 | |
| 51 | /// A complete MCP server definition including config and tool filter. |
| 52 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 53 | pub struct McpServerDefinition { |
| 54 | /// Server process configuration. |
| 55 | pub config: McpServerConfig, |
| 56 | /// Tool filter controlling which tools are exposed. |
| 57 | #[serde(default)] |
| 58 | pub filter: ToolFilter, |
| 59 | } |
| 60 | |
| 61 | /// Status of an individual MCP server during startup. |
| 62 | #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] |
| 63 | #[serde(rename_all = "snake_case")] |
| 64 | pub enum McpStartupStatus { |
| 65 | /// Server process is starting. |
| 66 | Starting, |
| 67 | /// Server is ready to accept tool calls. |
| 68 | Ready, |
| 69 | /// Server failed to start. |
| 70 | Failed { error: String }, |
| 71 | /// Server startup was cancelled (e.g., disabled in config). |
| 72 | Cancelled, |
| 73 | } |
| 74 | |
| 75 | /// Status update for a single MCP server during startup. |
| 76 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 77 | pub struct McpStartupUpdateEvent { |
| 78 | /// Name of the server this update pertains to. |
| 79 | pub server_name: String, |
| 80 | /// Current startup status. |
| 81 | pub status: McpStartupStatus, |
| 82 | } |
| 83 | |
| 84 | /// Record of an MCP server that failed to start. |
| 85 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 86 | pub struct McpStartupFailure { |
| 87 | /// Name of the server that failed. |
| 88 | pub server_name: String, |
| 89 | /// Error message describing the failure. |
| 90 | pub error: String, |
| 91 | } |
| 92 | |
| 93 | /// Summary emitted after all MCP servers have completed startup. |
| 94 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 95 | pub struct McpStartupCompleteEvent { |
| 96 | /// Names of servers that started successfully. |
| 97 | pub ready: Vec<String>, |
| 98 | /// Servers that failed with error details. |
| 99 | pub failed: Vec<McpStartupFailure>, |
| 100 | /// Names of servers that were skipped (disabled). |
| 101 | pub cancelled: Vec<String>, |
| 102 | } |
| 103 | |
| 104 | /// Describes a single tool provided by an MCP server. |
| 105 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 106 | pub struct McpToolDescriptor { |
| 107 | /// Name of the server providing this tool. |
| 108 | pub server_name: String, |
| 109 | /// Original tool name as reported by the server. |
| 110 | pub tool_name: String, |
| 111 | /// Fully qualified name (e.g., `mcp__server__tool`). |
| 112 | pub qualified_name: String, |
| 113 | /// Human-readable description of what the tool does. |
| 114 | pub description: Option<String>, |
| 115 | } |
| 116 | |
| 117 | /// Describes a resource provided by an MCP server. |
| 118 | #[derive(Debug, Clone, Serialize, Deserialize)] |
| 119 | pub struct McpResourceDescriptor { |
| 120 | /// Name of the server providing this resource. |
| 121 | pub server_name: String, |
| 122 | /// URI identifying the resource. |
| 123 | pub uri: String, |
| 124 | /// Human-readable description. |
| 125 | pub description: Option<String>, |
| 126 | } |
| 127 | |
| 128 | /// Trait abstracting an MCP client connection. |
| 129 | /// |
| 130 | /// Implementations handle communication with a single MCP server process. |
| 131 | pub trait McpManagedClient: Send + Sync { |
| 132 | /// List all tools provided by this server. |
| 133 | fn list_tools(&self) -> Result<Vec<McpToolDescriptor>>; |
| 134 | /// List tools together with their MCP argument schemas. |
| 135 | /// |
| 136 | /// The default keeps existing external client implementations source |
| 137 | /// compatible: clients that predate schema forwarding still advertise an |
| 138 | /// empty object schema, while transports that receive `inputSchema` can |
| 139 | /// override this method and preserve it. |
| 140 | fn list_tools_with_input_schemas(&self) -> Result<Vec<(McpToolDescriptor, Value)>> { |
| 141 | Ok(self |
| 142 | .list_tools()? |
| 143 | .into_iter() |
| 144 | .map(|tool| (tool, default_tool_input_schema())) |
| 145 | .collect()) |
| 146 | } |
| 147 | /// Invoke a tool by name with the given arguments. |
| 148 | fn call_tool(&self, tool_name: &str, arguments: Value) -> Result<Value>; |
| 149 | /// List all resources provided by this server. |
| 150 | fn list_resources(&self) -> Result<Vec<McpResourceDescriptor>>; |
| 151 | /// List resources together with their standard MCP metadata. |
| 152 | /// |
| 153 | /// `McpResourceDescriptor` predates the required MCP `name` field and the |
| 154 | /// optional `mimeType` field. Keeping those fields in additive metadata |
| 155 | /// preserves source compatibility for external trait implementations, |
| 156 | /// while transports that receive the standard fields can override this |
| 157 | /// method and forward them without loss. |
| 158 | fn list_resources_with_metadata(&self) -> Result<Vec<(McpResourceDescriptor, Value)>> { |
| 159 | Ok(self |
| 160 | .list_resources()? |
| 161 | .into_iter() |
| 162 | .map(|resource| { |
| 163 | let metadata = default_resource_metadata(&resource); |
| 164 | (resource, metadata) |
| 165 | }) |
| 166 | .collect()) |
| 167 | } |
| 168 | /// Read a resource by URI. |
| 169 | fn read_resource(&self, uri: &str) -> Result<Value>; |
| 170 | } |
| 171 | |
| 172 | /// A simple in-memory MCP client for tests and embedding callers. |
| 173 | /// |
| 174 | /// This is **not** wired into `codewhale mcp-server`: that path spawns |
| 175 | /// [`ChildProcessMcpClient`] and reports a typed error when the configured |
| 176 | /// command cannot be run. Serving canned values there made a broken |
| 177 | /// integration look identical to a working one (#4727). |
| 178 | #[derive(Debug, Default)] |
| 179 | pub struct InMemoryMcpClient { |
| 180 | tools: HashMap<String, Value>, |
| 181 | resources: HashMap<String, Value>, |
| 182 | } |
| 183 | |
| 184 | impl InMemoryMcpClient { |
| 185 | /// Register a tool with a fixed response value. |
| 186 | pub fn with_tool(mut self, name: &str, sample_result: Value) -> Self { |
| 187 | self.tools.insert(name.to_string(), sample_result); |
| 188 | self |
| 189 | } |
| 190 | |
| 191 | /// Register a resource with a fixed data value. |
| 192 | pub fn with_resource(mut self, uri: &str, data: Value) -> Self { |
| 193 | self.resources.insert(uri.to_string(), data); |
| 194 | self |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | impl McpManagedClient for InMemoryMcpClient { |
| 199 | fn list_tools(&self) -> Result<Vec<McpToolDescriptor>> { |
| 200 | Ok(self |
| 201 | .tools |
| 202 | .keys() |
| 203 | .map(|name| McpToolDescriptor { |
| 204 | server_name: "in-memory".to_string(), |
| 205 | tool_name: name.clone(), |
| 206 | qualified_name: name.clone(), |
| 207 | description: None, |
| 208 | }) |
| 209 | .collect()) |
| 210 | } |
| 211 | |
| 212 | fn call_tool(&self, tool_name: &str, _arguments: Value) -> Result<Value> { |
| 213 | self.tools |
| 214 | .get(tool_name) |
| 215 | .cloned() |
| 216 | .with_context(|| format!("tool '{tool_name}' not found")) |
| 217 | } |
| 218 | |
| 219 | fn list_resources(&self) -> Result<Vec<McpResourceDescriptor>> { |
| 220 | Ok(self |
| 221 | .resources |
| 222 | .keys() |
| 223 | .map(|uri| McpResourceDescriptor { |
| 224 | server_name: "in-memory".to_string(), |
| 225 | uri: uri.clone(), |
| 226 | description: None, |
| 227 | }) |
| 228 | .collect()) |
| 229 | } |
| 230 | |
| 231 | fn read_resource(&self, uri: &str) -> Result<Value> { |
| 232 | self.resources |
| 233 | .get(uri) |
| 234 | .cloned() |
| 235 | .with_context(|| format!("resource '{uri}' not found")) |
| 236 | } |
| 237 | } |
| 238 | |
| 239 | /// Manages multiple MCP server connections and their tool/resource registrations. |
| 240 | #[derive(Default)] |
| 241 | pub struct McpManager { |
| 242 | configs: HashMap<String, (McpServerConfig, ToolFilter)>, |
| 243 | clients: HashMap<String, Box<dyn McpManagedClient>>, |
| 244 | } |
| 245 | |
| 246 | impl McpManager { |
| 247 | /// Register an MCP server with its config, tool filter, and client implementation. |
| 248 | /// |
| 249 | /// Fails when the server's name collides with an already-registered server |
| 250 | /// after `sanitize_component` folding. Qualified tool names are built |
| 251 | /// from the sanitized name, so `my-server`, `my_server`, and `My.Server` |
| 252 | /// all produce `mcp__my_server__*`: registering two of them would let |
| 253 | /// either server answer a qualified name meant for the other. Re-registering |
| 254 | /// the same name replaces it, which is how restart works. |
| 255 | pub fn register_server( |
| 256 | &mut self, |
| 257 | config: McpServerConfig, |
| 258 | filter: ToolFilter, |
| 259 | client: Box<dyn McpManagedClient>, |
| 260 | ) -> Result<()> { |
| 261 | if let Some(existing) = self.colliding_server_name(&config.name) { |
| 262 | bail!( |
| 263 | "MCP server '{}' collides with already-registered server '{existing}': \ |
| 264 | both qualify tools as 'mcp__{}__*'", |
| 265 | config.name, |
| 266 | sanitize_component(&config.name) |
| 267 | ); |
| 268 | } |
| 269 | self.clients.insert(config.name.clone(), client); |
| 270 | self.configs.insert(config.name.clone(), (config, filter)); |
| 271 | Ok(()) |
| 272 | } |
| 273 | |
| 274 | /// Returns a registered server whose sanitized name matches `name`'s but |
| 275 | /// which is not `name` itself. |
| 276 | fn colliding_server_name(&self, name: &str) -> Option<&str> { |
| 277 | let sanitized = sanitize_component(name); |
| 278 | self.configs |
| 279 | .keys() |
| 280 | .find(|existing| existing.as_str() != name && sanitize_component(existing) == sanitized) |
| 281 | .map(String::as_str) |
| 282 | } |
| 283 | |
| 284 | /// Resolve a sanitized tool segment from a qualified name back to the |
| 285 | /// server's original tool name. |
| 286 | /// |
| 287 | /// `qualify_tool_name` folds `-`, `.`, and case into `_`, so the segment |
| 288 | /// carried by `mcp__server__segment` is not necessarily the name the |
| 289 | /// server expects. Exactly one advertised, allowed match resolves to its |
| 290 | /// original name. Multiple matches fail closed; no match preserves the |
| 291 | /// direct-call behavior for clients whose catalog is not exhaustive. |
| 292 | fn resolve_original_tool_name( |
| 293 | &self, |
| 294 | server_name: &str, |
| 295 | tool_segment: &str, |
| 296 | qualified_tool_name: &str, |
| 297 | ) -> Result<String> { |
| 298 | let Some(client) = self.clients.get(server_name) else { |
| 299 | return Ok(tool_segment.to_string()); |
| 300 | }; |
| 301 | let Ok(tools) = client.list_tools() else { |
| 302 | return Ok(tool_segment.to_string()); |
| 303 | }; |
| 304 | let filter = self.configs.get(server_name).map(|(_, filter)| filter); |
| 305 | let mut matches = tools.iter().filter(|tool| { |
| 306 | filter.is_none_or(|filter| allowed_by_filter(&tool.tool_name, filter)) |
| 307 | && qualify_tool_name(server_name, &tool.tool_name) == qualified_tool_name |
| 308 | }); |
| 309 | match (matches.next(), matches.next()) { |
| 310 | (Some(tool), None) => Ok(tool.tool_name.clone()), |
| 311 | (None, _) => Ok(tool_segment.to_string()), |
| 312 | (Some(_), Some(_)) => bail!( |
| 313 | "qualified MCP tool name '{qualified_tool_name}' is ambiguous within server '{server_name}'" |
| 314 | ), |
| 315 | } |
| 316 | } |
| 317 | |
| 318 | /// Start all registered servers, emitting status updates via the callback. |
| 319 | /// |
| 320 | /// Returns a summary of which servers are ready, failed, or cancelled. |
| 321 | pub fn start_all<F>(&self, mut emit: F) -> McpStartupCompleteEvent |
| 322 | where |
| 323 | F: FnMut(McpStartupUpdateEvent), |
| 324 | { |
| 325 | let mut ready = Vec::new(); |
| 326 | let mut failed = Vec::new(); |
| 327 | let mut cancelled = Vec::new(); |
| 328 | for (server_name, (cfg, _)) in &self.configs { |
| 329 | if !cfg.enabled { |
| 330 | emit(McpStartupUpdateEvent { |
| 331 | server_name: server_name.clone(), |
| 332 | status: McpStartupStatus::Cancelled, |
| 333 | }); |
| 334 | cancelled.push(server_name.clone()); |
| 335 | continue; |
| 336 | } |
| 337 | emit(McpStartupUpdateEvent { |
| 338 | server_name: server_name.clone(), |
| 339 | status: McpStartupStatus::Starting, |
| 340 | }); |
| 341 | if self.clients.contains_key(server_name) { |
| 342 | emit(McpStartupUpdateEvent { |
| 343 | server_name: server_name.clone(), |
| 344 | status: McpStartupStatus::Ready, |
| 345 | }); |
| 346 | ready.push(server_name.clone()); |
| 347 | } else { |
| 348 | let error = "client not registered".to_string(); |
| 349 | emit(McpStartupUpdateEvent { |
| 350 | server_name: server_name.clone(), |
| 351 | status: McpStartupStatus::Failed { |
| 352 | error: error.clone(), |
| 353 | }, |
| 354 | }); |
| 355 | failed.push(McpStartupFailure { |
| 356 | server_name: server_name.clone(), |
| 357 | error, |
| 358 | }); |
| 359 | } |
| 360 | } |
| 361 | McpStartupCompleteEvent { |
| 362 | ready, |
| 363 | failed, |
| 364 | cancelled, |
| 365 | } |
| 366 | } |
| 367 | |
| 368 | /// Stop a running server by removing its client. |
| 369 | pub fn stop_server(&mut self, server_name: &str) -> Result<()> { |
| 370 | self.clients |
| 371 | .remove(server_name) |
| 372 | .with_context(|| format!("server '{server_name}' is not running"))?; |
| 373 | Ok(()) |
| 374 | } |
| 375 | |
| 376 | /// Remove a server entirely (config and client). |
| 377 | pub fn unregister_server(&mut self, server_name: &str) -> Result<()> { |
| 378 | let had_config = self.configs.remove(server_name).is_some(); |
| 379 | self.clients.remove(server_name); |
| 380 | if !had_config { |
| 381 | bail!("server '{server_name}' is not registered"); |
| 382 | } |
| 383 | Ok(()) |
| 384 | } |
| 385 | |
| 386 | /// List all tools from all running servers, applying tool filters. |
| 387 | pub fn list_tools(&self) -> Result<Vec<McpToolDescriptor>> { |
| 388 | Ok(self |
| 389 | .list_tools_with_input_schemas()? |
| 390 | .into_iter() |
| 391 | .map(|(tool, _)| tool) |
| 392 | .collect()) |
| 393 | } |
| 394 | |
| 395 | fn list_tools_with_input_schemas(&self) -> Result<Vec<(McpToolDescriptor, Value)>> { |
| 396 | let mut out = Vec::new(); |
| 397 | let mut qualified_origins: HashMap<String, (String, String)> = HashMap::new(); |
| 398 | for (server_name, (_, filter)) in &self.configs { |
| 399 | let Some(client) = self.clients.get(server_name) else { |
| 400 | continue; |
| 401 | }; |
| 402 | let tools = client.list_tools_with_input_schemas()?; |
| 403 | for (tool, input_schema) in tools { |
| 404 | if !allowed_by_filter(&tool.tool_name, filter) { |
| 405 | continue; |
| 406 | } |
| 407 | let qualified_name = qualify_tool_name(server_name, &tool.tool_name); |
| 408 | if let Some((prior_server, prior_tool)) = qualified_origins.get(&qualified_name) |
| 409 | && (prior_server != server_name || prior_tool != &tool.tool_name) |
| 410 | { |
| 411 | let mut origins = [ |
| 412 | format!("{prior_server}:{prior_tool}"), |
| 413 | format!("{server_name}:{}", tool.tool_name), |
| 414 | ]; |
| 415 | origins.sort(); |
| 416 | bail!( |
| 417 | "qualified MCP tool name '{qualified_name}' is ambiguous between {}", |
| 418 | origins.join(" and ") |
| 419 | ); |
| 420 | } |
| 421 | qualified_origins.insert( |
| 422 | qualified_name.clone(), |
| 423 | (server_name.clone(), tool.tool_name.clone()), |
| 424 | ); |
| 425 | out.push(( |
| 426 | McpToolDescriptor { |
| 427 | server_name: server_name.clone(), |
| 428 | tool_name: tool.tool_name, |
| 429 | qualified_name, |
| 430 | description: tool.description, |
| 431 | }, |
| 432 | input_schema, |
| 433 | )); |
| 434 | } |
| 435 | } |
| 436 | Ok(out) |
| 437 | } |
| 438 | |
| 439 | /// Call a tool on a specific server by name. |
| 440 | /// |
| 441 | /// The server's [`ToolFilter`] is enforced on invocation, not just at |
| 442 | /// listing time: a denied (or not-allowed) tool cannot be executed by |
| 443 | /// addressing the server directly, whether by bare or qualified name. |
| 444 | pub fn call_tool(&self, server_name: &str, tool_name: &str, arguments: Value) -> Result<Value> { |
| 445 | let client = self |
| 446 | .clients |
| 447 | .get(server_name) |
| 448 | .with_context(|| format!("MCP server '{server_name}' not available"))?; |
| 449 | if let Some((_, filter)) = self.configs.get(server_name) |
| 450 | && !allowed_by_filter(tool_name, filter) |
| 451 | { |
| 452 | bail!("tool '{tool_name}' on MCP server '{server_name}' is blocked by the tool filter"); |
| 453 | } |
| 454 | client.call_tool(tool_name, arguments) |
| 455 | } |
| 456 | |
| 457 | /// Call a tool using its fully qualified name (e.g., `mcp__server__tool`). |
| 458 | pub fn call_qualified_tool( |
| 459 | &self, |
| 460 | qualified_tool_name: &str, |
| 461 | arguments: Value, |
| 462 | ) -> Result<Value> { |
| 463 | let parsed = parse_qualified_tool_name(qualified_tool_name) |
| 464 | .with_context(|| format!("invalid qualified MCP tool name: {qualified_tool_name}")); |
| 465 | |
| 466 | // An exact registration is the answer. Whatever the tool returns — |
| 467 | // including an error — is returned as-is: falling through to the scan |
| 468 | // below on a *call* failure would re-execute the same tool, and for a |
| 469 | // file write, a commit, or a paid API call that second invocation is a |
| 470 | // second real side effect. Only a failed *lookup* falls through. |
| 471 | // |
| 472 | // The parsed tool segment is the *sanitized* name (qualify_tool_name |
| 473 | // folds `-`, `.`, and case into `_`), so resolve it back to the |
| 474 | // server's original tool name before dispatching — otherwise tools |
| 475 | // like `my-tool` are un-callable through their advertised qualified |
| 476 | // name `mcp__server__my_tool`. |
| 477 | if let Ok((server_name, tool_name)) = &parsed |
| 478 | && self.clients.contains_key(server_name) |
| 479 | { |
| 480 | let resolved = |
| 481 | self.resolve_original_tool_name(server_name, tool_name, qualified_tool_name)?; |
| 482 | return self.call_tool(server_name, &resolved, arguments); |
| 483 | } |
| 484 | |
| 485 | // No exact registration: resolve by scanning qualified names. Collect |
| 486 | // every match rather than returning the first, because `configs` is a |
| 487 | // HashMap — returning early would make the choice depend on iteration |
| 488 | // order when two servers' names collide after sanitizing. |
| 489 | let mut matches: Vec<(&String, String)> = Vec::new(); |
| 490 | for (server_name, (_, filter)) in &self.configs { |
| 491 | let Some(client) = self.clients.get(server_name) else { |
| 492 | continue; |
| 493 | }; |
| 494 | for tool in client.list_tools()? { |
| 495 | if !allowed_by_filter(&tool.tool_name, filter) { |
| 496 | continue; |
| 497 | } |
| 498 | if qualify_tool_name(server_name, &tool.tool_name) == qualified_tool_name { |
| 499 | matches.push((server_name, tool.tool_name)); |
| 500 | } |
| 501 | } |
| 502 | } |
| 503 | match matches.len() { |
| 504 | 0 => {} |
| 505 | 1 => { |
| 506 | let (server_name, tool_name) = &matches[0]; |
| 507 | let client = self |
| 508 | .clients |
| 509 | .get(*server_name) |
| 510 | .with_context(|| format!("MCP server '{server_name}' not available"))?; |
| 511 | return client.call_tool(tool_name, arguments); |
| 512 | } |
| 513 | _ => { |
| 514 | matches.sort(); |
| 515 | let servers: Vec<&str> = matches |
| 516 | .iter() |
| 517 | .map(|(server_name, _)| server_name.as_str()) |
| 518 | .collect(); |
| 519 | bail!( |
| 520 | "qualified MCP tool name '{qualified_tool_name}' is ambiguous across servers: \ |
| 521 | {}", |
| 522 | servers.join(", ") |
| 523 | ); |
| 524 | } |
| 525 | } |
| 526 | |
| 527 | let (server_name, tool_name) = parsed?; |
| 528 | self.call_tool(&server_name, &tool_name, arguments) |
| 529 | } |
| 530 | |
| 531 | /// List all resources from all running servers. |
| 532 | pub fn list_resources(&self) -> Result<Vec<McpResourceDescriptor>> { |
| 533 | Ok(self |
| 534 | .list_resources_with_metadata()? |
| 535 | .into_iter() |
| 536 | .map(|(resource, _)| resource) |
| 537 | .collect()) |
| 538 | } |
| 539 | |
| 540 | fn list_resources_with_metadata(&self) -> Result<Vec<(McpResourceDescriptor, Value)>> { |
| 541 | let mut out = Vec::new(); |
| 542 | for server_name in self.configs.keys() { |
| 543 | let Some(client) = self.clients.get(server_name) else { |
| 544 | continue; |
| 545 | }; |
| 546 | for (mut resource, metadata) in client.list_resources_with_metadata()? { |
| 547 | resource.server_name = server_name.clone(); |
| 548 | out.push((resource, metadata)); |
| 549 | } |
| 550 | } |
| 551 | Ok(out) |
| 552 | } |
| 553 | |
| 554 | /// Read a resource from a specific server. |
| 555 | pub fn read_resource(&self, server_name: &str, uri: &str) -> Result<Value> { |
| 556 | let client = self |
| 557 | .clients |
| 558 | .get(server_name) |
| 559 | .with_context(|| format!("MCP server '{server_name}' not available"))?; |
| 560 | client.read_resource(uri) |
| 561 | } |
| 562 | |
| 563 | /// Resolve a standard URI-only resource read to exactly one child server. |
| 564 | /// |
| 565 | /// Older Codewhale clients supplied a non-standard `server` parameter (or |
| 566 | /// encoded it as the authority in an `mcp://server/...` URI). Standard MCP |
| 567 | /// clients send only the URI, so discover its owner from resources/list. |
| 568 | /// Never pick the first HashMap entry when more than one server advertises |
| 569 | /// the same URI. |
| 570 | fn read_resource_by_uri(&self, uri: &str) -> Result<Value> { |
| 571 | let mut matches = Vec::new(); |
| 572 | for server_name in self.configs.keys() { |
| 573 | let Some(client) = self.clients.get(server_name) else { |
| 574 | continue; |
| 575 | }; |
| 576 | if client |
| 577 | .list_resources()? |
| 578 | .iter() |
| 579 | .any(|resource| resource.uri == uri) |
| 580 | { |
| 581 | matches.push(server_name.clone()); |
| 582 | } |
| 583 | } |
| 584 | |
| 585 | match matches.len() { |
| 586 | 1 => self.read_resource(&matches[0], uri), |
| 587 | 0 => { |
| 588 | // Preserve the pre-standard URI convention for clients whose |
| 589 | // server does not implement resources/list. |
| 590 | if let Some(server_name) = parse_server_from_uri(uri) |
| 591 | && self.clients.contains_key(&server_name) |
| 592 | { |
| 593 | return self.read_resource(&server_name, uri); |
| 594 | } |
| 595 | bail!("resource URI '{uri}' was not advertised by any running MCP server") |
| 596 | } |
| 597 | _ => { |
| 598 | matches.sort(); |
| 599 | bail!( |
| 600 | "resource URI '{uri}' is ambiguous across MCP servers: {}; pass the legacy \ |
| 601 | server parameter to disambiguate", |
| 602 | matches.join(", ") |
| 603 | ) |
| 604 | } |
| 605 | } |
| 606 | } |
| 607 | |
| 608 | /// Generate sandbox state update notices for all registered servers. |
| 609 | pub fn update_sandbox_state(&self, sandbox_mode: &str, cwd: &str) -> Result<Vec<Value>> { |
| 610 | let mut notices = Vec::new(); |
| 611 | for server_name in self.configs.keys() { |
| 612 | notices.push(json!({ |
| 613 | "server_name": server_name, |
| 614 | "method": "codex/sandbox-state/update", |
| 615 | "params": { |
| 616 | "sandbox_mode": sandbox_mode, |
| 617 | "cwd": cwd |
| 618 | } |
| 619 | })); |
| 620 | } |
| 621 | Ok(notices) |
| 622 | } |
| 623 | } |
| 624 | |
| 625 | fn default_true() -> bool { |
| 626 | true |
| 627 | } |
| 628 | |
| 629 | fn default_tool_input_schema() -> Value { |
| 630 | json!({"type": "object", "properties": {}}) |
| 631 | } |
| 632 | |
| 633 | fn default_resource_metadata(resource: &McpResourceDescriptor) -> Value { |
| 634 | // The URI is a stable, non-empty fallback name for clients implementing |
| 635 | // the older descriptor-only trait surface. |
| 636 | json!({"name": resource.uri}) |
| 637 | } |
| 638 | |
| 639 | fn allowed_by_filter(name: &str, filter: &ToolFilter) -> bool { |
| 640 | if filter.deny.iter().any(|pattern| pattern == name) { |
| 641 | return false; |
| 642 | } |
| 643 | if filter.allow.is_empty() { |
| 644 | return true; |
| 645 | } |
| 646 | filter.allow.iter().any(|pattern| pattern == name) |
| 647 | } |
| 648 | |
| 649 | fn sanitize_component(value: &str) -> String { |
| 650 | value |
| 651 | .chars() |
| 652 | .map(|ch| { |
| 653 | if ch.is_ascii_alphanumeric() || ch == '_' { |
| 654 | ch.to_ascii_lowercase() |
| 655 | } else { |
| 656 | '_' |
| 657 | } |
| 658 | }) |
| 659 | .collect() |
| 660 | } |
| 661 | |
| 662 | fn qualify_tool_name(server: &str, tool: &str) -> String { |
| 663 | let server = sanitize_component(server); |
| 664 | let tool = sanitize_component(tool); |
| 665 | let mut name = format!("mcp__{server}__{tool}"); |
| 666 | if name.len() > 64 { |
| 667 | let mut hasher = DefaultHasher::new(); |
| 668 | name.hash(&mut hasher); |
| 669 | let hash = format!("{:x}", hasher.finish()); |
| 670 | let suffix = format!("_{}", &hash[..12]); |
| 671 | let component_budget = 64 - "mcp__".len() - "__".len() - suffix.len(); |
| 672 | let mut server_len = server.len().min(component_budget / 2); |
| 673 | let mut tool_len = tool.len().min(component_budget - server_len); |
| 674 | let remaining = component_budget - server_len - tool_len; |
| 675 | if remaining > 0 { |
| 676 | let server_extra = (server.len() - server_len).min(remaining); |
| 677 | server_len += server_extra; |
| 678 | tool_len += (tool.len() - tool_len).min(remaining - server_extra); |
| 679 | } |
| 680 | name = format!( |
| 681 | "mcp__{}__{}{}", |
| 682 | &server[..server_len], |
| 683 | &tool[..tool_len], |
| 684 | suffix |
| 685 | ); |
| 686 | } |
| 687 | name |
| 688 | } |
| 689 | |
| 690 | fn parse_qualified_tool_name(value: &str) -> Result<(String, String)> { |
| 691 | let Some(stripped) = value.strip_prefix("mcp__") else { |
| 692 | bail!("missing mcp__ prefix"); |
| 693 | }; |
| 694 | let mut split = stripped.splitn(2, "__"); |
| 695 | let server = split |
| 696 | .next() |
| 697 | .filter(|s| !s.is_empty()) |
| 698 | .map(ToOwned::to_owned) |
| 699 | .context("missing server segment")?; |
| 700 | let tool = split |
| 701 | .next() |
| 702 | .filter(|s| !s.is_empty()) |
| 703 | .map(ToOwned::to_owned) |
| 704 | .context("missing tool segment")?; |
| 705 | Ok((server, tool)) |
| 706 | } |
| 707 | |
| 708 | #[derive(Debug, Deserialize)] |
| 709 | struct JsonRpcRequest { |
| 710 | jsonrpc: String, |
| 711 | #[serde(default)] |
| 712 | id: JsonRpcRequestId, |
| 713 | method: String, |
| 714 | #[serde(default)] |
| 715 | params: Value, |
| 716 | } |
| 717 | |
| 718 | /// JSON-RPC defines a notification by an *absent* id. An explicit `null` id is |
| 719 | /// discouraged but still present and must receive a response carrying null; |
| 720 | /// `Option<Value>` cannot preserve that distinction during deserialization. |
| 721 | #[derive(Debug, Clone, Default)] |
| 722 | enum JsonRpcRequestId { |
| 723 | #[default] |
| 724 | Missing, |
| 725 | Present(Value), |
| 726 | } |
| 727 | |
| 728 | impl<'de> Deserialize<'de> for JsonRpcRequestId { |
| 729 | fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error> |
| 730 | where |
| 731 | D: serde::Deserializer<'de>, |
| 732 | { |
| 733 | let value = Value::deserialize(deserializer)?; |
| 734 | if !(value.is_null() || value.is_string() || value.is_number()) { |
| 735 | return Err(serde::de::Error::custom( |
| 736 | "JSON-RPC id must be a string, number, or null", |
| 737 | )); |
| 738 | } |
| 739 | Ok(Self::Present(value)) |
| 740 | } |
| 741 | } |
| 742 | |
| 743 | impl JsonRpcRequestId { |
| 744 | fn should_respond(&self) -> bool { |
| 745 | matches!(self, Self::Present(_)) |
| 746 | } |
| 747 | |
| 748 | fn response_id(&self) -> Option<Value> { |
| 749 | match self { |
| 750 | Self::Missing => None, |
| 751 | Self::Present(id) => Some(id.clone()), |
| 752 | } |
| 753 | } |
| 754 | } |
| 755 | |
| 756 | #[derive(Debug)] |
| 757 | struct JsonRpcError { |
| 758 | code: i64, |
| 759 | message: String, |
| 760 | data: Option<Value>, |
| 761 | } |
| 762 | |
| 763 | #[derive(Debug, Deserialize)] |
| 764 | struct ToolsListParams { |
| 765 | #[serde(default)] |
| 766 | server: Option<String>, |
| 767 | } |
| 768 | |
| 769 | #[derive(Debug, Deserialize)] |
| 770 | struct McpImplementationInfo { |
| 771 | name: String, |
| 772 | version: String, |
| 773 | } |
| 774 | |
| 775 | #[derive(Debug, Deserialize)] |
| 776 | struct InitializeParams { |
| 777 | #[serde(rename = "protocolVersion")] |
| 778 | protocol_version: String, |
| 779 | #[serde(rename = "clientInfo")] |
| 780 | client_info: McpImplementationInfo, |
| 781 | capabilities: serde_json::Map<String, Value>, |
| 782 | } |
| 783 | |
| 784 | #[derive(Debug, Deserialize)] |
| 785 | struct ToolsCallParams { |
| 786 | #[serde(default)] |
| 787 | name: Option<String>, |
| 788 | #[serde(default)] |
| 789 | tool: Option<String>, |
| 790 | #[serde(default)] |
| 791 | server: Option<String>, |
| 792 | #[serde(default = "default_tool_arguments")] |
| 793 | arguments: Value, |
| 794 | } |
| 795 | |
| 796 | fn default_tool_arguments() -> Value { |
| 797 | json!({}) |
| 798 | } |
| 799 | |
| 800 | #[derive(Debug, Deserialize)] |
| 801 | struct ResourcesListParams { |
| 802 | #[serde(default)] |
| 803 | server: Option<String>, |
| 804 | } |
| 805 | |
| 806 | #[derive(Debug, Deserialize)] |
| 807 | struct ResourcesReadParams { |
| 808 | #[serde(default)] |
| 809 | server: Option<String>, |
| 810 | uri: String, |
| 811 | } |
| 812 | |
| 813 | #[derive(Debug, Deserialize)] |
| 814 | struct ServerRegisterParams { |
| 815 | server: McpServerConfig, |
| 816 | #[serde(default)] |
| 817 | filter: ToolFilter, |
| 818 | #[serde(default = "default_true")] |
| 819 | start: bool, |
| 820 | } |
| 821 | |
| 822 | #[derive(Debug, Deserialize)] |
| 823 | struct ServerNameParams { |
| 824 | name: String, |
| 825 | } |
| 826 | |
| 827 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 828 | enum McpSessionPhase { |
| 829 | Uninitialized, |
| 830 | InitializeResponded, |
| 831 | Ready, |
| 832 | } |
| 833 | |
| 834 | struct StdioMcpState { |
| 835 | manager: McpManager, |
| 836 | definitions: HashMap<String, McpServerDefinition>, |
| 837 | running: HashMap<String, bool>, |
| 838 | /// Why a defined server is not running, surfaced in every lifecycle |
| 839 | /// snapshot so a failed spawn cannot be mistaken for a healthy server. |
| 840 | errors: HashMap<String, String>, |
| 841 | lifecycle_state: String, |
| 842 | session_phase: McpSessionPhase, |
| 843 | } |
| 844 | |
| 845 | impl StdioMcpState { |
| 846 | /// Spawn `definition`'s configured command and register the resulting |
| 847 | /// connection, recording the failure reason when the server cannot be |
| 848 | /// brought up. |
| 849 | /// |
| 850 | /// This is the only way a server enters `manager`, and it has no stub |
| 851 | /// branch: there is no configuration under which a registered server |
| 852 | /// answers from anything but its own process. |
| 853 | fn start_definition(&mut self, definition: &McpServerDefinition) -> Result<()> { |
| 854 | let name = definition.config.name.clone(); |
| 855 | let outcome = ChildProcessMcpClient::spawn(&definition.config).and_then(|client| { |
| 856 | self.manager.register_server( |
| 857 | definition.config.clone(), |
| 858 | definition.filter.clone(), |
| 859 | Box::new(client), |
| 860 | ) |
| 861 | }); |
| 862 | match outcome { |
| 863 | Ok(()) => { |
| 864 | self.errors.remove(&name); |
| 865 | self.running.insert(name, true); |
| 866 | Ok(()) |
| 867 | } |
| 868 | Err(err) => { |
| 869 | let message = format!("{err:#}"); |
| 870 | self.errors.insert(name.clone(), message.clone()); |
| 871 | self.running.insert(name, false); |
| 872 | Err(err) |
| 873 | } |
| 874 | } |
| 875 | } |
| 876 | } |
| 877 | |
| 878 | /// Run an MCP stdio server that reads JSON-RPC requests from stdin and writes responses to stdout. |
| 879 | /// |
| 880 | /// Returns the final server definitions after the session ends (useful for persisting |
| 881 | /// runtime changes like server registrations). |
| 882 | pub fn run_stdio_server( |
| 883 | initial_definitions: Vec<McpServerDefinition>, |
| 884 | ) -> Result<Vec<McpServerDefinition>> { |
| 885 | use std::io::{self, Write}; |
| 886 | |
| 887 | let stdin = io::stdin(); |
| 888 | let mut stdout = io::stdout(); |
| 889 | let mut stderr = io::stderr(); |
| 890 | let mut state = build_stdio_state(initial_definitions); |
| 891 | let mut input = stdin.lock(); |
| 892 | |
| 893 | loop { |
| 894 | let line = |
| 895 | match stdio_client::read_bounded_line(&mut input, stdio_client::MAX_JSONRPC_LINE_BYTES) |
| 896 | { |
| 897 | Ok(Some(line)) => line, |
| 898 | Ok(None) => break, |
| 899 | Err(err) => { |
| 900 | let response = jsonrpc_error( |
| 901 | None, |
| 902 | JsonRpcError::parse_error(format!("invalid JSON-RPC frame: {err}")), |
| 903 | ); |
| 904 | writeln!(stdout, "{response}")?; |
| 905 | stdout.flush()?; |
| 906 | bail!("failed to read bounded stdio JSON-RPC frame: {err}"); |
| 907 | } |
| 908 | }; |
| 909 | if line.trim().is_empty() { |
| 910 | continue; |
| 911 | } |
| 912 | |
| 913 | let value: Value = match serde_json::from_str(&line) { |
| 914 | Ok(value) => value, |
| 915 | Err(err) => { |
| 916 | let msg = jsonrpc_error( |
| 917 | None, |
| 918 | JsonRpcError::parse_error(format!("invalid json: {err}")), |
| 919 | ); |
| 920 | writeln!(stdout, "{msg}")?; |
| 921 | stdout.flush()?; |
| 922 | continue; |
| 923 | } |
| 924 | }; |
| 925 | let request: JsonRpcRequest = match serde_json::from_value(value) { |
| 926 | Ok(request) => request, |
| 927 | Err(err) => { |
| 928 | let response = jsonrpc_error( |
| 929 | None, |
| 930 | JsonRpcError::invalid_request(format!("invalid JSON-RPC request: {err}")), |
| 931 | ); |
| 932 | writeln!(stdout, "{response}")?; |
| 933 | stdout.flush()?; |
| 934 | continue; |
| 935 | } |
| 936 | }; |
| 937 | let should_respond = request.id.should_respond(); |
| 938 | let response_id = request.id.response_id(); |
| 939 | |
| 940 | if request.jsonrpc != "2.0" { |
| 941 | let response = jsonrpc_error( |
| 942 | response_id, |
| 943 | JsonRpcError::invalid_request("jsonrpc version must be exactly 2.0"), |
| 944 | ); |
| 945 | writeln!(stdout, "{response}")?; |
| 946 | stdout.flush()?; |
| 947 | continue; |
| 948 | } |
| 949 | |
| 950 | if !should_respond { |
| 951 | // `initialize` requires a response carrying the negotiated |
| 952 | // protocol version. Treating an initialize notification as a real |
| 953 | // handshake would advance state even though the client could not |
| 954 | // observe that negotiation. |
| 955 | if request.method == "initialize" { |
| 956 | continue; |
| 957 | } |
| 958 | match dispatch_stdio_request(&mut state, &request.method, request.params) { |
| 959 | Ok((_, should_exit)) if should_exit => break, |
| 960 | Ok(_) | Err(_) => {} |
| 961 | } |
| 962 | continue; |
| 963 | } |
| 964 | |
| 965 | let response = match dispatch_stdio_request(&mut state, &request.method, request.params) { |
| 966 | Ok((result, should_exit)) => { |
| 967 | let payload = jsonrpc_result(response_id, result); |
| 968 | writeln!(stdout, "{payload}")?; |
| 969 | stdout.flush()?; |
| 970 | if should_exit { |
| 971 | break; |
| 972 | } |
| 973 | continue; |
| 974 | } |
| 975 | Err(err) => jsonrpc_error(response_id, err), |
| 976 | }; |
| 977 | |
| 978 | writeln!(stdout, "{response}")?; |
| 979 | stdout.flush()?; |
| 980 | } |
| 981 | |
| 982 | state.lifecycle_state = "stopped".to_string(); |
| 983 | let _ = writeln!(stderr, "codewhale mcp-server: stdio server exited"); |
| 984 | let mut definitions: Vec<McpServerDefinition> = state.definitions.into_values().collect(); |
| 985 | definitions.sort_by(|a, b| a.config.name.cmp(&b.config.name)); |
| 986 | Ok(definitions) |
| 987 | } |
| 988 | |
| 989 | fn build_stdio_state(initial_definitions: Vec<McpServerDefinition>) -> StdioMcpState { |
| 990 | let mut state = StdioMcpState { |
| 991 | manager: McpManager::default(), |
| 992 | definitions: HashMap::new(), |
| 993 | running: HashMap::new(), |
| 994 | errors: HashMap::new(), |
| 995 | lifecycle_state: "running".to_string(), |
| 996 | session_phase: McpSessionPhase::Uninitialized, |
| 997 | }; |
| 998 | |
| 999 | for definition in initial_definitions { |
| 1000 | let name = definition.config.name.clone(); |
| 1001 | state.definitions.insert(name.clone(), definition.clone()); |
| 1002 | if !definition.config.enabled { |
| 1003 | state.running.insert(name, false); |
| 1004 | continue; |
| 1005 | } |
| 1006 | // A server that cannot be spawned stays stopped and says so on stderr. |
| 1007 | // stdout is the JSON-RPC channel, so the warning goes to stderr where |
| 1008 | // it will not corrupt the protocol stream but is still visible to the |
| 1009 | // operator; `lifecycle` carries the same text for programmatic clients. |
| 1010 | if let Err(err) = state.start_definition(&definition) { |
| 1011 | tracing::warn!("MCP server '{name}' is not available: {err:#}"); |
| 1012 | eprintln!("codewhale mcp-server: server '{name}' is not available: {err:#}"); |
| 1013 | } |
| 1014 | } |
| 1015 | |
| 1016 | state |
| 1017 | } |
| 1018 | |
| 1019 | fn default_rpc_methods() -> Vec<&'static str> { |
| 1020 | vec![ |
| 1021 | "initialize", |
| 1022 | "notifications/initialized", |
| 1023 | "ping", |
| 1024 | "healthz", |
| 1025 | "capabilities", |
| 1026 | "tools/list", |
| 1027 | "tools/call", |
| 1028 | "resources/list", |
| 1029 | "resources/read", |
| 1030 | "server/list", |
| 1031 | "server/register", |
| 1032 | "server/start", |
| 1033 | "server/stop", |
| 1034 | "server/unregister", |
| 1035 | "shutdown", |
| 1036 | ] |
| 1037 | } |
| 1038 | |
| 1039 | /// Latest dated MCP protocol revision this server implements. Codewhale's |
| 1040 | /// MCP clients advertise the same revision at `initialize`. |
| 1041 | pub(crate) const MCP_PROTOCOL_VERSION: &str = "2025-06-18"; |
| 1042 | /// Dated MCP revisions accepted during protocol negotiation, newest first. |
| 1043 | /// Servers answering an older supported revision get it echoed back. |
| 1044 | pub(crate) const MCP_SUPPORTED_PROTOCOL_VERSIONS: &[&str] = |
| 1045 | &[MCP_PROTOCOL_VERSION, "2025-03-26", "2024-11-05"]; |
| 1046 | const MCP_SERVER_NAME: &str = "codewhale-mcp-server"; |
| 1047 | |
| 1048 | fn initialize_response(state: &StdioMcpState, protocol_version: &str) -> Value { |
| 1049 | json!({ |
| 1050 | // Standard MCP initialize result. Keep the management metadata below |
| 1051 | // as additive compatibility fields for existing Codewhale clients. |
| 1052 | "protocolVersion": protocol_version, |
| 1053 | "capabilities": { |
| 1054 | "tools": {}, |
| 1055 | "resources": {} |
| 1056 | }, |
| 1057 | "serverInfo": { |
| 1058 | "name": MCP_SERVER_NAME, |
| 1059 | "version": env!("CARGO_PKG_VERSION") |
| 1060 | }, |
| 1061 | "server": MCP_SERVER_NAME, |
| 1062 | "transport": "stdio", |
| 1063 | "methods": default_rpc_methods(), |
| 1064 | "lifecycle": lifecycle_snapshot(state) |
| 1065 | }) |
| 1066 | } |
| 1067 | |
| 1068 | fn valid_mcp_annotations(value: &Value) -> bool { |
| 1069 | let Some(annotations) = value.as_object() else { |
| 1070 | return false; |
| 1071 | }; |
| 1072 | if let Some(audience) = annotations.get("audience") { |
| 1073 | let Some(audience) = audience.as_array() else { |
| 1074 | return false; |
| 1075 | }; |
| 1076 | if audience |
| 1077 | .iter() |
| 1078 | .any(|role| !matches!(role.as_str(), Some("user") | Some("assistant"))) |
| 1079 | { |
| 1080 | return false; |
| 1081 | } |
| 1082 | } |
| 1083 | if let Some(priority) = annotations.get("priority") { |
| 1084 | let Some(priority) = priority.as_f64() else { |
| 1085 | return false; |
| 1086 | }; |
| 1087 | if !priority.is_finite() || !(0.0..=1.0).contains(&priority) { |
| 1088 | return false; |
| 1089 | } |
| 1090 | } |
| 1091 | true |
| 1092 | } |
| 1093 | |
| 1094 | fn valid_optional_annotations(fields: &serde_json::Map<String, Value>) -> bool { |
| 1095 | fields.get("annotations").is_none_or(valid_mcp_annotations) |
| 1096 | } |
| 1097 | |
| 1098 | fn valid_resource_content(value: &Value) -> bool { |
| 1099 | let Some(content) = value.as_object() else { |
| 1100 | return false; |
| 1101 | }; |
| 1102 | let has_uri = content.get("uri").and_then(Value::as_str).is_some(); |
| 1103 | let has_payload = content.get("text").and_then(Value::as_str).is_some() |
| 1104 | || content.get("blob").and_then(Value::as_str).is_some(); |
| 1105 | let valid_mime_type = content.get("mimeType").is_none_or(Value::is_string); |
| 1106 | has_uri && has_payload && valid_mime_type |
| 1107 | } |
| 1108 | |
| 1109 | fn valid_tool_content(value: &Value) -> bool { |
| 1110 | let Some(content) = value.as_object() else { |
| 1111 | return false; |
| 1112 | }; |
| 1113 | if !valid_optional_annotations(content) { |
| 1114 | return false; |
| 1115 | } |
| 1116 | match content.get("type").and_then(Value::as_str) { |
| 1117 | Some("text") => content.get("text").is_some_and(Value::is_string), |
| 1118 | Some("image") => { |
| 1119 | content.get("data").is_some_and(Value::is_string) |
| 1120 | && content.get("mimeType").is_some_and(Value::is_string) |
| 1121 | } |
| 1122 | Some("resource") => content.get("resource").is_some_and(valid_resource_content), |
| 1123 | _ => false, |
| 1124 | } |
| 1125 | } |
| 1126 | |
| 1127 | fn valid_call_tool_result(fields: &serde_json::Map<String, Value>) -> bool { |
| 1128 | fields |
| 1129 | .get("content") |
| 1130 | .and_then(Value::as_array) |
| 1131 | .is_some_and(|content| content.iter().all(valid_tool_content)) |
| 1132 | && fields.get("isError").is_none_or(Value::is_boolean) |
| 1133 | && fields.get("_meta").is_none_or(Value::is_object) |
| 1134 | } |
| 1135 | |
| 1136 | fn looks_like_call_tool_result(fields: &serde_json::Map<String, Value>) -> bool { |
| 1137 | fields.contains_key("content") || fields.contains_key("isError") || fields.contains_key("_meta") |
| 1138 | } |
| 1139 | |
| 1140 | fn legacy_value_text(value: &Value) -> String { |
| 1141 | match value { |
| 1142 | Value::String(text) => text.clone(), |
| 1143 | value => value.to_string(), |
| 1144 | } |
| 1145 | } |
| 1146 | |
| 1147 | fn stdio_tool_descriptor((tool, input_schema): (McpToolDescriptor, Value)) -> Value { |
| 1148 | let McpToolDescriptor { |
| 1149 | server_name, |
| 1150 | tool_name, |
| 1151 | qualified_name, |
| 1152 | description, |
| 1153 | } = tool; |
| 1154 | let mut value = json!({ |
| 1155 | // Standard MCP fields. The qualified name is the only collision-safe |
| 1156 | // public name once several child servers are aggregated. |
| 1157 | "name": qualified_name.clone(), |
| 1158 | "inputSchema": input_schema, |
| 1159 | // Retain the pre-0.9.11 management fields for compatibility. |
| 1160 | "server_name": server_name, |
| 1161 | "tool_name": tool_name, |
| 1162 | "qualified_name": qualified_name, |
| 1163 | }); |
| 1164 | if let Some(description) = description { |
| 1165 | value["description"] = Value::String(description); |
| 1166 | } |
| 1167 | value |
| 1168 | } |
| 1169 | |
| 1170 | fn stdio_tool_call_result(result: Value) -> Result<Value> { |
| 1171 | let legacy_result = result.clone(); |
| 1172 | match result { |
| 1173 | Value::Object(mut fields) if looks_like_call_tool_result(&fields) => { |
| 1174 | if !valid_call_tool_result(&fields) { |
| 1175 | bail!("child returned a malformed MCP CallToolResult"); |
| 1176 | } |
| 1177 | // The child already returned a standard MCP CallToolResult. Expose |
| 1178 | // it directly, while retaining the old nested result for clients |
| 1179 | // that used the proxy before its MCP envelope was corrected. |
| 1180 | fields.insert("result".to_string(), legacy_result); |
| 1181 | Ok(Value::Object(fields)) |
| 1182 | } |
| 1183 | value => Ok(json!({ |
| 1184 | "content": [{"type": "text", "text": legacy_value_text(&value)}], |
| 1185 | "result": legacy_result |
| 1186 | })), |
| 1187 | } |
| 1188 | } |
| 1189 | |
| 1190 | fn stdio_resource_descriptor((resource, metadata): (McpResourceDescriptor, Value)) -> Value { |
| 1191 | let McpResourceDescriptor { |
| 1192 | server_name, |
| 1193 | uri, |
| 1194 | description, |
| 1195 | } = resource; |
| 1196 | let name = metadata |
| 1197 | .get("name") |
| 1198 | .and_then(Value::as_str) |
| 1199 | .filter(|name| !name.trim().is_empty()) |
| 1200 | .unwrap_or(&uri) |
| 1201 | .to_string(); |
| 1202 | let mut value = json!({ |
| 1203 | // Standard MCP Resource fields. |
| 1204 | "uri": uri, |
| 1205 | "name": name, |
| 1206 | // Retain the pre-0.9.11 server selector as additive metadata. |
| 1207 | "server_name": server_name, |
| 1208 | }); |
| 1209 | if let Some(description) = description { |
| 1210 | value["description"] = Value::String(description); |
| 1211 | } |
| 1212 | if let Some(mime_type) = metadata.get("mimeType").and_then(Value::as_str) { |
| 1213 | value["mimeType"] = Value::String(mime_type.to_string()); |
| 1214 | } |
| 1215 | if let Some(size) = metadata |
| 1216 | .get("size") |
| 1217 | .filter(|size| size.as_i64().is_some() || size.as_u64().is_some()) |
| 1218 | { |
| 1219 | value["size"] = size.clone(); |
| 1220 | } |
| 1221 | if let Some(annotations) = metadata |
| 1222 | .get("annotations") |
| 1223 | .filter(|annotations| valid_mcp_annotations(annotations)) |
| 1224 | { |
| 1225 | value["annotations"] = annotations.clone(); |
| 1226 | } |
| 1227 | value |
| 1228 | } |
| 1229 | |
| 1230 | fn valid_resource_contents(contents: &Value) -> bool { |
| 1231 | contents |
| 1232 | .as_array() |
| 1233 | .is_some_and(|contents| contents.iter().all(valid_resource_content)) |
| 1234 | } |
| 1235 | |
| 1236 | fn stdio_resource_read_result(uri: &str, result: Value) -> Value { |
| 1237 | let legacy_resource = result.clone(); |
| 1238 | match result { |
| 1239 | Value::Object(mut fields) |
| 1240 | if fields.get("contents").is_some_and(valid_resource_contents) |
| 1241 | && fields.get("_meta").is_none_or(Value::is_object) => |
| 1242 | { |
| 1243 | // Pass through a valid standard ReadResourceResult and retain the |
| 1244 | // old nested value for existing Codewhale management clients. |
| 1245 | fields.insert("resource".to_string(), legacy_resource); |
| 1246 | Value::Object(fields) |
| 1247 | } |
| 1248 | Value::String(text) => json!({ |
| 1249 | "contents": [{"uri": uri, "text": text}], |
| 1250 | "resource": legacy_resource |
| 1251 | }), |
| 1252 | value => json!({ |
| 1253 | "contents": [{"uri": uri, "text": legacy_value_text(&value)}], |
| 1254 | "resource": legacy_resource |
| 1255 | }), |
| 1256 | } |
| 1257 | } |
| 1258 | |
| 1259 | fn lifecycle_snapshot(state: &StdioMcpState) -> Value { |
| 1260 | let mut servers: Vec<Value> = state |
| 1261 | .definitions |
| 1262 | .iter() |
| 1263 | .map(|(name, definition)| { |
| 1264 | let is_running = state.running.get(name).copied().unwrap_or(false); |
| 1265 | json!({ |
| 1266 | "name": name, |
| 1267 | "enabled": definition.config.enabled, |
| 1268 | "running": is_running, |
| 1269 | "command": definition.config.command.clone(), |
| 1270 | "args": definition.config.args.clone(), |
| 1271 | // Null when the server is healthy. A client polling |
| 1272 | // `server/list` must be able to tell "up" from "never |
| 1273 | // started" without guessing. |
| 1274 | "error": state.errors.get(name).cloned(), |
| 1275 | }) |
| 1276 | }) |
| 1277 | .collect(); |
| 1278 | servers.sort_by(|a, b| { |
| 1279 | let a_name = a.get("name").and_then(Value::as_str).unwrap_or_default(); |
| 1280 | let b_name = b.get("name").and_then(Value::as_str).unwrap_or_default(); |
| 1281 | a_name.cmp(b_name) |
| 1282 | }); |
| 1283 | |
| 1284 | let running_count = state.running.values().filter(|running| **running).count(); |
| 1285 | json!({ |
| 1286 | "status": state.lifecycle_state, |
| 1287 | "servers": servers, |
| 1288 | "counts": { |
| 1289 | "defined": state.definitions.len(), |
| 1290 | "running": running_count |
| 1291 | } |
| 1292 | }) |
| 1293 | } |
| 1294 | |
| 1295 | fn params_or_object(params: Value) -> Value { |
| 1296 | if params.is_null() { json!({}) } else { params } |
| 1297 | } |
| 1298 | |
| 1299 | fn parse_params<T: DeserializeOwned>(params: Value) -> std::result::Result<T, JsonRpcError> { |
| 1300 | serde_json::from_value(params).map_err(|err| JsonRpcError::invalid_params(err.to_string())) |
| 1301 | } |
| 1302 | |
| 1303 | fn parse_server_from_uri(uri: &str) -> Option<String> { |
| 1304 | let stripped = uri.strip_prefix("mcp://")?; |
| 1305 | let server = stripped.split('/').next()?; |
| 1306 | if server.is_empty() { |
| 1307 | None |
| 1308 | } else { |
| 1309 | Some(server.to_string()) |
| 1310 | } |
| 1311 | } |
| 1312 | |
| 1313 | fn require_ready_session( |
| 1314 | state: &StdioMcpState, |
| 1315 | method: &str, |
| 1316 | ) -> std::result::Result<(), JsonRpcError> { |
| 1317 | if state.session_phase != McpSessionPhase::Ready { |
| 1318 | return Err(JsonRpcError::invalid_request(format!( |
| 1319 | "{method} requires a completed initialize / notifications/initialized handshake" |
| 1320 | ))); |
| 1321 | } |
| 1322 | Ok(()) |
| 1323 | } |
| 1324 | |
| 1325 | fn dispatch_stdio_request( |
| 1326 | state: &mut StdioMcpState, |
| 1327 | method: &str, |
| 1328 | params: Value, |
| 1329 | ) -> std::result::Result<(Value, bool), JsonRpcError> { |
| 1330 | match method { |
| 1331 | "initialize" => { |
| 1332 | if state.session_phase != McpSessionPhase::Uninitialized { |
| 1333 | return Err(JsonRpcError::invalid_request( |
| 1334 | "initialize may only be sent once per stdio session", |
| 1335 | )); |
| 1336 | } |
| 1337 | let parsed: InitializeParams = parse_params(params_or_object(params))?; |
| 1338 | if parsed.protocol_version.trim().is_empty() { |
| 1339 | return Err(JsonRpcError::invalid_params( |
| 1340 | "protocolVersion must not be empty", |
| 1341 | )); |
| 1342 | } |
| 1343 | if parsed.client_info.name.trim().is_empty() |
| 1344 | || parsed.client_info.version.trim().is_empty() |
| 1345 | { |
| 1346 | return Err(JsonRpcError::invalid_params( |
| 1347 | "clientInfo.name and clientInfo.version must not be empty", |
| 1348 | )); |
| 1349 | } |
| 1350 | // Deserializing into a Map above is the object-shape check. The |
| 1351 | // proxy does not currently consume any client capability. |
| 1352 | let _client_capabilities = parsed.capabilities; |
| 1353 | // Per spec, echo the requested revision when we support it; |
| 1354 | // otherwise answer with the newest revision we do support and |
| 1355 | // let the client decide whether to continue. |
| 1356 | let negotiated = |
| 1357 | if MCP_SUPPORTED_PROTOCOL_VERSIONS.contains(&parsed.protocol_version.as_str()) { |
| 1358 | parsed.protocol_version |
| 1359 | } else { |
| 1360 | MCP_PROTOCOL_VERSION.to_string() |
| 1361 | }; |
| 1362 | state.session_phase = McpSessionPhase::InitializeResponded; |
| 1363 | Ok((initialize_response(state, &negotiated), false)) |
| 1364 | } |
| 1365 | // Pre-standard Codewhale management alias; it intentionally requires |
| 1366 | // no MCP initialize envelope. |
| 1367 | "capabilities" => Ok((initialize_response(state, MCP_PROTOCOL_VERSION), false)), |
| 1368 | "notifications/initialized" => { |
| 1369 | if state.session_phase != McpSessionPhase::InitializeResponded { |
| 1370 | return Err(JsonRpcError::invalid_request( |
| 1371 | "notifications/initialized requires a successful initialize request", |
| 1372 | )); |
| 1373 | } |
| 1374 | state.session_phase = McpSessionPhase::Ready; |
| 1375 | Ok((json!({}), false)) |
| 1376 | } |
| 1377 | "ping" => Ok((json!({}), false)), |
| 1378 | "healthz" => Ok(( |
| 1379 | json!({ |
| 1380 | "status": "ok", |
| 1381 | "service": MCP_SERVER_NAME, |
| 1382 | "transport": "stdio", |
| 1383 | "lifecycle": lifecycle_snapshot(state) |
| 1384 | }), |
| 1385 | false, |
| 1386 | )), |
| 1387 | "tools/list" => { |
| 1388 | require_ready_session(state, method)?; |
| 1389 | let parsed: ToolsListParams = parse_params(params_or_object(params))?; |
| 1390 | let mut tools = state |
| 1391 | .manager |
| 1392 | .list_tools_with_input_schemas() |
| 1393 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 1394 | if let Some(server) = parsed.server { |
| 1395 | tools.retain(|(tool, _)| tool.server_name == server); |
| 1396 | } |
| 1397 | let tools = tools |
| 1398 | .into_iter() |
| 1399 | .map(stdio_tool_descriptor) |
| 1400 | .collect::<Vec<_>>(); |
| 1401 | Ok((json!({ "tools": tools }), false)) |
| 1402 | } |
| 1403 | "tools/call" => { |
| 1404 | require_ready_session(state, method)?; |
| 1405 | let parsed: ToolsCallParams = parse_params(params_or_object(params))?; |
| 1406 | let ToolsCallParams { |
| 1407 | name, |
| 1408 | tool, |
| 1409 | server, |
| 1410 | arguments, |
| 1411 | } = parsed; |
| 1412 | let tool_name = name |
| 1413 | .or(tool) |
| 1414 | .context("missing tool name") |
| 1415 | .map_err(|err| JsonRpcError::invalid_params(err.to_string()))?; |
| 1416 | if !arguments.is_object() { |
| 1417 | return Err(JsonRpcError::invalid_params( |
| 1418 | "tools/call arguments must be an object", |
| 1419 | )); |
| 1420 | } |
| 1421 | let result = if tool_name.starts_with("mcp__") { |
| 1422 | state |
| 1423 | .manager |
| 1424 | .call_qualified_tool(&tool_name, arguments) |
| 1425 | .map_err(|err| JsonRpcError::internal(err.to_string()))? |
| 1426 | } else { |
| 1427 | let server = server |
| 1428 | .context("missing server for unqualified tool") |
| 1429 | .map_err(|err| JsonRpcError::invalid_params(err.to_string()))?; |
| 1430 | state |
| 1431 | .manager |
| 1432 | .call_tool(&server, &tool_name, arguments) |
| 1433 | .map_err(|err| JsonRpcError::internal(err.to_string()))? |
| 1434 | }; |
| 1435 | let result = stdio_tool_call_result(result) |
| 1436 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 1437 | Ok((result, false)) |
| 1438 | } |
| 1439 | "resources/list" => { |
| 1440 | require_ready_session(state, method)?; |
| 1441 | let parsed: ResourcesListParams = parse_params(params_or_object(params))?; |
| 1442 | let mut resources = state |
| 1443 | .manager |
| 1444 | .list_resources_with_metadata() |
| 1445 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 1446 | if let Some(server) = parsed.server { |
| 1447 | resources.retain(|(resource, _)| resource.server_name == server); |
| 1448 | } |
| 1449 | let resources = resources |
| 1450 | .into_iter() |
| 1451 | .map(stdio_resource_descriptor) |
| 1452 | .collect::<Vec<_>>(); |
| 1453 | Ok((json!({ "resources": resources }), false)) |
| 1454 | } |
| 1455 | "resources/read" => { |
| 1456 | require_ready_session(state, method)?; |
| 1457 | let parsed: ResourcesReadParams = parse_params(params_or_object(params))?; |
| 1458 | let ResourcesReadParams { server, uri } = parsed; |
| 1459 | let value = match server { |
| 1460 | Some(server_name) => state.manager.read_resource(&server_name, &uri), |
| 1461 | None => state.manager.read_resource_by_uri(&uri), |
| 1462 | } |
| 1463 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 1464 | Ok((stdio_resource_read_result(&uri, value), false)) |
| 1465 | } |
| 1466 | "server/list" | "servers/list" => { |
| 1467 | Ok((json!({ "lifecycle": lifecycle_snapshot(state) }), false)) |
| 1468 | } |
| 1469 | "server/register" | "servers/register" => { |
| 1470 | let parsed: ServerRegisterParams = parse_params(params_or_object(params))?; |
| 1471 | let name = parsed.server.name.clone(); |
| 1472 | if name.trim().is_empty() { |
| 1473 | return Err(JsonRpcError::invalid_params( |
| 1474 | "server.name must not be empty", |
| 1475 | )); |
| 1476 | } |
| 1477 | |
| 1478 | if state.definitions.contains_key(&name) { |
| 1479 | let _ = state.manager.unregister_server(&name); |
| 1480 | } |
| 1481 | let definition = McpServerDefinition { |
| 1482 | config: parsed.server.clone(), |
| 1483 | filter: parsed.filter.clone(), |
| 1484 | }; |
| 1485 | state.definitions.insert(name.clone(), definition.clone()); |
| 1486 | state.errors.remove(&name); |
| 1487 | if parsed.start && parsed.server.enabled { |
| 1488 | // Registration is only "ok" if the configured command actually |
| 1489 | // came up. Reporting success here and answering later tool |
| 1490 | // calls from a stub is what #4727 was. |
| 1491 | state |
| 1492 | .start_definition(&definition) |
| 1493 | .map_err(|err| JsonRpcError::internal(format!("{err:#}")))?; |
| 1494 | } else { |
| 1495 | state.running.insert(name, false); |
| 1496 | } |
| 1497 | Ok((json!({ "lifecycle": lifecycle_snapshot(state) }), false)) |
| 1498 | } |
| 1499 | "server/start" | "servers/start" => { |
| 1500 | let parsed: ServerNameParams = parse_params(params_or_object(params))?; |
| 1501 | let definition = state |
| 1502 | .definitions |
| 1503 | .get(&parsed.name) |
| 1504 | .cloned() |
| 1505 | .with_context(|| format!("server '{}' is not defined", parsed.name)) |
| 1506 | .map_err(|err| JsonRpcError::invalid_params(err.to_string()))?; |
| 1507 | if !definition.config.enabled { |
| 1508 | return Err(JsonRpcError::invalid_params(format!( |
| 1509 | "server '{}' is disabled", |
| 1510 | parsed.name |
| 1511 | ))); |
| 1512 | } |
| 1513 | if !state.running.get(&parsed.name).copied().unwrap_or(false) { |
| 1514 | state |
| 1515 | .start_definition(&definition) |
| 1516 | .map_err(|err| JsonRpcError::internal(format!("{err:#}")))?; |
| 1517 | } |
| 1518 | Ok((json!({ "lifecycle": lifecycle_snapshot(state) }), false)) |
| 1519 | } |
| 1520 | "server/stop" | "servers/stop" => { |
| 1521 | let parsed: ServerNameParams = parse_params(params_or_object(params))?; |
| 1522 | if state.running.get(&parsed.name).copied().unwrap_or(false) { |
| 1523 | state |
| 1524 | .manager |
| 1525 | .stop_server(&parsed.name) |
| 1526 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 1527 | } |
| 1528 | // A deliberate stop is not a failure, so it clears any recorded |
| 1529 | // startup error rather than leaving a stale one on display. |
| 1530 | state.errors.remove(&parsed.name); |
| 1531 | state.running.insert(parsed.name, false); |
| 1532 | Ok((json!({ "lifecycle": lifecycle_snapshot(state) }), false)) |
| 1533 | } |
| 1534 | "server/unregister" | "servers/unregister" => { |
| 1535 | let parsed: ServerNameParams = parse_params(params_or_object(params))?; |
| 1536 | if state.definitions.remove(&parsed.name).is_none() { |
| 1537 | return Err(JsonRpcError::invalid_params(format!( |
| 1538 | "server '{}' is not defined", |
| 1539 | parsed.name |
| 1540 | ))); |
| 1541 | } |
| 1542 | let _ = state.manager.unregister_server(&parsed.name); |
| 1543 | state.running.remove(&parsed.name); |
| 1544 | state.errors.remove(&parsed.name); |
| 1545 | Ok((json!({ "lifecycle": lifecycle_snapshot(state) }), false)) |
| 1546 | } |
| 1547 | "shutdown" => { |
| 1548 | state.lifecycle_state = "shutting_down".to_string(); |
| 1549 | Ok(( |
| 1550 | json!({ |
| 1551 | "ok": true, |
| 1552 | "lifecycle": lifecycle_snapshot(state) |
| 1553 | }), |
| 1554 | true, |
| 1555 | )) |
| 1556 | } |
| 1557 | _ => Err(JsonRpcError::method_not_found(method)), |
| 1558 | } |
| 1559 | } |
| 1560 | |
| 1561 | fn jsonrpc_result(id: Option<Value>, result: Value) -> Value { |
| 1562 | json!({ |
| 1563 | "jsonrpc": "2.0", |
| 1564 | "id": id.unwrap_or(Value::Null), |
| 1565 | "result": result |
| 1566 | }) |
| 1567 | } |
| 1568 | |
| 1569 | fn jsonrpc_error(id: Option<Value>, err: JsonRpcError) -> Value { |
| 1570 | json!({ |
| 1571 | "jsonrpc": "2.0", |
| 1572 | "id": id.unwrap_or(Value::Null), |
| 1573 | "error": { |
| 1574 | "code": err.code, |
| 1575 | "message": err.message, |
| 1576 | "data": err.data |
| 1577 | } |
| 1578 | }) |
| 1579 | } |
| 1580 | |
| 1581 | impl JsonRpcError { |
| 1582 | fn parse_error(message: impl Into<String>) -> Self { |
| 1583 | Self { |
| 1584 | code: -32700, |
| 1585 | message: message.into(), |
| 1586 | data: None, |
| 1587 | } |
| 1588 | } |
| 1589 | |
| 1590 | fn invalid_request(message: impl Into<String>) -> Self { |
| 1591 | Self { |
| 1592 | code: -32600, |
| 1593 | message: message.into(), |
| 1594 | data: None, |
| 1595 | } |
| 1596 | } |
| 1597 | |
| 1598 | fn method_not_found(method: &str) -> Self { |
| 1599 | Self { |
| 1600 | code: -32601, |
| 1601 | message: format!("unsupported method: {method}"), |
| 1602 | data: None, |
| 1603 | } |
| 1604 | } |
| 1605 | |
| 1606 | fn invalid_params(message: impl Into<String>) -> Self { |
| 1607 | Self { |
| 1608 | code: -32602, |
| 1609 | message: message.into(), |
| 1610 | data: None, |
| 1611 | } |
| 1612 | } |
| 1613 | |
| 1614 | fn internal(message: impl Into<String>) -> Self { |
| 1615 | Self { |
| 1616 | code: -32603, |
| 1617 | message: message.into(), |
| 1618 | data: None, |
| 1619 | } |
| 1620 | } |
| 1621 | } |
| 1622 | |
| 1623 | #[cfg(test)] |
| 1624 | mod tests { |
| 1625 | use std::sync::Arc; |
| 1626 | use std::sync::atomic::{AtomicUsize, Ordering}; |
| 1627 | |
| 1628 | use super::*; |
| 1629 | |
| 1630 | fn complete_stdio_handshake(state: &mut StdioMcpState) { |
| 1631 | dispatch_stdio_request( |
| 1632 | state, |
| 1633 | "initialize", |
| 1634 | json!({ |
| 1635 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 1636 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 1637 | "capabilities": {} |
| 1638 | }), |
| 1639 | ) |
| 1640 | .expect("initialize"); |
| 1641 | dispatch_stdio_request(state, "notifications/initialized", Value::Null) |
| 1642 | .expect("initialized notification"); |
| 1643 | } |
| 1644 | |
| 1645 | struct EchoMcpClient; |
| 1646 | |
| 1647 | impl McpManagedClient for EchoMcpClient { |
| 1648 | fn list_tools(&self) -> Result<Vec<McpToolDescriptor>> { |
| 1649 | Ok(vec![]) |
| 1650 | } |
| 1651 | |
| 1652 | fn call_tool(&self, tool_name: &str, arguments: Value) -> Result<Value> { |
| 1653 | if tool_name == "error" { |
| 1654 | bail!("intentional error for testing"); |
| 1655 | } |
| 1656 | Ok(arguments) |
| 1657 | } |
| 1658 | |
| 1659 | fn list_resources(&self) -> Result<Vec<McpResourceDescriptor>> { |
| 1660 | Ok(vec![]) |
| 1661 | } |
| 1662 | |
| 1663 | fn read_resource(&self, _uri: &str) -> Result<Value> { |
| 1664 | bail!("not supported") |
| 1665 | } |
| 1666 | } |
| 1667 | |
| 1668 | // ── InMemoryMcpClient ────────────────────────────────────────────── |
| 1669 | |
| 1670 | #[test] |
| 1671 | fn in_memory_client_list_tools_returns_registered() { |
| 1672 | let client = InMemoryMcpClient::default() |
| 1673 | .with_tool("echo", json!({"output": "hi"})) |
| 1674 | .with_tool("greet", json!({"msg": "hello"})); |
| 1675 | let tools = client.list_tools().unwrap(); |
| 1676 | assert_eq!(tools.len(), 2); |
| 1677 | let names: Vec<&str> = tools.iter().map(|t| t.tool_name.as_str()).collect(); |
| 1678 | assert!(names.contains(&"echo")); |
| 1679 | assert!(names.contains(&"greet")); |
| 1680 | } |
| 1681 | |
| 1682 | #[test] |
| 1683 | fn in_memory_client_call_tool_returns_value() { |
| 1684 | let client = InMemoryMcpClient::default().with_tool("echo", json!({"output": "hi"})); |
| 1685 | let result = client.call_tool("echo", json!({})).unwrap(); |
| 1686 | assert_eq!(result["output"], "hi"); |
| 1687 | } |
| 1688 | |
| 1689 | #[test] |
| 1690 | fn in_memory_client_call_tool_errors_on_missing() { |
| 1691 | let client = InMemoryMcpClient::default(); |
| 1692 | let err = client.call_tool("nope", json!({})).unwrap_err(); |
| 1693 | assert!(err.to_string().contains("not found")); |
| 1694 | } |
| 1695 | |
| 1696 | #[test] |
| 1697 | fn in_memory_client_list_resources_returns_registered() { |
| 1698 | let client = InMemoryMcpClient::default() |
| 1699 | .with_resource("mcp://s/health", json!({"ok": true})) |
| 1700 | .with_resource("mcp://s/caps", json!({"tools": []})); |
| 1701 | let resources = client.list_resources().unwrap(); |
| 1702 | assert_eq!(resources.len(), 2); |
| 1703 | } |
| 1704 | |
| 1705 | #[test] |
| 1706 | fn in_memory_client_read_resource_returns_value() { |
| 1707 | let client = |
| 1708 | InMemoryMcpClient::default().with_resource("mcp://s/health", json!({"ok": true})); |
| 1709 | let result = client.read_resource("mcp://s/health").unwrap(); |
| 1710 | assert_eq!(result["ok"], true); |
| 1711 | } |
| 1712 | |
| 1713 | #[test] |
| 1714 | fn in_memory_client_read_resource_errors_on_missing() { |
| 1715 | let client = InMemoryMcpClient::default(); |
| 1716 | let err = client.read_resource("mcp://s/nope").unwrap_err(); |
| 1717 | assert!(err.to_string().contains("not found")); |
| 1718 | } |
| 1719 | |
| 1720 | // ── McpManager ───────────────────────────────────────────────────── |
| 1721 | |
| 1722 | fn make_server_config(name: &str) -> McpServerConfig { |
| 1723 | McpServerConfig { |
| 1724 | name: name.to_string(), |
| 1725 | command: "test".to_string(), |
| 1726 | args: vec![], |
| 1727 | env: HashMap::new(), |
| 1728 | enabled: true, |
| 1729 | } |
| 1730 | } |
| 1731 | |
| 1732 | /// Client that counts `call_tool` invocations and always fails, so a |
| 1733 | /// retry shows up as a count rather than as a swallowed error. |
| 1734 | #[derive(Default)] |
| 1735 | struct CountingFailingClient { |
| 1736 | calls: Arc<AtomicUsize>, |
| 1737 | } |
| 1738 | |
| 1739 | impl McpManagedClient for CountingFailingClient { |
| 1740 | fn list_tools(&self) -> Result<Vec<McpToolDescriptor>> { |
| 1741 | Ok(vec![McpToolDescriptor { |
| 1742 | server_name: "counting".to_string(), |
| 1743 | tool_name: "write".to_string(), |
| 1744 | qualified_name: "write".to_string(), |
| 1745 | description: None, |
| 1746 | }]) |
| 1747 | } |
| 1748 | |
| 1749 | fn call_tool(&self, _tool_name: &str, _arguments: Value) -> Result<Value> { |
| 1750 | self.calls.fetch_add(1, Ordering::SeqCst); |
| 1751 | bail!("transient upstream failure") |
| 1752 | } |
| 1753 | |
| 1754 | fn list_resources(&self) -> Result<Vec<McpResourceDescriptor>> { |
| 1755 | Ok(Vec::new()) |
| 1756 | } |
| 1757 | |
| 1758 | fn read_resource(&self, _uri: &str) -> Result<Value> { |
| 1759 | bail!("no resources") |
| 1760 | } |
| 1761 | } |
| 1762 | |
| 1763 | #[test] |
| 1764 | fn failed_qualified_tool_call_is_not_retried() { |
| 1765 | // #4728: the fast path used to fall through to a re-resolution loop |
| 1766 | // whenever the *call* errored, re-invoking the same tool. For a file |
| 1767 | // write, a commit, or a paid API call, that second invocation is a |
| 1768 | // second real side effect. |
| 1769 | let calls = Arc::new(AtomicUsize::new(0)); |
| 1770 | let mut manager = McpManager::default(); |
| 1771 | manager |
| 1772 | .register_server( |
| 1773 | make_server_config("writer"), |
| 1774 | ToolFilter::default(), |
| 1775 | Box::new(CountingFailingClient { |
| 1776 | calls: Arc::clone(&calls), |
| 1777 | }), |
| 1778 | ) |
| 1779 | .unwrap(); |
| 1780 | |
| 1781 | let err = manager |
| 1782 | .call_qualified_tool("mcp__writer__write", json!({})) |
| 1783 | .unwrap_err(); |
| 1784 | |
| 1785 | assert_eq!( |
| 1786 | calls.load(Ordering::SeqCst), |
| 1787 | 1, |
| 1788 | "tool must be invoked exactly once, got {} invocations", |
| 1789 | calls.load(Ordering::SeqCst) |
| 1790 | ); |
| 1791 | // The original error is propagated, not discarded in favour of a |
| 1792 | // later attempt's. |
| 1793 | assert!( |
| 1794 | err.to_string().contains("transient upstream failure"), |
| 1795 | "unexpected error: {err}" |
| 1796 | ); |
| 1797 | } |
| 1798 | |
| 1799 | #[test] |
| 1800 | fn register_server_rejects_a_name_that_collides_after_sanitizing() { |
| 1801 | // #4729: `my-server` and `my_server` both qualify as `mcp__my_server__*`, |
| 1802 | // so registering both would let either answer a call meant for the |
| 1803 | // other, decided by HashMap iteration order. |
| 1804 | let mut manager = McpManager::default(); |
| 1805 | manager |
| 1806 | .register_server( |
| 1807 | make_server_config("my_server"), |
| 1808 | ToolFilter::default(), |
| 1809 | Box::new(InMemoryMcpClient::default().with_tool("t", json!("trusted"))), |
| 1810 | ) |
| 1811 | .unwrap(); |
| 1812 | |
| 1813 | for colliding in ["my-server", "My.Server"] { |
| 1814 | let err = manager |
| 1815 | .register_server( |
| 1816 | make_server_config(colliding), |
| 1817 | ToolFilter::default(), |
| 1818 | Box::new(InMemoryMcpClient::default().with_tool("t", json!("hostile"))), |
| 1819 | ) |
| 1820 | .unwrap_err(); |
| 1821 | assert!( |
| 1822 | err.to_string().contains("collides"), |
| 1823 | "expected collision error for {colliding}, got: {err}" |
| 1824 | ); |
| 1825 | } |
| 1826 | |
| 1827 | // The trusted server keeps answering its own qualified name. |
| 1828 | assert_eq!( |
| 1829 | manager |
| 1830 | .call_qualified_tool("mcp__my_server__t", json!({})) |
| 1831 | .unwrap(), |
| 1832 | json!("trusted") |
| 1833 | ); |
| 1834 | } |
| 1835 | |
| 1836 | #[test] |
| 1837 | fn same_server_tool_name_collisions_fail_closed_in_list_and_call() { |
| 1838 | let mut manager = McpManager::default(); |
| 1839 | manager |
| 1840 | .register_server( |
| 1841 | make_server_config("s1"), |
| 1842 | ToolFilter::default(), |
| 1843 | Box::new( |
| 1844 | InMemoryMcpClient::default() |
| 1845 | .with_tool("foo-bar", json!("hyphen")) |
| 1846 | .with_tool("foo_bar", json!("underscore")), |
| 1847 | ), |
| 1848 | ) |
| 1849 | .unwrap(); |
| 1850 | |
| 1851 | let list_error = manager.list_tools().unwrap_err(); |
| 1852 | assert!( |
| 1853 | list_error.to_string().contains("ambiguous"), |
| 1854 | "unexpected list error: {list_error}" |
| 1855 | ); |
| 1856 | |
| 1857 | let call_error = manager |
| 1858 | .call_qualified_tool("mcp__s1__foo_bar", json!({})) |
| 1859 | .unwrap_err(); |
| 1860 | assert!( |
| 1861 | call_error |
| 1862 | .to_string() |
| 1863 | .contains("ambiguous within server 's1'"), |
| 1864 | "unexpected call error: {call_error}" |
| 1865 | ); |
| 1866 | } |
| 1867 | |
| 1868 | #[test] |
| 1869 | fn re_registering_the_same_server_name_replaces_it() { |
| 1870 | // Collision rejection must not break restart, which re-registers the |
| 1871 | // same name with a fresh client. |
| 1872 | let mut manager = McpManager::default(); |
| 1873 | for value in ["first", "second"] { |
| 1874 | manager |
| 1875 | .register_server( |
| 1876 | make_server_config("s1"), |
| 1877 | ToolFilter::default(), |
| 1878 | Box::new(InMemoryMcpClient::default().with_tool("t", json!(value))), |
| 1879 | ) |
| 1880 | .unwrap(); |
| 1881 | } |
| 1882 | assert_eq!( |
| 1883 | manager |
| 1884 | .call_qualified_tool("mcp__s1__t", json!({})) |
| 1885 | .unwrap(), |
| 1886 | json!("second") |
| 1887 | ); |
| 1888 | } |
| 1889 | |
| 1890 | #[test] |
| 1891 | fn manager_start_all_marks_ready_for_registered_clients() { |
| 1892 | let mut manager = McpManager::default(); |
| 1893 | manager |
| 1894 | .register_server( |
| 1895 | make_server_config("s1"), |
| 1896 | ToolFilter::default(), |
| 1897 | Box::new(InMemoryMcpClient::default().with_tool("t", json!(null))), |
| 1898 | ) |
| 1899 | .unwrap(); |
| 1900 | let mut events = Vec::new(); |
| 1901 | let summary = manager.start_all(|e| events.push(e)); |
| 1902 | assert_eq!(summary.ready, vec!["s1"]); |
| 1903 | assert!(summary.failed.is_empty()); |
| 1904 | assert!(events.iter().any(|event| { |
| 1905 | event.server_name == "s1" && event.status == McpStartupStatus::Starting |
| 1906 | })); |
| 1907 | assert!( |
| 1908 | events.iter().any(|event| { |
| 1909 | event.server_name == "s1" && event.status == McpStartupStatus::Ready |
| 1910 | }) |
| 1911 | ); |
| 1912 | } |
| 1913 | |
| 1914 | #[test] |
| 1915 | fn manager_start_all_marks_failed_when_client_missing() { |
| 1916 | let mut manager = McpManager::default(); |
| 1917 | manager |
| 1918 | .register_server( |
| 1919 | make_server_config("s1"), |
| 1920 | ToolFilter::default(), |
| 1921 | Box::new(InMemoryMcpClient::default()), |
| 1922 | ) |
| 1923 | .unwrap(); |
| 1924 | manager.stop_server("s1").unwrap(); |
| 1925 | let summary = manager.start_all(|_| {}); |
| 1926 | assert!(summary.ready.is_empty()); |
| 1927 | assert_eq!(summary.failed.len(), 1); |
| 1928 | assert_eq!(summary.failed[0].server_name, "s1"); |
| 1929 | } |
| 1930 | |
| 1931 | #[test] |
| 1932 | fn manager_start_all_cancels_disabled_servers() { |
| 1933 | let mut manager = McpManager::default(); |
| 1934 | let mut cfg = make_server_config("s1"); |
| 1935 | cfg.enabled = false; |
| 1936 | manager |
| 1937 | .register_server( |
| 1938 | cfg, |
| 1939 | ToolFilter::default(), |
| 1940 | Box::new(InMemoryMcpClient::default()), |
| 1941 | ) |
| 1942 | .unwrap(); |
| 1943 | let summary = manager.start_all(|_| {}); |
| 1944 | assert!(summary.ready.is_empty()); |
| 1945 | assert_eq!(summary.cancelled, vec!["s1"]); |
| 1946 | } |
| 1947 | |
| 1948 | #[test] |
| 1949 | fn manager_list_tools_applies_filter() { |
| 1950 | let mut manager = McpManager::default(); |
| 1951 | let client = InMemoryMcpClient::default() |
| 1952 | .with_tool("allowed", json!(null)) |
| 1953 | .with_tool("denied", json!(null)); |
| 1954 | manager |
| 1955 | .register_server( |
| 1956 | make_server_config("s1"), |
| 1957 | ToolFilter { |
| 1958 | allow: vec!["allowed".to_string()], |
| 1959 | deny: vec![], |
| 1960 | }, |
| 1961 | Box::new(client), |
| 1962 | ) |
| 1963 | .unwrap(); |
| 1964 | let tools = manager.list_tools().unwrap(); |
| 1965 | assert_eq!(tools.len(), 1); |
| 1966 | assert_eq!(tools[0].tool_name, "allowed"); |
| 1967 | } |
| 1968 | |
| 1969 | #[test] |
| 1970 | fn manager_list_tools_deny_overrides_allow() { |
| 1971 | let mut manager = McpManager::default(); |
| 1972 | let client = InMemoryMcpClient::default() |
| 1973 | .with_tool("a", json!(null)) |
| 1974 | .with_tool("b", json!(null)); |
| 1975 | manager |
| 1976 | .register_server( |
| 1977 | make_server_config("s1"), |
| 1978 | ToolFilter { |
| 1979 | allow: vec!["a".to_string(), "b".to_string()], |
| 1980 | deny: vec!["b".to_string()], |
| 1981 | }, |
| 1982 | Box::new(client), |
| 1983 | ) |
| 1984 | .unwrap(); |
| 1985 | let tools = manager.list_tools().unwrap(); |
| 1986 | assert_eq!(tools.len(), 1); |
| 1987 | assert_eq!(tools[0].tool_name, "a"); |
| 1988 | } |
| 1989 | |
| 1990 | #[test] |
| 1991 | fn manager_call_tool_delegates_to_client() { |
| 1992 | let mut manager = McpManager::default(); |
| 1993 | manager |
| 1994 | .register_server( |
| 1995 | make_server_config("s1"), |
| 1996 | ToolFilter::default(), |
| 1997 | Box::new(InMemoryMcpClient::default().with_tool("t", json!({"v": 42}))), |
| 1998 | ) |
| 1999 | .unwrap(); |
| 2000 | let result = manager.call_tool("s1", "t", json!({})).unwrap(); |
| 2001 | assert_eq!(result["v"], 42); |
| 2002 | } |
| 2003 | |
| 2004 | #[test] |
| 2005 | fn manager_call_tool_passes_arguments_to_client() { |
| 2006 | let mut manager = McpManager::default(); |
| 2007 | manager |
| 2008 | .register_server( |
| 2009 | make_server_config("s1"), |
| 2010 | ToolFilter::default(), |
| 2011 | Box::new(EchoMcpClient), |
| 2012 | ) |
| 2013 | .unwrap(); |
| 2014 | let args = json!({"hello": "world", "num": 100}); |
| 2015 | let result = manager.call_tool("s1", "echo", args.clone()).unwrap(); |
| 2016 | assert_eq!(result, args); |
| 2017 | } |
| 2018 | |
| 2019 | #[test] |
| 2020 | fn manager_call_tool_propagates_client_error() { |
| 2021 | let mut manager = McpManager::default(); |
| 2022 | manager |
| 2023 | .register_server( |
| 2024 | make_server_config("s1"), |
| 2025 | ToolFilter::default(), |
| 2026 | Box::new(EchoMcpClient), |
| 2027 | ) |
| 2028 | .unwrap(); |
| 2029 | let err = manager.call_tool("s1", "error", json!({})).unwrap_err(); |
| 2030 | assert!(err.to_string().contains("intentional error for testing")); |
| 2031 | } |
| 2032 | |
| 2033 | #[test] |
| 2034 | fn manager_call_tool_errors_on_missing_server() { |
| 2035 | let manager = McpManager::default(); |
| 2036 | let err = manager.call_tool("nope", "t", json!({})).unwrap_err(); |
| 2037 | assert!(err.to_string().contains("not available")); |
| 2038 | } |
| 2039 | |
| 2040 | #[test] |
| 2041 | fn manager_call_tool_enforces_deny_filter() { |
| 2042 | // The filter used to be consulted only when listing tools; a denied |
| 2043 | // tool stayed callable by addressing the server directly. |
| 2044 | let mut manager = McpManager::default(); |
| 2045 | manager |
| 2046 | .register_server( |
| 2047 | make_server_config("s1"), |
| 2048 | ToolFilter { |
| 2049 | allow: vec![], |
| 2050 | deny: vec!["secret".to_string()], |
| 2051 | }, |
| 2052 | Box::new(InMemoryMcpClient::default().with_tool("secret", json!({"ok": true}))), |
| 2053 | ) |
| 2054 | .unwrap(); |
| 2055 | let err = manager.call_tool("s1", "secret", json!({})).unwrap_err(); |
| 2056 | assert!( |
| 2057 | err.to_string().contains("blocked by the tool filter"), |
| 2058 | "unexpected error: {err}" |
| 2059 | ); |
| 2060 | } |
| 2061 | |
| 2062 | #[test] |
| 2063 | fn manager_call_tool_enforces_allow_filter() { |
| 2064 | let mut manager = McpManager::default(); |
| 2065 | manager |
| 2066 | .register_server( |
| 2067 | make_server_config("s1"), |
| 2068 | ToolFilter { |
| 2069 | allow: vec!["allowed".to_string()], |
| 2070 | deny: vec![], |
| 2071 | }, |
| 2072 | Box::new( |
| 2073 | InMemoryMcpClient::default() |
| 2074 | .with_tool("allowed", json!({"ok": true})) |
| 2075 | .with_tool("other", json!({"ok": false})), |
| 2076 | ), |
| 2077 | ) |
| 2078 | .unwrap(); |
| 2079 | let err = manager.call_tool("s1", "other", json!({})).unwrap_err(); |
| 2080 | assert!( |
| 2081 | err.to_string().contains("blocked by the tool filter"), |
| 2082 | "unexpected error: {err}" |
| 2083 | ); |
| 2084 | // The allowed tool still runs. |
| 2085 | assert_eq!( |
| 2086 | manager.call_tool("s1", "allowed", json!({})).unwrap(), |
| 2087 | json!({"ok": true}) |
| 2088 | ); |
| 2089 | } |
| 2090 | |
| 2091 | #[test] |
| 2092 | fn denied_tool_cannot_be_called_by_qualified_name() { |
| 2093 | // Security: `mcp__s1__secret` must be as unreachable as `secret`. |
| 2094 | let mut manager = McpManager::default(); |
| 2095 | manager |
| 2096 | .register_server( |
| 2097 | make_server_config("s1"), |
| 2098 | ToolFilter { |
| 2099 | allow: vec![], |
| 2100 | deny: vec!["secret".to_string()], |
| 2101 | }, |
| 2102 | Box::new(InMemoryMcpClient::default().with_tool("secret", json!({"ok": true}))), |
| 2103 | ) |
| 2104 | .unwrap(); |
| 2105 | let err = manager |
| 2106 | .call_qualified_tool("mcp__s1__secret", json!({})) |
| 2107 | .unwrap_err(); |
| 2108 | assert!( |
| 2109 | err.to_string().contains("blocked by the tool filter"), |
| 2110 | "unexpected error: {err}" |
| 2111 | ); |
| 2112 | } |
| 2113 | |
| 2114 | #[test] |
| 2115 | fn manager_call_qualified_tool_parses_name() { |
| 2116 | let mut manager = McpManager::default(); |
| 2117 | manager |
| 2118 | .register_server( |
| 2119 | make_server_config("my_server"), |
| 2120 | ToolFilter::default(), |
| 2121 | Box::new(InMemoryMcpClient::default().with_tool("my_tool", json!({"ok": true}))), |
| 2122 | ) |
| 2123 | .unwrap(); |
| 2124 | let result = manager |
| 2125 | .call_qualified_tool("mcp__my_server__my_tool", json!({})) |
| 2126 | .unwrap(); |
| 2127 | assert_eq!(result["ok"], true); |
| 2128 | } |
| 2129 | |
| 2130 | #[test] |
| 2131 | fn manager_call_qualified_tool_resolves_sanitized_segment_to_original_name() { |
| 2132 | // `qualify_tool_name` folds `-`/`.`/case into `_`, so the qualified |
| 2133 | // name advertised for `my-tool` is `mcp__s1__my_tool`. The exact-match |
| 2134 | // fast path used to dispatch that sanitized segment verbatim, and the |
| 2135 | // server (which only knows `my-tool`) rejected the call. |
| 2136 | let mut manager = McpManager::default(); |
| 2137 | manager |
| 2138 | .register_server( |
| 2139 | make_server_config("s1"), |
| 2140 | ToolFilter::default(), |
| 2141 | Box::new( |
| 2142 | InMemoryMcpClient::default() |
| 2143 | .with_tool("my-tool", json!({"via": "hyphen"})) |
| 2144 | .with_tool("other.thing", json!({"via": "dot"})), |
| 2145 | ), |
| 2146 | ) |
| 2147 | .unwrap(); |
| 2148 | |
| 2149 | let hyphen = manager |
| 2150 | .call_qualified_tool("mcp__s1__my_tool", json!({})) |
| 2151 | .unwrap(); |
| 2152 | assert_eq!(hyphen, json!({"via": "hyphen"})); |
| 2153 | |
| 2154 | let dot = manager |
| 2155 | .call_qualified_tool("mcp__s1__other_thing", json!({})) |
| 2156 | .unwrap(); |
| 2157 | assert_eq!(dot, json!({"via": "dot"})); |
| 2158 | } |
| 2159 | |
| 2160 | #[test] |
| 2161 | fn manager_call_qualified_tool_handles_truncated_names() { |
| 2162 | let long_server = "server".repeat(20); |
| 2163 | let long_tool = "tool".repeat(20); |
| 2164 | let mut manager = McpManager::default(); |
| 2165 | manager |
| 2166 | .register_server( |
| 2167 | make_server_config(&long_server), |
| 2168 | ToolFilter::default(), |
| 2169 | Box::new(InMemoryMcpClient::default().with_tool(&long_tool, json!({"ok": true}))), |
| 2170 | ) |
| 2171 | .unwrap(); |
| 2172 | let tools = manager.list_tools().unwrap(); |
| 2173 | let qualified = &tools[0].qualified_name; |
| 2174 | assert!(qualified.len() <= 64); |
| 2175 | assert!(parse_qualified_tool_name(qualified).is_ok()); |
| 2176 | |
| 2177 | let result = manager.call_qualified_tool(qualified, json!({})).unwrap(); |
| 2178 | assert_eq!(result["ok"], true); |
| 2179 | } |
| 2180 | |
| 2181 | #[test] |
| 2182 | fn manager_unregister_removes_server() { |
| 2183 | let mut manager = McpManager::default(); |
| 2184 | manager |
| 2185 | .register_server( |
| 2186 | make_server_config("s1"), |
| 2187 | ToolFilter::default(), |
| 2188 | Box::new(InMemoryMcpClient::default()), |
| 2189 | ) |
| 2190 | .unwrap(); |
| 2191 | manager.unregister_server("s1").unwrap(); |
| 2192 | assert!(manager.configs.is_empty()); |
| 2193 | } |
| 2194 | |
| 2195 | #[test] |
| 2196 | fn manager_unregister_errors_on_unknown() { |
| 2197 | let mut manager = McpManager::default(); |
| 2198 | let err = manager.unregister_server("nope").unwrap_err(); |
| 2199 | assert!(err.to_string().contains("not registered")); |
| 2200 | } |
| 2201 | |
| 2202 | #[test] |
| 2203 | fn manager_stop_server_errors_on_unknown() { |
| 2204 | let mut manager = McpManager::default(); |
| 2205 | let err = manager.stop_server("nope").unwrap_err(); |
| 2206 | assert!(err.to_string().contains("not running")); |
| 2207 | } |
| 2208 | |
| 2209 | #[test] |
| 2210 | fn manager_list_resources_returns_from_clients() { |
| 2211 | let mut manager = McpManager::default(); |
| 2212 | manager |
| 2213 | .register_server( |
| 2214 | make_server_config("s1"), |
| 2215 | ToolFilter::default(), |
| 2216 | Box::new( |
| 2217 | InMemoryMcpClient::default() |
| 2218 | .with_resource("mcp://s1/health", json!({"ok": true})), |
| 2219 | ), |
| 2220 | ) |
| 2221 | .unwrap(); |
| 2222 | let resources = manager.list_resources().unwrap(); |
| 2223 | assert_eq!(resources.len(), 1); |
| 2224 | assert_eq!(resources[0].server_name, "s1"); |
| 2225 | } |
| 2226 | |
| 2227 | #[test] |
| 2228 | fn manager_read_resource_delegates() { |
| 2229 | let mut manager = McpManager::default(); |
| 2230 | manager |
| 2231 | .register_server( |
| 2232 | make_server_config("s1"), |
| 2233 | ToolFilter::default(), |
| 2234 | Box::new( |
| 2235 | InMemoryMcpClient::default() |
| 2236 | .with_resource("mcp://s1/health", json!({"ok": true})), |
| 2237 | ), |
| 2238 | ) |
| 2239 | .unwrap(); |
| 2240 | let result = manager.read_resource("s1", "mcp://s1/health").unwrap(); |
| 2241 | assert_eq!(result["ok"], true); |
| 2242 | } |
| 2243 | |
| 2244 | #[test] |
| 2245 | fn manager_resolves_a_unique_standard_resource_uri() { |
| 2246 | let mut manager = McpManager::default(); |
| 2247 | manager |
| 2248 | .register_server( |
| 2249 | make_server_config("docs"), |
| 2250 | ToolFilter::default(), |
| 2251 | Box::new( |
| 2252 | InMemoryMcpClient::default() |
| 2253 | .with_resource("file:///guide.md", json!({"text": "guide"})), |
| 2254 | ), |
| 2255 | ) |
| 2256 | .unwrap(); |
| 2257 | |
| 2258 | let result = manager.read_resource_by_uri("file:///guide.md").unwrap(); |
| 2259 | assert_eq!(result["text"], "guide"); |
| 2260 | } |
| 2261 | |
| 2262 | #[test] |
| 2263 | fn manager_rejects_an_ambiguous_standard_resource_uri() { |
| 2264 | let mut manager = McpManager::default(); |
| 2265 | for server in ["alpha", "beta"] { |
| 2266 | manager |
| 2267 | .register_server( |
| 2268 | make_server_config(server), |
| 2269 | ToolFilter::default(), |
| 2270 | Box::new( |
| 2271 | InMemoryMcpClient::default() |
| 2272 | .with_resource("file:///shared.txt", json!({"server": server})), |
| 2273 | ), |
| 2274 | ) |
| 2275 | .unwrap(); |
| 2276 | } |
| 2277 | |
| 2278 | let err = manager |
| 2279 | .read_resource_by_uri("file:///shared.txt") |
| 2280 | .unwrap_err(); |
| 2281 | let message = err.to_string(); |
| 2282 | assert!(message.contains("ambiguous"), "unexpected error: {message}"); |
| 2283 | assert!( |
| 2284 | message.contains("alpha, beta"), |
| 2285 | "server names must be deterministic: {message}" |
| 2286 | ); |
| 2287 | } |
| 2288 | |
| 2289 | #[test] |
| 2290 | fn manager_update_sandbox_state_returns_notices() { |
| 2291 | let mut manager = McpManager::default(); |
| 2292 | manager |
| 2293 | .register_server( |
| 2294 | make_server_config("s1"), |
| 2295 | ToolFilter::default(), |
| 2296 | Box::new(InMemoryMcpClient::default()), |
| 2297 | ) |
| 2298 | .unwrap(); |
| 2299 | let notices = manager.update_sandbox_state("strict", "/tmp").unwrap(); |
| 2300 | assert_eq!(notices.len(), 1); |
| 2301 | assert_eq!(notices[0]["server_name"], "s1"); |
| 2302 | } |
| 2303 | |
| 2304 | // ── Tool filter ──────────────────────────────────────────────────── |
| 2305 | |
| 2306 | #[test] |
| 2307 | fn allowed_by_filter_empty_allow_permits_all() { |
| 2308 | let filter = ToolFilter { |
| 2309 | allow: vec![], |
| 2310 | deny: vec![], |
| 2311 | }; |
| 2312 | assert!(allowed_by_filter("anything", &filter)); |
| 2313 | } |
| 2314 | |
| 2315 | #[test] |
| 2316 | fn allowed_by_filter_deny_blocks() { |
| 2317 | let filter = ToolFilter { |
| 2318 | allow: vec![], |
| 2319 | deny: vec!["danger".to_string()], |
| 2320 | }; |
| 2321 | assert!(!allowed_by_filter("danger", &filter)); |
| 2322 | assert!(allowed_by_filter("safe", &filter)); |
| 2323 | } |
| 2324 | |
| 2325 | #[test] |
| 2326 | fn allowed_by_filter_allow_only_permits_listed() { |
| 2327 | let filter = ToolFilter { |
| 2328 | allow: vec!["a".to_string()], |
| 2329 | deny: vec![], |
| 2330 | }; |
| 2331 | assert!(allowed_by_filter("a", &filter)); |
| 2332 | assert!(!allowed_by_filter("b", &filter)); |
| 2333 | } |
| 2334 | |
| 2335 | // ── Helper functions ─────────────────────────────────────────────── |
| 2336 | |
| 2337 | #[test] |
| 2338 | fn sanitize_component_lowercases_and_replaces_specials() { |
| 2339 | assert_eq!(sanitize_component("My-Server.Name"), "my_server_name"); |
| 2340 | assert_eq!(sanitize_component("ABC123"), "abc123"); |
| 2341 | } |
| 2342 | |
| 2343 | #[test] |
| 2344 | fn qualify_tool_name_produces_mcp_prefix() { |
| 2345 | let name = qualify_tool_name("server", "tool"); |
| 2346 | assert!(name.starts_with("mcp__server__tool")); |
| 2347 | } |
| 2348 | |
| 2349 | #[test] |
| 2350 | fn qualify_tool_name_truncates_long_names() { |
| 2351 | let long_server = "a".repeat(100); |
| 2352 | let name = qualify_tool_name(&long_server, "tool"); |
| 2353 | assert!(name.len() <= 64); |
| 2354 | assert!(parse_qualified_tool_name(&name).is_ok()); |
| 2355 | } |
| 2356 | |
| 2357 | #[test] |
| 2358 | fn parse_qualified_tool_name_round_trip() { |
| 2359 | let qualified = qualify_tool_name("my_server", "my_tool"); |
| 2360 | let (server, tool) = parse_qualified_tool_name(&qualified).unwrap(); |
| 2361 | assert_eq!(server, "my_server"); |
| 2362 | assert_eq!(tool, "my_tool"); |
| 2363 | } |
| 2364 | |
| 2365 | #[test] |
| 2366 | fn parse_qualified_tool_name_rejects_missing_prefix() { |
| 2367 | let err = parse_qualified_tool_name("not_mcp__server__tool").unwrap_err(); |
| 2368 | assert!(err.to_string().contains("missing mcp__ prefix")); |
| 2369 | } |
| 2370 | |
| 2371 | #[test] |
| 2372 | fn parse_qualified_tool_name_rejects_empty_segments() { |
| 2373 | let err = parse_qualified_tool_name("mcp____tool").unwrap_err(); |
| 2374 | assert!(err.to_string().contains("missing server segment")); |
| 2375 | } |
| 2376 | |
| 2377 | #[test] |
| 2378 | fn parse_server_from_uri_extracts_server() { |
| 2379 | assert_eq!( |
| 2380 | parse_server_from_uri("mcp://my-server/capabilities"), |
| 2381 | Some("my-server".to_string()) |
| 2382 | ); |
| 2383 | } |
| 2384 | |
| 2385 | #[test] |
| 2386 | fn parse_server_from_uri_returns_none_for_invalid() { |
| 2387 | assert!(parse_server_from_uri("http://not-mcp").is_none()); |
| 2388 | assert!(parse_server_from_uri("mcp:///path").is_none()); |
| 2389 | } |
| 2390 | |
| 2391 | // ── JsonRpcError ─────────────────────────────────────────────────── |
| 2392 | |
| 2393 | #[test] |
| 2394 | fn jsonrpc_error_codes_are_correct() { |
| 2395 | assert_eq!(JsonRpcError::parse_error("").code, -32700); |
| 2396 | assert_eq!(JsonRpcError::invalid_request("").code, -32600); |
| 2397 | assert_eq!(JsonRpcError::method_not_found("x").code, -32601); |
| 2398 | assert_eq!(JsonRpcError::invalid_params("").code, -32602); |
| 2399 | assert_eq!(JsonRpcError::internal("").code, -32603); |
| 2400 | } |
| 2401 | |
| 2402 | #[test] |
| 2403 | fn jsonrpc_result_produces_valid_envelope() { |
| 2404 | let result = jsonrpc_result(Some(json!(1)), json!({"ok": true})); |
| 2405 | assert_eq!(result["jsonrpc"], "2.0"); |
| 2406 | assert_eq!(result["id"], 1); |
| 2407 | assert_eq!(result["result"]["ok"], true); |
| 2408 | } |
| 2409 | |
| 2410 | #[test] |
| 2411 | fn jsonrpc_error_produces_valid_envelope() { |
| 2412 | let err = jsonrpc_error(Some(json!(2)), JsonRpcError::invalid_params("bad")); |
| 2413 | assert_eq!(err["jsonrpc"], "2.0"); |
| 2414 | assert_eq!(err["id"], 2); |
| 2415 | assert_eq!(err["error"]["code"], -32602); |
| 2416 | } |
| 2417 | |
| 2418 | #[test] |
| 2419 | fn jsonrpc_missing_and_explicit_null_ids_remain_distinct() { |
| 2420 | let notification: JsonRpcRequest = serde_json::from_value(json!({ |
| 2421 | "jsonrpc": "2.0", |
| 2422 | "method": "ping" |
| 2423 | })) |
| 2424 | .unwrap(); |
| 2425 | assert!(!notification.id.should_respond()); |
| 2426 | assert_eq!(notification.id.response_id(), None); |
| 2427 | |
| 2428 | let null_id: JsonRpcRequest = serde_json::from_value(json!({ |
| 2429 | "jsonrpc": "2.0", |
| 2430 | "id": null, |
| 2431 | "method": "ping" |
| 2432 | })) |
| 2433 | .unwrap(); |
| 2434 | assert!(null_id.id.should_respond()); |
| 2435 | assert_eq!(null_id.id.response_id(), Some(Value::Null)); |
| 2436 | |
| 2437 | for invalid in [ |
| 2438 | json!({"method": "ping"}), |
| 2439 | json!({"jsonrpc": null, "method": "ping"}), |
| 2440 | json!({"jsonrpc": "2.0", "id": {}, "method": "ping"}), |
| 2441 | ] { |
| 2442 | assert!( |
| 2443 | serde_json::from_value::<JsonRpcRequest>(invalid).is_err(), |
| 2444 | "invalid envelope was accepted" |
| 2445 | ); |
| 2446 | } |
| 2447 | } |
| 2448 | |
| 2449 | #[test] |
| 2450 | fn stdio_initialize_uses_standard_mcp_shape_and_codewhale_identity() { |
| 2451 | let state = build_stdio_state(Vec::new()); |
| 2452 | let response = initialize_response(&state, MCP_PROTOCOL_VERSION); |
| 2453 | assert_eq!(response["protocolVersion"], MCP_PROTOCOL_VERSION); |
| 2454 | assert_eq!(response["serverInfo"]["name"], MCP_SERVER_NAME); |
| 2455 | assert_eq!(response["serverInfo"]["version"], env!("CARGO_PKG_VERSION")); |
| 2456 | assert!(response["capabilities"]["tools"].is_object()); |
| 2457 | assert!(response["capabilities"]["resources"].is_object()); |
| 2458 | assert_eq!(response["server"], MCP_SERVER_NAME); |
| 2459 | } |
| 2460 | |
| 2461 | #[test] |
| 2462 | fn stdio_initialize_validates_required_client_fields() { |
| 2463 | let mut state = build_stdio_state(Vec::new()); |
| 2464 | let valid = dispatch_stdio_request( |
| 2465 | &mut state, |
| 2466 | "initialize", |
| 2467 | json!({ |
| 2468 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 2469 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 2470 | "capabilities": {} |
| 2471 | }), |
| 2472 | ) |
| 2473 | .expect("valid initialize"); |
| 2474 | assert_eq!(valid.0["protocolVersion"], MCP_PROTOCOL_VERSION); |
| 2475 | |
| 2476 | for invalid in [ |
| 2477 | json!({}), |
| 2478 | json!({ |
| 2479 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 2480 | "clientInfo": {"name": "test-client", "version": "1"} |
| 2481 | }), |
| 2482 | json!({ |
| 2483 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 2484 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 2485 | "capabilities": [] |
| 2486 | }), |
| 2487 | json!({ |
| 2488 | "protocolVersion": "", |
| 2489 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 2490 | "capabilities": {} |
| 2491 | }), |
| 2492 | ] { |
| 2493 | let mut state = build_stdio_state(Vec::new()); |
| 2494 | let error = dispatch_stdio_request(&mut state, "initialize", invalid) |
| 2495 | .expect_err("malformed initialize must fail"); |
| 2496 | assert_eq!(error.code, -32602); |
| 2497 | } |
| 2498 | } |
| 2499 | |
| 2500 | #[test] |
| 2501 | fn standard_catalog_methods_require_the_complete_mcp_handshake() { |
| 2502 | let mut state = build_stdio_state(Vec::new()); |
| 2503 | let early_initialized = |
| 2504 | dispatch_stdio_request(&mut state, "notifications/initialized", Value::Null) |
| 2505 | .expect_err("initialized cannot precede initialize"); |
| 2506 | assert_eq!(early_initialized.code, -32600); |
| 2507 | |
| 2508 | // Explicit Codewhale management compatibility remains available before |
| 2509 | // MCP initialization. |
| 2510 | assert!(dispatch_stdio_request(&mut state, "capabilities", Value::Null).is_ok()); |
| 2511 | assert!(dispatch_stdio_request(&mut state, "server/list", Value::Null).is_ok()); |
| 2512 | |
| 2513 | for method in [ |
| 2514 | "tools/list", |
| 2515 | "tools/call", |
| 2516 | "resources/list", |
| 2517 | "resources/read", |
| 2518 | ] { |
| 2519 | let error = dispatch_stdio_request(&mut state, method, json!({})) |
| 2520 | .expect_err("standard MCP data methods must fail before initialize"); |
| 2521 | assert_eq!(error.code, -32600); |
| 2522 | assert!(error.message.contains("completed initialize")); |
| 2523 | } |
| 2524 | |
| 2525 | dispatch_stdio_request( |
| 2526 | &mut state, |
| 2527 | "initialize", |
| 2528 | json!({ |
| 2529 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 2530 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 2531 | "capabilities": {} |
| 2532 | }), |
| 2533 | ) |
| 2534 | .unwrap(); |
| 2535 | let too_early = dispatch_stdio_request(&mut state, "tools/list", json!({})) |
| 2536 | .expect_err("initialized notification is required"); |
| 2537 | assert_eq!(too_early.code, -32600); |
| 2538 | |
| 2539 | dispatch_stdio_request(&mut state, "notifications/initialized", Value::Null).unwrap(); |
| 2540 | let tools = dispatch_stdio_request(&mut state, "tools/list", json!({})).unwrap(); |
| 2541 | assert_eq!(tools.0["tools"], json!([])); |
| 2542 | |
| 2543 | let duplicate = dispatch_stdio_request( |
| 2544 | &mut state, |
| 2545 | "initialize", |
| 2546 | json!({ |
| 2547 | "protocolVersion": MCP_PROTOCOL_VERSION, |
| 2548 | "clientInfo": {"name": "test-client", "version": "1"}, |
| 2549 | "capabilities": {} |
| 2550 | }), |
| 2551 | ) |
| 2552 | .expect_err("duplicate initialize must fail"); |
| 2553 | assert_eq!(duplicate.code, -32600); |
| 2554 | } |
| 2555 | |
| 2556 | #[test] |
| 2557 | fn stdio_tools_call_rejects_non_object_arguments() { |
| 2558 | let mut state = build_stdio_state(Vec::new()); |
| 2559 | complete_stdio_handshake(&mut state); |
| 2560 | for arguments in [json!(null), json!("bad"), json!([]), json!(1)] { |
| 2561 | let error = dispatch_stdio_request( |
| 2562 | &mut state, |
| 2563 | "tools/call", |
| 2564 | json!({"name": "mcp__missing__tool", "arguments": arguments}), |
| 2565 | ) |
| 2566 | .expect_err("non-object arguments must fail before dispatch"); |
| 2567 | assert_eq!(error.code, -32602); |
| 2568 | assert!(error.message.contains("arguments must be an object")); |
| 2569 | } |
| 2570 | } |
| 2571 | |
| 2572 | #[test] |
| 2573 | fn stdio_tool_results_only_pass_through_valid_mcp_content_arrays() { |
| 2574 | let standard = json!({ |
| 2575 | "content": [{"type": "text", "text": "ok"}], |
| 2576 | "isError": true |
| 2577 | }); |
| 2578 | let standard_response = stdio_tool_call_result(standard.clone()).unwrap(); |
| 2579 | assert_eq!(standard_response["content"][0]["text"], "ok"); |
| 2580 | assert_eq!(standard_response["isError"], true); |
| 2581 | assert_eq!(standard_response["result"], standard); |
| 2582 | |
| 2583 | let legacy = json!({"answer": 42}); |
| 2584 | let legacy_response = stdio_tool_call_result(legacy.clone()).unwrap(); |
| 2585 | assert_eq!(legacy_response["content"][0]["type"], "text"); |
| 2586 | assert_eq!(legacy_response["result"], legacy); |
| 2587 | |
| 2588 | for malformed in [ |
| 2589 | json!({"content": {"type": "text", "text": "not-an-array"}}), |
| 2590 | json!({"content": [null]}), |
| 2591 | json!({"content": [{"type": "text"}]}), |
| 2592 | json!({"content": [{"type": "image", "data": "abc"}]}), |
| 2593 | json!({"content": [{"type": "resource", "resource": {"uri": "file:///x"}}]}), |
| 2594 | json!({"content": [], "isError": "false"}), |
| 2595 | json!({"content": [], "_meta": []}), |
| 2596 | json!({"isError": true}), |
| 2597 | json!({"_meta": {"trace": "x"}}), |
| 2598 | ] { |
| 2599 | let error = stdio_tool_call_result(malformed.clone()) |
| 2600 | .expect_err("CallToolResult-shaped malformed output must fail closed"); |
| 2601 | assert!(error.to_string().contains("malformed MCP CallToolResult")); |
| 2602 | } |
| 2603 | |
| 2604 | let legacy_string = stdio_tool_call_result(json!("plain text")).unwrap(); |
| 2605 | assert_eq!(legacy_string["content"][0]["text"], "plain text"); |
| 2606 | } |
| 2607 | |
| 2608 | #[test] |
| 2609 | fn malformed_call_tool_result_shape_becomes_a_protocol_error() { |
| 2610 | let mut state = build_stdio_state(Vec::new()); |
| 2611 | state |
| 2612 | .manager |
| 2613 | .register_server( |
| 2614 | make_server_config("malformed"), |
| 2615 | ToolFilter::default(), |
| 2616 | Box::new( |
| 2617 | InMemoryMcpClient::default().with_tool("broken", json!({"isError": true})), |
| 2618 | ), |
| 2619 | ) |
| 2620 | .unwrap(); |
| 2621 | complete_stdio_handshake(&mut state); |
| 2622 | |
| 2623 | let error = dispatch_stdio_request( |
| 2624 | &mut state, |
| 2625 | "tools/call", |
| 2626 | json!({"name": "mcp__malformed__broken", "arguments": {}}), |
| 2627 | ) |
| 2628 | .expect_err("malformed CallToolResult-shaped output must not become success text"); |
| 2629 | assert_eq!(error.code, -32603); |
| 2630 | assert!(error.message.contains("malformed MCP CallToolResult")); |
| 2631 | } |
| 2632 | |
| 2633 | #[test] |
| 2634 | fn stdio_resources_use_standard_shapes_with_legacy_metadata() { |
| 2635 | let descriptor = McpResourceDescriptor { |
| 2636 | server_name: "docs".to_string(), |
| 2637 | uri: "file:///guide.md".to_string(), |
| 2638 | description: Some("User guide".to_string()), |
| 2639 | }; |
| 2640 | let listed = stdio_resource_descriptor(( |
| 2641 | descriptor, |
| 2642 | json!({ |
| 2643 | "name": "Guide", |
| 2644 | "mimeType": "text/markdown", |
| 2645 | "size": 42, |
| 2646 | "annotations": {"audience": ["assistant"], "priority": 0.75} |
| 2647 | }), |
| 2648 | )); |
| 2649 | assert_eq!(listed["uri"], "file:///guide.md"); |
| 2650 | assert_eq!(listed["name"], "Guide"); |
| 2651 | assert_eq!(listed["mimeType"], "text/markdown"); |
| 2652 | assert_eq!(listed["size"], 42); |
| 2653 | assert_eq!(listed["annotations"]["audience"], json!(["assistant"])); |
| 2654 | assert_eq!(listed["annotations"]["priority"], 0.75); |
| 2655 | assert_eq!(listed["server_name"], "docs"); |
| 2656 | |
| 2657 | let standard = json!({ |
| 2658 | "contents": [{ |
| 2659 | "uri": "file:///guide.md", |
| 2660 | "mimeType": "text/markdown", |
| 2661 | "text": "hello" |
| 2662 | }] |
| 2663 | }); |
| 2664 | let standard_response = stdio_resource_read_result("file:///guide.md", standard.clone()); |
| 2665 | assert_eq!(standard_response["contents"][0]["text"], "hello"); |
| 2666 | assert_eq!(standard_response["resource"], standard); |
| 2667 | |
| 2668 | for legacy in [ |
| 2669 | json!({"body": "legacy"}), |
| 2670 | json!({"contents": [{"text": "missing URI"}]}), |
| 2671 | json!({"contents": [], "_meta": []}), |
| 2672 | ] { |
| 2673 | let response = stdio_resource_read_result("file:///guide.md", legacy.clone()); |
| 2674 | assert_eq!(response["contents"][0]["uri"], "file:///guide.md"); |
| 2675 | assert!(response["contents"][0]["text"].is_string()); |
| 2676 | assert_eq!(response["resource"], legacy); |
| 2677 | } |
| 2678 | } |
| 2679 | |
| 2680 | // ── stdio dispatch: no stub may answer for a configured server ───── |
| 2681 | |
| 2682 | fn definition(name: &str, command: &str, args: &[&str]) -> McpServerDefinition { |
| 2683 | McpServerDefinition { |
| 2684 | config: McpServerConfig { |
| 2685 | name: name.to_string(), |
| 2686 | command: command.to_string(), |
| 2687 | args: args.iter().map(|arg| (*arg).to_string()).collect(), |
| 2688 | env: HashMap::new(), |
| 2689 | enabled: true, |
| 2690 | }, |
| 2691 | filter: ToolFilter::default(), |
| 2692 | } |
| 2693 | } |
| 2694 | |
| 2695 | fn call(state: &mut StdioMcpState, method: &str, params: Value) -> Value { |
| 2696 | dispatch_stdio_request(state, method, params) |
| 2697 | .unwrap_or_else(|err| panic!("{method} failed: {}", err.message)) |
| 2698 | .0 |
| 2699 | } |
| 2700 | |
| 2701 | #[test] |
| 2702 | fn a_server_that_cannot_be_spawned_is_reported_not_running() { |
| 2703 | // #4727: this used to register a stub and report `running: true`, so a |
| 2704 | // typo in `command` was indistinguishable from a working server. |
| 2705 | let mut state = build_stdio_state(vec![definition( |
| 2706 | "broken", |
| 2707 | "codewhale-nonexistent-mcp-server-binary", |
| 2708 | &[], |
| 2709 | )]); |
| 2710 | |
| 2711 | let lifecycle = call(&mut state, "server/list", json!({}))["lifecycle"].clone(); |
| 2712 | assert_eq!(lifecycle["servers"][0]["running"], json!(false)); |
| 2713 | let error = lifecycle["servers"][0]["error"] |
| 2714 | .as_str() |
| 2715 | .expect("a stopped server must carry its failure reason"); |
| 2716 | assert!( |
| 2717 | error.contains("failed to spawn command"), |
| 2718 | "unexpected error: {error}" |
| 2719 | ); |
| 2720 | |
| 2721 | // And nothing answers on its behalf. |
| 2722 | complete_stdio_handshake(&mut state); |
| 2723 | let err = dispatch_stdio_request( |
| 2724 | &mut state, |
| 2725 | "tools/call", |
| 2726 | json!({"name": "mcp__broken__health", "arguments": {}}), |
| 2727 | ) |
| 2728 | .expect_err("a server that never started must not answer tool calls"); |
| 2729 | assert_eq!(err.code, -32603); |
| 2730 | } |
| 2731 | |
| 2732 | #[cfg(unix)] |
| 2733 | #[test] |
| 2734 | fn tools_come_from_the_spawned_process_not_a_stub() { |
| 2735 | let script = crate::test_support::write_fake_mcp_server("dispatch_tools"); |
| 2736 | let mut state = build_stdio_state(vec![definition( |
| 2737 | "fake", |
| 2738 | "/bin/sh", |
| 2739 | &[script.path().to_str().expect("utf-8 script path")], |
| 2740 | )]); |
| 2741 | complete_stdio_handshake(&mut state); |
| 2742 | |
| 2743 | let tools = call(&mut state, "tools/list", json!({}))["tools"].clone(); |
| 2744 | let names: Vec<&str> = tools |
| 2745 | .as_array() |
| 2746 | .expect("tools array") |
| 2747 | .iter() |
| 2748 | .filter_map(|tool| tool["name"].as_str()) |
| 2749 | .collect(); |
| 2750 | assert_eq!( |
| 2751 | names, |
| 2752 | vec!["mcp__fake__add"], |
| 2753 | "only the child's real tools may be listed, got {names:?}" |
| 2754 | ); |
| 2755 | assert_eq!(tools[0]["tool_name"], "add"); |
| 2756 | assert_eq!(tools[0]["inputSchema"]["required"], json!(["a", "b"])); |
| 2757 | |
| 2758 | let result = call( |
| 2759 | &mut state, |
| 2760 | "tools/call", |
| 2761 | json!({"name": "mcp__fake__add", "arguments": {"a": 2, "b": 3}}), |
| 2762 | ); |
| 2763 | assert_eq!(result["content"][0]["text"], "5"); |
| 2764 | assert_eq!(result["result"]["content"][0]["text"], "5"); |
| 2765 | |
| 2766 | let resources = call(&mut state, "resources/list", json!({}))["resources"].clone(); |
| 2767 | assert_eq!(resources[0]["uri"], "file:///fake/readme.txt"); |
| 2768 | assert_eq!(resources[0]["name"], "Fake readme"); |
| 2769 | assert_eq!(resources[0]["mimeType"], "text/plain"); |
| 2770 | assert_eq!(resources[0]["size"], 16); |
| 2771 | assert_eq!( |
| 2772 | resources[0]["annotations"]["audience"], |
| 2773 | json!(["assistant"]) |
| 2774 | ); |
| 2775 | assert_eq!(resources[0]["server_name"], "fake"); |
| 2776 | |
| 2777 | let read = call( |
| 2778 | &mut state, |
| 2779 | "resources/read", |
| 2780 | json!({"uri": "file:///fake/readme.txt"}), |
| 2781 | ); |
| 2782 | assert_eq!(read["contents"][0]["text"], "spawned-resource"); |
| 2783 | assert_eq!(read["resource"]["contents"][0]["text"], "spawned-resource"); |
| 2784 | } |
| 2785 | |
| 2786 | #[cfg(unix)] |
| 2787 | #[test] |
| 2788 | fn server_register_fails_when_the_command_cannot_be_started() { |
| 2789 | let mut state = build_stdio_state(Vec::new()); |
| 2790 | let err = dispatch_stdio_request( |
| 2791 | &mut state, |
| 2792 | "server/register", |
| 2793 | json!({"server": {"name": "late", "command": "/bin/sh", "args": ["-c", "exit 1"]}}), |
| 2794 | ) |
| 2795 | .expect_err("registering an unstartable server must not report success"); |
| 2796 | assert_eq!(err.code, -32603); |
| 2797 | assert!( |
| 2798 | err.message.contains("initialize"), |
| 2799 | "unexpected error: {}", |
| 2800 | err.message |
| 2801 | ); |
| 2802 | assert_eq!(state.running.get("late"), Some(&false)); |
| 2803 | } |
| 2804 | |
| 2805 | // ── McpServerConfig serialization ────────────────────────────────── |
| 2806 | |
| 2807 | #[test] |
| 2808 | fn mcp_server_config_defaults_enabled_to_true() { |
| 2809 | let json = json!({"name": "s", "command": "cmd"}); |
| 2810 | let config: McpServerConfig = serde_json::from_value(json).unwrap(); |
| 2811 | assert!(config.enabled); |
| 2812 | assert!(config.args.is_empty()); |
| 2813 | assert!(config.env.is_empty()); |
| 2814 | } |
| 2815 | |
| 2816 | #[test] |
| 2817 | fn mcp_startup_status_serializes_with_snake_case() { |
| 2818 | let status = McpStartupStatus::Failed { |
| 2819 | error: "oops".to_string(), |
| 2820 | }; |
| 2821 | let json = serde_json::to_value(&status).unwrap(); |
| 2822 | assert_eq!(json["failed"]["error"], "oops"); |
| 2823 | } |
| 2824 | } |
| 2825 |