返回 CodeWhale
mcp.rs
根目录 / crates / tui / src / conformance / mcp.rs
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
686 lines RUST