| 1 | //! Recorded-MCP family: one transcript, every dispatch implementation. |
| 2 | //! |
| 3 | //! A transcript (`<case>.case.json`) scripts an MCP server — what it answers |
| 4 | //! to `initialize`, `tools/list`, `resources/*`, `prompts/*`, and to each |
| 5 | //! `tools/call` (a result, an `isError` result, a JSON-RPC error, progress |
| 6 | //! notifications before the result, or no answer at all until the caller |
| 7 | //! cancels) — plus the steps a host takes against it. The harness serves the |
| 8 | //! transcript as a real Streamable HTTP MCP server on loopback |
| 9 | //! ([`server`]), so any client that can open a URL can be pointed at it, |
| 10 | //! including a child Node process; or, for `"transport": "stdio"`, as a |
| 11 | //! POSIX-shell child (`stdio_server.sh.fixture`, named so the repo's `*.sh` |
| 12 | //! ignore rule does not swallow it) the dispatch spawns itself. |
| 13 | //! |
| 14 | //! [`DISPATCHES`] is the seam for the TypeScript migration |
| 15 | //! (TS-EXTENSION-HOST-DESIGN §5.2, §9.3): today it holds the Rust pool path |
| 16 | //! production uses (`Engine::execute_mcp_tool_with_pool`); `HostMcpDispatch` |
| 17 | //! now joins it and must produce the *same* golden, because the golden |
| 18 | //! does not name the dispatch. It uses the real pinned SDK and Rust broker. |
| 19 | //! |
| 20 | //! Normalization: the server URL/port is masked; results are |
| 21 | //! `{"ok": {success, content, metadata, content_blocks}}` with JSON content |
| 22 | //! parsed, or `{"err": {kind, detail}}` with exact detail bytes. |
| 23 | //! `server_received` lists the side-effecting requests the server saw |
| 24 | //! (`tools/call`, `resources/read`, `prompts/get`) — a call that must not be |
| 25 | //! sent, or must not be replayed, shows up there. |
| 26 | //! |
| 27 | //! # Case fields beyond the original two cases |
| 28 | //! |
| 29 | //! All optional; an absent field is the original behaviour, so the original |
| 30 | //! goldens did not move. |
| 31 | //! |
| 32 | //! - `transport`: `"http"` (default) or `"stdio"`. A stdio case's `server` |
| 33 | //! object is written to files for `stdio_server.sh.fixture`; only `initialize`, |
| 34 | //! `tools/list` and `tools/call` (by `match.name`, with `"exit": true` for |
| 35 | //! a server that dies after reading the call) are supported. Unix only. |
| 36 | //! - `config`: extra `McpServerConfig` fields merged over the harness's own. |
| 37 | //! - `disallowed_tools`: deny rules handed to the dispatch (pool-level and |
| 38 | //! per-call, as the engine does). |
| 39 | //! - `expect_boot`: `"error"` records a failed boot as |
| 40 | //! `{"op":"boot","outcome":{"err":detail}}` instead of failing the harness; |
| 41 | //! steps still run, so the catalog after a failed boot is pinned too. A |
| 42 | //! case that expects an error and boots cleanly fails the harness. |
| 43 | //! - `server_options`: `numbered_sessions`, `require_bearer` (see [`server`]). |
| 44 | //! - `bearer_secret`: a fake token the harness exports as the server's |
| 45 | //! `bearer_token_env_var`. It must appear in no recorded byte: the outcome |
| 46 | //! is scanned for it before a golden is compared or written, in update mode |
| 47 | //! too. |
| 48 | //! - `sink`: start a second server and put its URL where the spec says |
| 49 | //! `{{sink_url}}`; the outcome records how many requests reached it. |
| 50 | //! - `record_requests`: `"full"` adds every POST the server saw (session, |
| 51 | //! protocol version, authorization shape, status) to the outcome as |
| 52 | //! `requests`; `"summary"` adds only the distinct methods, authorization |
| 53 | //! shapes and statuses (`requests_summary`), for cases that must not pin how |
| 54 | //! often a client retries a connect. |
| 55 | //! |
| 56 | //! Steps: `catalog`, `call`, and `approval_hints` (`{"op":"approval_hints", |
| 57 | //! "tools":[model names]}`, read after a `catalog` step). |
| 58 | //! |
| 59 | //! What a host-backed dispatch is *not* tested on by this family is listed in |
| 60 | //! the fixture README's coverage boundaries: reviewed-plugin launches (hash |
| 61 | //! refusal, trusted-read-only approval hints), real OAuth login, the SSE |
| 62 | //! legacy transport, the 32 MiB aggregate catalog byte cap, concurrent |
| 63 | //! calls, and a held stdio call. |
| 64 | |
| 65 | mod controls; |
| 66 | mod server; |
| 67 | |
| 68 | use std::path::{Path, PathBuf}; |
| 69 | use std::sync::{Arc, Mutex}; |
| 70 | use std::time::Duration; |
| 71 | |
| 72 | use serde_json::{Value, json}; |
| 73 | use tokio::sync::{Mutex as AsyncMutex, mpsc}; |
| 74 | use tokio_util::sync::CancellationToken; |
| 75 | |
| 76 | use self::server::{ServerOptions, TranscriptServer}; |
| 77 | use super::golden::{self, Failures, Sandbox}; |
| 78 | use crate::core::engine::Engine; |
| 79 | use crate::core::events::Event; |
| 80 | use crate::mcp::{McpConfig, McpPool, McpToolApprovalHint}; |
| 81 | use crate::test_support::EnvVarGuard; |
| 82 | use crate::tools::spec::{RichToolResult, ToolError}; |
| 83 | |
| 84 | const FAMILY: &str = "mcp"; |
| 85 | const CALL_DEADLINE: Duration = Duration::from_secs(30); |
| 86 | /// The environment variable a case's `bearer_secret` is exported through. |
| 87 | const BEARER_ENV: &str = "CODEWHALE_CONFORMANCE_MCP_TOKEN"; |
| 88 | |
| 89 | // === The dispatch seam === |
| 90 | |
| 91 | /// Everything a dispatch is built from: the connection config, and the deny |
| 92 | /// rules the engine would also hand to the pool (`--disallowed-tools`). |
| 93 | struct DispatchSetup { |
| 94 | config: McpConfig, |
| 95 | disallowed_tools: Vec<String>, |
| 96 | } |
| 97 | |
| 98 | /// One implementation of "call a tool on an external MCP server", as the |
| 99 | /// engine sees it. Planning, hooks and approval happen before `call`; |
| 100 | /// implementors never gate (design §5.2) — except that deny rules are part of |
| 101 | /// the dispatch's setup and refuse a call before anything reaches a server. |
| 102 | #[async_trait::async_trait] |
| 103 | trait McpDispatchUnderTest: Send + Sync { |
| 104 | /// Connect to every configured server. An `Err` is a connect failure |
| 105 | /// (including needs-auth); the dispatch must still answer `catalog` and |
| 106 | /// `call` afterwards, with whatever it can still advertise. |
| 107 | async fn boot(&self) -> Result<(), String>; |
| 108 | /// The model-visible tool catalog this dispatch advertises. |
| 109 | async fn catalog(&self) -> Vec<codewhale_models::Tool>; |
| 110 | /// Call one model-facing tool; `cancel` is the turn interrupt. |
| 111 | async fn call( |
| 112 | &self, |
| 113 | model_name: &str, |
| 114 | input: Value, |
| 115 | cancel: CancellationToken, |
| 116 | ) -> Result<RichToolResult, ToolError>; |
| 117 | /// Addition for the approval-hint case: how the approval path may treat |
| 118 | /// `model_name` given what the server's catalog declared, read after |
| 119 | /// `catalog()`. `"destructive"` (each call keeps its prompt), |
| 120 | /// `"trusted_read_only"` (a reviewed plugin's read-only tool runs |
| 121 | /// unprompted), or `None`. Implementation-neutral: any dispatch that |
| 122 | /// consumes MCP tool annotations has this answer. |
| 123 | async fn approval_hint(&self, model_name: &str) -> Option<&'static str>; |
| 124 | async fn shutdown(&self); |
| 125 | } |
| 126 | |
| 127 | type DispatchFactory = fn(setup: DispatchSetup) -> Box<dyn McpDispatchUnderTest>; |
| 128 | |
| 129 | /// Every dispatch runs every transcript against the same golden. |
| 130 | const DISPATCHES: &[(&str, DispatchFactory)] = &[ |
| 131 | ("mcp_pool", McpPoolDispatch::boxed), |
| 132 | ("host_sdk", HostMcpDispatch::boxed), |
| 133 | ]; |
| 134 | |
| 135 | /// The Rust path production uses today: the shared pool behind the engine's |
| 136 | /// direct MCP execution seam. |
| 137 | struct McpPoolDispatch { |
| 138 | pool: Arc<AsyncMutex<McpPool>>, |
| 139 | tx_event: mpsc::Sender<Event>, |
| 140 | disallowed_tools: Vec<String>, |
| 141 | // Kept open so status events the dispatch emits never hit a closed channel. |
| 142 | _rx_event: Mutex<mpsc::Receiver<Event>>, |
| 143 | } |
| 144 | |
| 145 | impl McpPoolDispatch { |
| 146 | fn boxed(setup: DispatchSetup) -> Box<dyn McpDispatchUnderTest> { |
| 147 | let (tx_event, rx_event) = mpsc::channel(64); |
| 148 | let pool = McpPool::new(setup.config).with_disallowed_tools(setup.disallowed_tools.clone()); |
| 149 | Box::new(Self { |
| 150 | pool: Arc::new(AsyncMutex::new(pool)), |
| 151 | tx_event, |
| 152 | disallowed_tools: setup.disallowed_tools, |
| 153 | _rx_event: Mutex::new(rx_event), |
| 154 | }) |
| 155 | } |
| 156 | } |
| 157 | |
| 158 | #[async_trait::async_trait] |
| 159 | impl McpDispatchUnderTest for McpPoolDispatch { |
| 160 | async fn boot(&self) -> Result<(), String> { |
| 161 | let errors = self.pool.lock().await.connect_all().await; |
| 162 | if errors.is_empty() { |
| 163 | Ok(()) |
| 164 | } else { |
| 165 | Err(errors |
| 166 | .iter() |
| 167 | .map(|(name, error)| format!("{name}: {error:#}")) |
| 168 | .collect::<Vec<_>>() |
| 169 | .join("; ")) |
| 170 | } |
| 171 | } |
| 172 | |
| 173 | async fn catalog(&self) -> Vec<codewhale_models::Tool> { |
| 174 | self.pool.lock().await.to_api_tools() |
| 175 | } |
| 176 | |
| 177 | async fn call( |
| 178 | &self, |
| 179 | model_name: &str, |
| 180 | input: Value, |
| 181 | cancel: CancellationToken, |
| 182 | ) -> Result<RichToolResult, ToolError> { |
| 183 | // The engine interrupts a tool by dropping its future; do the same. |
| 184 | tokio::select! { |
| 185 | biased; |
| 186 | () = cancel.cancelled() => Err(ToolError::cancelled("tool call interrupted by the host")), |
| 187 | result = Engine::execute_mcp_tool_with_pool( |
| 188 | Arc::clone(&self.pool), |
| 189 | &self.tx_event, |
| 190 | model_name, |
| 191 | input, |
| 192 | // The engine hands the registry context's rules to every call |
| 193 | // as well as to the pool; do both. |
| 194 | &self.disallowed_tools, |
| 195 | // Conformance probes cannot supply a person's card decision. |
| 196 | None, |
| 197 | ) => result, |
| 198 | } |
| 199 | } |
| 200 | |
| 201 | async fn approval_hint(&self, model_name: &str) -> Option<&'static str> { |
| 202 | crate::mcp::mcp_tool_approval_hint(model_name).map(|hint| match hint { |
| 203 | McpToolApprovalHint::TrustedReadOnly => "trusted_read_only", |
| 204 | McpToolApprovalHint::Destructive => "destructive", |
| 205 | }) |
| 206 | } |
| 207 | |
| 208 | async fn shutdown(&self) { |
| 209 | self.pool.lock().await.shutdown_all().await; |
| 210 | } |
| 211 | } |
| 212 | |
| 213 | /// The selected production Host pool, not a semantic fake. The sandbox |
| 214 | /// already seals CODEWHALE_HOME; manager materialization and native child |
| 215 | /// launch use that same root. Guards live through shutdown on this runner's |
| 216 | /// current-thread runtime, so host workers never inherit another case's policy. |
| 217 | struct HostMcpDispatch { |
| 218 | inner: McpPoolDispatch, |
| 219 | manager: Arc<crate::extension_host::ExtensionHostManager>, |
| 220 | _manager: crate::extension_host::TestManagerGuard, |
| 221 | _policy: crate::plugins::activation::TestPolicyGuard, |
| 222 | } |
| 223 | impl HostMcpDispatch { |
| 224 | fn boxed(setup: DispatchSetup) -> Box<dyn McpDispatchUnderTest> { |
| 225 | let node = crate::extension_host::tests::node_for_tests("recorded MCP Host SDK parity") |
| 226 | .expect("recorded Host SDK parity requires a supported local Node runtime"); |
| 227 | let root = |
| 228 | PathBuf::from(std::env::var_os("CODEWHALE_HOME").expect("sealed conformance home")); |
| 229 | let manager = Arc::new(crate::extension_host::ExtensionHostManager::new( |
| 230 | crate::extension_host::ExtensionHostOptions { |
| 231 | node_override: Some(node), |
| 232 | root: Some(root), |
| 233 | ..Default::default() |
| 234 | }, |
| 235 | )); |
| 236 | let policy = crate::plugins::activation::TestPolicyGuard::extension_host(true); |
| 237 | let manager_guard = crate::extension_host::TestManagerGuard::install(Arc::clone(&manager)); |
| 238 | let (tx_event, rx_event) = mpsc::channel(64); |
| 239 | let pool = McpPool::new(setup.config) |
| 240 | .with_disallowed_tools(setup.disallowed_tools.clone()) |
| 241 | .with_backend(crate::mcp::McpBackend::Host); |
| 242 | Box::new(Self { |
| 243 | inner: McpPoolDispatch { |
| 244 | pool: Arc::new(AsyncMutex::new(pool)), |
| 245 | tx_event, |
| 246 | disallowed_tools: setup.disallowed_tools, |
| 247 | _rx_event: Mutex::new(rx_event), |
| 248 | }, |
| 249 | manager, |
| 250 | _manager: manager_guard, |
| 251 | _policy: policy, |
| 252 | }) |
| 253 | } |
| 254 | } |
| 255 | #[async_trait::async_trait] |
| 256 | impl McpDispatchUnderTest for HostMcpDispatch { |
| 257 | async fn boot(&self) -> Result<(), String> { |
| 258 | self.inner.boot().await |
| 259 | } |
| 260 | async fn catalog(&self) -> Vec<codewhale_models::Tool> { |
| 261 | self.inner.catalog().await |
| 262 | } |
| 263 | async fn call( |
| 264 | &self, |
| 265 | name: &str, |
| 266 | input: Value, |
| 267 | cancel: CancellationToken, |
| 268 | ) -> Result<RichToolResult, ToolError> { |
| 269 | self.inner.call(name, input, cancel).await |
| 270 | } |
| 271 | async fn approval_hint(&self, name: &str) -> Option<&'static str> { |
| 272 | self.inner.approval_hint(name).await |
| 273 | } |
| 274 | async fn shutdown(&self) { |
| 275 | self.inner.shutdown().await; |
| 276 | self.manager.shutdown().await; |
| 277 | } |
| 278 | } |
| 279 | |
| 280 | // === The server a case runs against === |
| 281 | |
| 282 | struct Backend { |
| 283 | server: Option<TranscriptServer>, |
| 284 | sink: Option<TranscriptServer>, |
| 285 | stdio_dir: Option<PathBuf>, |
| 286 | } |
| 287 | |
| 288 | impl Backend { |
| 289 | /// Start what the case scripts and return it with the `McpServerConfig` |
| 290 | /// the dispatch connects with. |
| 291 | async fn start(case: &Value, workspace: &Path) -> (Self, Value) { |
| 292 | let mut config = if case["transport"].as_str() == Some("stdio") { |
| 293 | let dir = workspace.join("stdio_server"); |
| 294 | write_stdio_server(&dir, &case["server"]); |
| 295 | let script = dir.join("stdio_server.sh"); |
| 296 | let config = json!({ |
| 297 | "command": "sh", |
| 298 | "args": [script, dir], |
| 299 | "connect_timeout": 10, |
| 300 | "execute_timeout": 10, |
| 301 | }); |
| 302 | return ( |
| 303 | Self { |
| 304 | server: None, |
| 305 | sink: None, |
| 306 | stdio_dir: Some(dir), |
| 307 | }, |
| 308 | merged(config, &case["config"]), |
| 309 | ); |
| 310 | } else { |
| 311 | json!({ "connect_timeout": 10, "execute_timeout": 10 }) |
| 312 | }; |
| 313 | let sink = if case["sink"].as_bool() == Some(true) { |
| 314 | Some(TranscriptServer::start(json!({}), ServerOptions::default()).await) |
| 315 | } else { |
| 316 | None |
| 317 | }; |
| 318 | // The sink's URL is not known until it is listening; the spec names |
| 319 | // it as `{{sink_url}}`. |
| 320 | let spec = match &sink { |
| 321 | Some(sink) => serde_json::from_str( |
| 322 | &case["server"] |
| 323 | .to_string() |
| 324 | .replace("{{sink_url}}", &sink.url), |
| 325 | ) |
| 326 | .expect("spec with sink url"), |
| 327 | None => case["server"].clone(), |
| 328 | }; |
| 329 | let secret = case["bearer_secret"].as_str(); |
| 330 | let options = ServerOptions { |
| 331 | numbered_sessions: case["server_options"]["numbered_sessions"].as_bool() == Some(true), |
| 332 | require_bearer: case["server_options"]["require_bearer"] |
| 333 | .as_bool() |
| 334 | .filter(|required| *required) |
| 335 | .and_then(|_| secret.map(str::to_string)), |
| 336 | }; |
| 337 | let server = TranscriptServer::start(spec, options).await; |
| 338 | config["url"] = json!(server.url); |
| 339 | if secret.is_some() { |
| 340 | config["bearer_token_env_var"] = json!(BEARER_ENV); |
| 341 | } |
| 342 | ( |
| 343 | Self { |
| 344 | server: Some(server), |
| 345 | sink, |
| 346 | stdio_dir: None, |
| 347 | }, |
| 348 | merged(config, &case["config"]), |
| 349 | ) |
| 350 | } |
| 351 | |
| 352 | /// Server URL/port literals the golden masks. |
| 353 | fn literals(&self) -> Vec<(String, String)> { |
| 354 | let mut literals = Vec::new(); |
| 355 | for (server, url_label, addr_label) in [ |
| 356 | (&self.server, "<SERVER_URL>", "<SERVER_ADDR>"), |
| 357 | (&self.sink, "<SINK_URL>", "<SINK_ADDR>"), |
| 358 | ] { |
| 359 | if let Some(server) = server { |
| 360 | literals.push((server.url.clone(), url_label.to_string())); |
| 361 | literals.push((server.addr.clone(), addr_label.to_string())); |
| 362 | } |
| 363 | } |
| 364 | literals |
| 365 | } |
| 366 | |
| 367 | /// Side-effecting requests the server was asked to run. |
| 368 | fn received(&self) -> Vec<Value> { |
| 369 | if let Some(server) = &self.server { |
| 370 | return server.received(); |
| 371 | } |
| 372 | let dir = self.stdio_dir.as_ref().expect("a backend has a server"); |
| 373 | std::fs::read_to_string(dir.join("calls.log")) |
| 374 | .unwrap_or_default() |
| 375 | .lines() |
| 376 | .map(|line| { |
| 377 | let request: Value = serde_json::from_str(line).expect("logged request"); |
| 378 | json!({ |
| 379 | "method": "tools/call", |
| 380 | "name": request["params"]["name"], |
| 381 | "arguments": request["params"].get("arguments").cloned().unwrap_or(Value::Null), |
| 382 | }) |
| 383 | }) |
| 384 | .collect() |
| 385 | } |
| 386 | } |
| 387 | |
| 388 | fn merged(mut base: Value, extra: &Value) -> Value { |
| 389 | if let Some(extra) = extra.as_object() { |
| 390 | for (key, value) in extra { |
| 391 | base[key] = value.clone(); |
| 392 | } |
| 393 | } |
| 394 | base |
| 395 | } |
| 396 | |
| 397 | /// Materialize a transcript for `stdio_server.sh.fixture`: the script itself plus one |
| 398 | /// reply file per scripted answer (see the script's header). |
| 399 | fn write_stdio_server(dir: &Path, spec: &Value) { |
| 400 | std::fs::create_dir_all(dir).expect("stdio server dir"); |
| 401 | std::fs::copy( |
| 402 | golden::family_dir(FAMILY).join("stdio_server.sh.fixture"), |
| 403 | dir.join("stdio_server.sh"), |
| 404 | ) |
| 405 | .expect("copy stdio_server.sh.fixture"); |
| 406 | // `"result":{…}` / `"error":{…}`: a JSON-RPC reply minus its id. |
| 407 | let fragment = |entry: &Value| match entry.get("error") { |
| 408 | Some(error) => format!("\"error\":{error}"), |
| 409 | None => format!("\"result\":{}", entry["result"]), |
| 410 | }; |
| 411 | for (method, file) in [ |
| 412 | ("initialize", "initialize.reply"), |
| 413 | ("tools/list", "tools_list.reply"), |
| 414 | ] { |
| 415 | let entry = spec |
| 416 | .get(method) |
| 417 | .unwrap_or_else(|| panic!("a stdio case scripts `{method}`")); |
| 418 | std::fs::write(dir.join(file), fragment(entry)).expect("write reply"); |
| 419 | } |
| 420 | for call in spec["tools/call"].as_array().into_iter().flatten() { |
| 421 | let tool = call["match"]["name"] |
| 422 | .as_str() |
| 423 | .expect("tools/call match.name"); |
| 424 | if call["exit"].as_bool() == Some(true) { |
| 425 | std::fs::write(dir.join(format!("call.{tool}.exit")), "").expect("write exit"); |
| 426 | } else { |
| 427 | std::fs::write(dir.join(format!("call.{tool}.reply")), fragment(call)) |
| 428 | .expect("write reply"); |
| 429 | } |
| 430 | } |
| 431 | } |
| 432 | |
| 433 | // === Running a transcript === |
| 434 | |
| 435 | #[test] |
| 436 | fn harness_timeout_rejects_a_real_unanswered_mcp_call() { |
| 437 | let mut case = golden::read_case(FAMILY, "tools_resources_prompts"); |
| 438 | let held = case["steps"] |
| 439 | .as_array() |
| 440 | .expect("steps") |
| 441 | .iter() |
| 442 | .find(|step| step["cancel"] == "after_server_holds") |
| 443 | .expect("held-call fixture") |
| 444 | .clone(); |
| 445 | let mut held = held; |
| 446 | held.as_object_mut().expect("step").remove("cancel"); |
| 447 | case["steps"] = json!([held]); |
| 448 | let outcome = run_case( |
| 449 | "unanswered_call", |
| 450 | &case, |
| 451 | McpPoolDispatch::boxed, |
| 452 | Duration::from_secs(1), |
| 453 | true, |
| 454 | ); |
| 455 | let mut failures = Failures::default(); |
| 456 | match outcome { |
| 457 | Err(error) => failures.push("unanswered_call", error), |
| 458 | Ok(()) => panic!("an unanswered MCP call became recordable output"), |
| 459 | } |
| 460 | assert!(failures.contains("harness timeout: MCP call")); |
| 461 | } |
| 462 | |
| 463 | fn normalize_call(result: &Result<RichToolResult, ToolError>) -> Value { |
| 464 | match result { |
| 465 | Ok(rich) => json!({ "ok": { |
| 466 | "success": rich.result.success, |
| 467 | "content": serde_json::from_str::<Value>(&rich.result.content) |
| 468 | .unwrap_or_else(|_| Value::String(rich.result.content.clone())), |
| 469 | "metadata": rich.result.metadata, |
| 470 | "content_blocks": rich.content_blocks, |
| 471 | } }), |
| 472 | Err(error) => { |
| 473 | json!({ "err": { "kind": golden::tool_error_kind(error), "detail": error.to_string() } }) |
| 474 | } |
| 475 | } |
| 476 | } |
| 477 | |
| 478 | async fn run_transcript( |
| 479 | case: &Value, |
| 480 | factory: DispatchFactory, |
| 481 | workspace: &Path, |
| 482 | deadline: Duration, |
| 483 | ) -> Result<(Value, Vec<(String, String)>), String> { |
| 484 | let scripted_steps = case["steps"].as_array().expect("case.steps"); |
| 485 | if scripted_steps.is_empty() { |
| 486 | return Err("MCP transcript has no host steps".to_string()); |
| 487 | } |
| 488 | let (backend, server_config) = Backend::start(case, workspace).await; |
| 489 | let server_name = case["server_name"].as_str().unwrap_or("conformance"); |
| 490 | let config: McpConfig = |
| 491 | serde_json::from_value(json!({ "servers": { server_name: server_config } })) |
| 492 | .expect("transcript MCP config"); |
| 493 | let dispatch = factory(DispatchSetup { |
| 494 | config, |
| 495 | disallowed_tools: case["disallowed_tools"] |
| 496 | .as_array() |
| 497 | .into_iter() |
| 498 | .flatten() |
| 499 | .map(|rule| rule.as_str().expect("deny rule").to_string()) |
| 500 | .collect(), |
| 501 | }); |
| 502 | let mut steps = Vec::new(); |
| 503 | let expect_boot_error = case["expect_boot"].as_str() == Some("error"); |
| 504 | let boot = match golden::complete_within("MCP boot", deadline, dispatch.boot()).await? { |
| 505 | Ok(()) if expect_boot_error => { |
| 506 | return Err("MCP transcript expected boot to fail but it connected".to_string()); |
| 507 | } |
| 508 | Ok(()) => json!("ok"), |
| 509 | Err(detail) if expect_boot_error => json!({ "err": detail }), |
| 510 | Err(detail) => return Err(format!("MCP transcript could not boot: {detail}")), |
| 511 | }; |
| 512 | steps.push(json!({ "op": "boot", "outcome": boot })); |
| 513 | for step in scripted_steps { |
| 514 | match step["op"].as_str().expect("step.op") { |
| 515 | "catalog" => { |
| 516 | let catalog = serde_json::to_value( |
| 517 | golden::complete_within("MCP catalog", deadline, dispatch.catalog()).await?, |
| 518 | ) |
| 519 | .expect("catalog"); |
| 520 | steps.push(json!({ "op": "catalog", "tools": catalog })); |
| 521 | } |
| 522 | "approval_hints" => { |
| 523 | let mut hints = serde_json::Map::new(); |
| 524 | for tool in step["tools"].as_array().expect("step.tools") { |
| 525 | let tool = tool.as_str().expect("tool name"); |
| 526 | let hint = golden::complete_within( |
| 527 | "MCP approval hint", |
| 528 | deadline, |
| 529 | dispatch.approval_hint(tool), |
| 530 | ) |
| 531 | .await?; |
| 532 | hints.insert(tool.to_string(), json!(hint)); |
| 533 | } |
| 534 | steps.push(json!({ "op": "approval_hints", "hints": hints })); |
| 535 | } |
| 536 | "call" => { |
| 537 | let tool = step["tool"].as_str().expect("step.tool"); |
| 538 | let cancel = CancellationToken::new(); |
| 539 | let canceller = |
| 540 | (step["cancel"].as_str() == Some("after_server_holds")).then(|| { |
| 541 | let state = Arc::clone( |
| 542 | &backend |
| 543 | .server |
| 544 | .as_ref() |
| 545 | .expect("only an HTTP case can hold a call") |
| 546 | .state, |
| 547 | ); |
| 548 | let cancel = cancel.clone(); |
| 549 | tokio::spawn(async move { |
| 550 | state.held.notified().await; |
| 551 | cancel.cancel(); |
| 552 | }) |
| 553 | }); |
| 554 | let outcome = golden::complete_within( |
| 555 | "MCP call", |
| 556 | deadline, |
| 557 | dispatch.call(tool, step["input"].clone(), cancel), |
| 558 | ) |
| 559 | .await; |
| 560 | if let Some(canceller) = canceller { |
| 561 | canceller.abort(); |
| 562 | } |
| 563 | let outcome = normalize_call(&outcome?); |
| 564 | steps.push(json!({ "op": "call", "tool": tool, "outcome": outcome })); |
| 565 | } |
| 566 | other => panic!("unknown transcript op `{other}`"), |
| 567 | } |
| 568 | } |
| 569 | golden::complete_within("MCP shutdown", deadline, dispatch.shutdown()).await?; |
| 570 | let mut outcome = json!({ "steps": steps, "server_received": backend.received() }); |
| 571 | let requests = backend |
| 572 | .server |
| 573 | .as_ref() |
| 574 | .map(TranscriptServer::requests) |
| 575 | .unwrap_or_default(); |
| 576 | match case["record_requests"].as_str() { |
| 577 | Some("full") => outcome["requests"] = Value::Array(requests), |
| 578 | Some("summary") => outcome["requests_summary"] = summarize_requests(&requests), |
| 579 | Some(other) => panic!("unknown record_requests `{other}`"), |
| 580 | None => {} |
| 581 | } |
| 582 | if let Some(sink) = &backend.sink { |
| 583 | outcome["sink_requests"] = json!(sink.requests().len()); |
| 584 | } |
| 585 | Ok((outcome, backend.literals())) |
| 586 | } |
| 587 | |
| 588 | /// The distinct RPC methods, `Authorization` shapes and HTTP statuses the |
| 589 | /// server saw, without their order or count. For cases where *that* a request |
| 590 | /// carried a credential (or never got past `initialize`) is the contract, but |
| 591 | /// how often a client retries a connect is its own business. |
| 592 | fn summarize_requests(requests: &[Value]) -> Value { |
| 593 | let distinct = |field: &str| { |
| 594 | let mut values: Vec<Value> = Vec::new(); |
| 595 | for request in requests { |
| 596 | if !values.contains(&request[field]) { |
| 597 | values.push(request[field].clone()); |
| 598 | } |
| 599 | } |
| 600 | values.sort_by_key(Value::to_string); |
| 601 | Value::Array(values) |
| 602 | }; |
| 603 | json!({ |
| 604 | "rpcs": distinct("rpc"), |
| 605 | "authorization_shapes": distinct("authorization"), |
| 606 | "statuses": distinct("status"), |
| 607 | }) |
| 608 | } |
| 609 | |
| 610 | /// No recorded byte may carry a credential. The needle is the secret itself, |
| 611 | /// so a bearer header, a URL, a JSON string or an error detail all trip it. |
| 612 | fn assert_no_secret(secret: &str, recorded: &str) -> Result<(), String> { |
| 613 | if recorded.contains(secret) { |
| 614 | return Err(format!( |
| 615 | "secret leak: the case's bearer secret appears in the recorded output:\n{}", |
| 616 | recorded |
| 617 | .lines() |
| 618 | .filter(|line| line.contains(secret)) |
| 619 | .collect::<Vec<_>>() |
| 620 | .join("\n") |
| 621 | )); |
| 622 | } |
| 623 | Ok(()) |
| 624 | } |
| 625 | |
| 626 | /// Run one case through one dispatch and compare with its golden. |
| 627 | /// |
| 628 | /// `compare_only` forces compare mode even under `CODEWHALE_CONFORMANCE_UPDATE`: |
| 629 | /// the negative controls feed deliberately changed cases through here and must |
| 630 | /// never rewrite a golden. |
| 631 | fn run_case( |
| 632 | name: &str, |
| 633 | case: &Value, |
| 634 | factory: DispatchFactory, |
| 635 | deadline: Duration, |
| 636 | compare_only: bool, |
| 637 | ) -> Result<(), String> { |
| 638 | if case["transport"].as_str() == Some("stdio") && !cfg!(unix) { |
| 639 | eprintln!("conformance: `{name}` needs a POSIX shell; skipped on this platform"); |
| 640 | return Ok(()); |
| 641 | } |
| 642 | // The pool may touch its state directory; keep that hermetic. Guards are |
| 643 | // declared after the sandbox so they restore the environment before the |
| 644 | // sandbox releases the lock. |
| 645 | let sandbox = Sandbox::new(&Value::Null); |
| 646 | let _compare_only = compare_only.then(|| EnvVarGuard::set(golden::UPDATE_ENV, "0")); |
| 647 | let secret = case["bearer_secret"].as_str(); |
| 648 | let _bearer = secret.map(|secret| EnvVarGuard::set(BEARER_ENV, secret)); |
| 649 | let runtime = tokio::runtime::Builder::new_current_thread() |
| 650 | .enable_all() |
| 651 | .build() |
| 652 | .expect("runtime"); |
| 653 | let (mut outcome, literals) = |
| 654 | runtime.block_on(run_transcript(case, factory, &sandbox.workspace, deadline))?; |
| 655 | drop(runtime); |
| 656 | let mut masker = sandbox.masker(&[]); |
| 657 | for (literal, label) in literals { |
| 658 | masker = masker.literal(&literal, &label); |
| 659 | } |
| 660 | masker.value(&mut outcome); |
| 661 | let recorded = golden::pretty(&golden::canonical(&outcome)); |
| 662 | if let Some(secret) = secret { |
| 663 | assert_no_secret(secret, &recorded)?; |
| 664 | } |
| 665 | golden::check_golden( |
| 666 | &golden::family_dir(FAMILY).join(format!("{name}.golden.json")), |
| 667 | &recorded, |
| 668 | ) |
| 669 | } |
| 670 | |
| 671 | #[test] |
| 672 | fn mcp_transcripts_match_goldens_for_every_dispatch() { |
| 673 | let names = golden::case_names(FAMILY); |
| 674 | let mut failures = Failures::default(); |
| 675 | for name in &names { |
| 676 | let case = golden::read_case(FAMILY, name); |
| 677 | for (dispatch, factory) in DISPATCHES { |
| 678 | failures.record( |
| 679 | &format!("{name} via {dispatch}"), |
| 680 | run_case(name, &case, *factory, CALL_DEADLINE, false), |
| 681 | ); |
| 682 | } |
| 683 | } |
| 684 | failures.finish(FAMILY, names.len() * DISPATCHES.len()); |
| 685 | } |
| 686 |