1use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
2use std::path::PathBuf;
3use std::sync::atomic::{AtomicU64, Ordering};
4use std::sync::{Arc, Mutex as StdMutex};
5use std::time::Duration;
6
7use async_trait::async_trait;
8use serde_json::{json, Value};
9use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
10use tokio::process::{Child, ChildStderr, ChildStdin, ChildStdout, Command};
11use tokio::time::timeout;
12
13use orchestral_core::agent_protocol::wire::{Digest, ToolActivityEvidence};
14use orchestral_core::mcp_protocol::{
15 McpProtocolEra, McpServerId, McpServerSnapshot, McpToolAnnotations, McpToolSnapshot,
16 McpTransportAuthority, McpTransportCancellation, McpTransportConnection, McpTransportError,
17 McpTransportFactory, McpTransportKind, McpTransportRequest, MCP_LATEST_LEGACY_PROTOCOL,
18 MCP_STATELESS_PROTOCOL_2026_07_28,
19};
20use orchestral_core::tool_protocol::{
21 ApprovalCapabilityStore, CapabilityRequest, CapabilitySelector, EffectScope, ModelToolSchema,
22 ToolConcurrency, ToolDescriptor, ToolId, ToolIdempotency, ToolInvocation, ToolOperationPlan,
23 ToolOperationRisk, ToolOutcome, ToolRestriction,
24};
25use tokio_util::sync::CancellationToken;
26
27use crate::tool_runtime::{
28 GuardedToolExecution, GuardedToolExecutor, GuardedToolRuntime, ToolRuntimeError,
29};
30use crate::tools::shell_sandbox::{
31 normalize_network_targets, sandbox_command, SandboxNetworkAccess, ShellSandboxPolicy,
32};
33
34const DEFAULT_MCP_MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
35const MAX_SAFE_MCP_HEADER_INTEGER: u64 = 9_007_199_254_740_991;
36const MCP_TRANSPORT_CLOSE_TIMEOUT: Duration = Duration::from_millis(750);
37const MCP_PROCESS_REAP_TIMEOUT: Duration = Duration::from_millis(500);
38const MCP_STDERR_CAPTURE_BYTES: usize = 16 * 1024;
39pub const MCP_STDIO_SANDBOX_PROFILE: &str = "orchestral.mcp.stdio.v1";
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq)]
42enum McpHeaderValueKind {
43 String,
44 Integer,
45 Boolean,
46}
47
48#[derive(Debug, Clone, PartialEq, Eq)]
49struct McpHeaderBinding {
50 property_path: Vec<String>,
51 header_suffix: String,
52 value_kind: McpHeaderValueKind,
53}
54
55fn compile_mcp_header_bindings(schema: &Value) -> Result<Vec<McpHeaderBinding>, String> {
56 if !schema.is_object() {
57 return Err("MCP Tool inputSchema must be an object".to_owned());
58 }
59 let mut bindings = Vec::new();
60 let mut header_names = BTreeSet::new();
61 scan_mcp_header_annotations(schema, &[], true, &mut header_names, &mut bindings)?;
62 bindings.sort_by(|left, right| {
63 left.header_suffix
64 .to_ascii_lowercase()
65 .cmp(&right.header_suffix.to_ascii_lowercase())
66 .then_with(|| left.property_path.cmp(&right.property_path))
67 });
68 Ok(bindings)
69}
70
71fn scan_mcp_header_annotations(
72 value: &Value,
73 property_path: &[String],
74 statically_reachable: bool,
75 header_names: &mut BTreeSet<String>,
76 bindings: &mut Vec<McpHeaderBinding>,
77) -> Result<(), String> {
78 match value {
79 Value::Object(object) => {
80 if let Some(annotation) = object.get("x-mcp-header") {
81 if !statically_reachable || property_path.is_empty() {
82 return Err(
83 "x-mcp-header is not on a property statically reachable from the schema root"
84 .to_owned(),
85 );
86 }
87 let header_suffix = annotation.as_str().ok_or_else(|| {
88 "x-mcp-header must be a non-empty HTTP field-name token".to_owned()
89 })?;
90 if !is_mcp_http_token(header_suffix) {
91 return Err(format!(
92 "x-mcp-header '{header_suffix}' is not an HTTP field-name token"
93 ));
94 }
95 if !header_names.insert(header_suffix.to_ascii_lowercase()) {
96 return Err(format!(
97 "x-mcp-header '{header_suffix}' is not case-insensitively unique"
98 ));
99 }
100 let value_kind = match object.get("type").and_then(Value::as_str) {
101 Some("string") => McpHeaderValueKind::String,
102 Some("integer") => McpHeaderValueKind::Integer,
103 Some("boolean") => McpHeaderValueKind::Boolean,
104 _ => {
105 return Err(format!(
106 "x-mcp-header '{header_suffix}' must annotate string, integer, or boolean"
107 ))
108 }
109 };
110 bindings.push(McpHeaderBinding {
111 property_path: property_path.to_vec(),
112 header_suffix: header_suffix.to_owned(),
113 value_kind,
114 });
115 }
116
117 for (keyword, child) in object {
118 if keyword == "properties" && statically_reachable {
119 let properties = child.as_object().ok_or_else(|| {
120 "JSON Schema properties containing x-mcp-header must be an object"
121 .to_owned()
122 })?;
123 for (name, property_schema) in properties {
124 let mut child_path = property_path.to_vec();
125 child_path.push(name.clone());
126 scan_mcp_header_annotations(
127 property_schema,
128 &child_path,
129 true,
130 header_names,
131 bindings,
132 )?;
133 }
134 } else if keyword != "x-mcp-header" {
135 scan_mcp_header_annotations(
136 child,
137 property_path,
138 false,
139 header_names,
140 bindings,
141 )?;
142 }
143 }
144 }
145 Value::Array(values) => {
146 for child in values {
147 scan_mcp_header_annotations(child, property_path, false, header_names, bindings)?;
148 }
149 }
150 _ => {}
151 }
152 Ok(())
153}
154
155fn is_mcp_http_token(value: &str) -> bool {
156 !value.is_empty()
157 && value.bytes().all(|byte| {
158 byte.is_ascii_alphanumeric()
159 || matches!(
160 byte,
161 b'!' | b'#'
162 | b'$'
163 | b'%'
164 | b'&'
165 | b'\''
166 | b'*'
167 | b'+'
168 | b'-'
169 | b'.'
170 | b'^'
171 | b'_'
172 | b'`'
173 | b'|'
174 | b'~'
175 )
176 })
177}
178
179fn build_mcp_parameter_headers(
180 bindings: &[McpHeaderBinding],
181 arguments: &Value,
182) -> Result<BTreeMap<String, String>, String> {
183 let mut headers = BTreeMap::new();
184 for binding in bindings {
185 let mut value = arguments;
186 let mut missing = false;
187 for segment in &binding.property_path {
188 let Some(next) = value.as_object().and_then(|object| object.get(segment)) else {
189 missing = true;
190 break;
191 };
192 value = next;
193 }
194 if missing || value.is_null() {
195 continue;
196 }
197 let value = match binding.value_kind {
198 McpHeaderValueKind::String => value
199 .as_str()
200 .map(str::to_owned)
201 .ok_or_else(|| invalid_mcp_header_value(binding))?,
202 McpHeaderValueKind::Boolean => value
203 .as_bool()
204 .map(|value| value.to_string())
205 .ok_or_else(|| invalid_mcp_header_value(binding))?,
206 McpHeaderValueKind::Integer => canonical_safe_mcp_integer(value)
207 .ok_or_else(|| invalid_mcp_header_value(binding))?,
208 };
209 headers.insert(binding.header_suffix.clone(), value);
210 }
211 Ok(headers)
212}
213
214fn canonical_safe_mcp_integer(value: &Value) -> Option<String> {
215 if let Some(value) = value.as_i64() {
216 if value.unsigned_abs() <= MAX_SAFE_MCP_HEADER_INTEGER {
217 return Some(value.to_string());
218 }
219 } else if let Some(value) = value.as_u64() {
220 if value <= MAX_SAFE_MCP_HEADER_INTEGER {
221 return Some(value.to_string());
222 }
223 }
224 None
225}
226
227fn invalid_mcp_header_value(binding: &McpHeaderBinding) -> String {
228 format!(
229 "MCP Tool argument at '{}' does not match x-mcp-header '{}' primitive type or safe integer range",
230 binding.property_path.join("."),
231 binding.header_suffix
232 )
233}
234
235#[derive(Debug, Clone)]
237pub struct StdioMcpSandboxPolicy {
238 cwd: PathBuf,
239 readable_roots: BTreeSet<PathBuf>,
240 writable_roots: BTreeSet<PathBuf>,
241 network_targets: BTreeSet<String>,
242 allow_unrestricted_network: bool,
243 private_runtime_home: Option<PathBuf>,
244 allow_child_processes: bool,
245 allow_host_ui: bool,
246}
247
248impl StdioMcpSandboxPolicy {
249 pub fn workspace(root: impl Into<PathBuf>) -> Self {
250 let root = root.into();
251 Self {
252 cwd: root.clone(),
253 readable_roots: BTreeSet::from([root.clone()]),
254 writable_roots: BTreeSet::from([root]),
255 network_targets: BTreeSet::new(),
256 allow_unrestricted_network: false,
257 private_runtime_home: None,
258 allow_child_processes: false,
259 allow_host_ui: false,
260 }
261 }
262
263 pub fn scoped(
267 cwd: impl Into<PathBuf>,
268 readable_roots: BTreeSet<PathBuf>,
269 writable_roots: BTreeSet<PathBuf>,
270 network_targets: BTreeSet<String>,
271 ) -> Self {
272 Self {
273 cwd: cwd.into(),
274 readable_roots,
275 writable_roots,
276 network_targets,
277 allow_unrestricted_network: false,
278 private_runtime_home: None,
279 allow_child_processes: false,
280 allow_host_ui: false,
281 }
282 }
283
284 pub fn with_child_processes(mut self, allow: bool) -> Self {
287 self.allow_child_processes = allow;
288 self
289 }
290
291 pub fn with_host_ui(mut self, allow: bool) -> Self {
295 self.allow_host_ui = allow;
296 self
297 }
298
299 pub fn with_unrestricted_network(mut self, allow: bool) -> Self {
303 self.allow_unrestricted_network = allow;
304 self
305 }
306
307 pub fn with_private_runtime_home(mut self, root: impl Into<PathBuf>) -> Self {
311 self.private_runtime_home = Some(root.into());
312 self
313 }
314
315 fn normalize(self) -> Result<Self, McpToolsAdapterError> {
316 let cwd = canonical_mcp_directory(&self.cwd, "cwd")?;
317 let mut readable_roots = self
318 .readable_roots
319 .iter()
320 .map(|root| canonical_mcp_directory(root, "readable root"))
321 .collect::<Result<BTreeSet<_>, _>>()?;
322 let writable_roots = self
323 .writable_roots
324 .iter()
325 .map(|root| canonical_mcp_directory(root, "writable root"))
326 .collect::<Result<BTreeSet<_>, _>>()?;
327 let network_targets = normalize_network_targets(&self.network_targets)
328 .map_err(McpToolsAdapterError::InvalidConfig)?;
329 if self.allow_unrestricted_network && !network_targets.is_empty() {
330 return Err(McpToolsAdapterError::InvalidConfig(
331 "MCP stdio network authority must be exact targets or unrestricted, not both"
332 .to_owned(),
333 ));
334 }
335 if self.allow_host_ui && !self.allow_child_processes {
336 return Err(McpToolsAdapterError::InvalidConfig(
337 "MCP stdio Host UI access requires child-process authority".to_owned(),
338 ));
339 }
340 let private_runtime_home = self
341 .private_runtime_home
342 .as_deref()
343 .map(|root| canonical_mcp_directory(root, "private runtime home"))
344 .transpose()?;
345 if let Some(home) = &private_runtime_home {
346 readable_roots.insert(home.clone());
347 }
348 if readable_roots.is_empty()
349 || writable_roots.is_empty()
350 || !readable_roots.iter().any(|root| cwd.starts_with(root))
351 || private_runtime_home.as_ref().is_some_and(|home| {
352 !writable_roots
353 .iter()
354 .any(|writable| home == writable || home.starts_with(writable))
355 })
356 {
357 return Err(McpToolsAdapterError::InvalidConfig(
358 "MCP stdio sandbox requires readable/writable roots and a readable cwd".to_owned(),
359 ));
360 }
361 Ok(Self {
362 cwd,
363 readable_roots,
364 writable_roots,
365 network_targets,
366 allow_unrestricted_network: self.allow_unrestricted_network,
367 private_runtime_home,
368 allow_child_processes: self.allow_child_processes,
369 allow_host_ui: self.allow_host_ui,
370 })
371 }
372}
373
374fn canonical_mcp_directory(
375 path: &std::path::Path,
376 label: &str,
377) -> Result<PathBuf, McpToolsAdapterError> {
378 let canonical = std::fs::canonicalize(path).map_err(|error| {
379 McpToolsAdapterError::InvalidConfig(format!(
380 "canonicalize MCP stdio sandbox {label} '{}' failed: {error}",
381 path.display()
382 ))
383 })?;
384 if !canonical.is_dir() {
385 return Err(McpToolsAdapterError::InvalidConfig(format!(
386 "MCP stdio sandbox {label} '{}' is not a directory",
387 path.display()
388 )));
389 }
390 Ok(canonical)
391}
392
393#[derive(Debug, Clone)]
396pub struct StdioMcpTransportFactory {
397 program: PathBuf,
398 args: Vec<String>,
399 environment: BTreeMap<String, String>,
400 sandbox: StdioMcpSandboxPolicy,
401 authority: McpTransportAuthority,
402}
403
404impl StdioMcpTransportFactory {
405 pub fn new(
406 program: PathBuf,
407 args: Vec<String>,
408 environment: BTreeMap<String, String>,
409 sandbox: StdioMcpSandboxPolicy,
410 ) -> Result<Self, McpToolsAdapterError> {
411 let canonical_program = program.canonicalize().ok();
412 if !program.is_absolute()
413 || !program.is_file()
414 || canonical_program.as_ref() != Some(&program)
415 || environment.keys().any(|name| {
416 name.trim().is_empty() || name.contains('=') || name.chars().any(char::is_control)
417 })
418 {
419 return Err(McpToolsAdapterError::InvalidConfig(
420 "invalid guarded MCP stdio transport".to_owned(),
421 ));
422 }
423 let sandbox = sandbox.normalize()?;
424 let mut effect_scopes = BTreeSet::from([
425 EffectScope::Process,
426 EffectScope::FilesystemRead,
427 EffectScope::FilesystemWrite,
428 EffectScope::ExternalSideEffect,
429 ]);
430 if !environment.is_empty() {
431 effect_scopes.insert(EffectScope::SecretRead);
432 }
433 if sandbox.allow_unrestricted_network || !sandbox.network_targets.is_empty() {
434 effect_scopes.insert(EffectScope::Network);
435 }
436 let binding = json!({
437 "transport": "stdio",
438 "program": program.to_string_lossy(),
439 "args": args,
440 "environmentNames": environment.keys().collect::<Vec<_>>(),
441 "cwd": sandbox.cwd.to_string_lossy(),
442 "readableRoots": sandbox.readable_roots.iter().map(|root| root.to_string_lossy()).collect::<Vec<_>>(),
443 "writableRoots": sandbox.writable_roots.iter().map(|root| root.to_string_lossy()).collect::<Vec<_>>(),
444 "networkTargets": &sandbox.network_targets,
445 "allowUnrestrictedNetwork": sandbox.allow_unrestricted_network,
446 "privateRuntimeHome": sandbox.private_runtime_home.as_ref().map(|path| path.to_string_lossy()),
447 "allowChildProcesses": sandbox.allow_child_processes,
448 "allowHostUi": sandbox.allow_host_ui,
449 "sandboxProfile": MCP_STDIO_SANDBOX_PROFILE,
450 "maxFrameBytes": DEFAULT_MCP_MAX_FRAME_BYTES,
451 });
452 let binding_digest = Digest::sha256(
453 serde_jcs::to_vec(&binding)
454 .map_err(|error| McpToolsAdapterError::InvalidConfig(error.to_string()))?,
455 );
456 let authority = McpTransportAuthority {
457 kind: McpTransportKind::Stdio,
458 binding_digest,
459 effect_scopes,
460 process_programs: BTreeSet::from([program.to_string_lossy().to_string()]),
461 allow_child_processes: sandbox.allow_child_processes,
462 allow_host_ui: sandbox.allow_host_ui,
463 filesystem_read_roots: sandbox
464 .readable_roots
465 .iter()
466 .map(|root| root.to_string_lossy().into_owned())
467 .collect(),
468 filesystem_write_roots: sandbox
469 .writable_roots
470 .iter()
471 .map(|root| root.to_string_lossy().into_owned())
472 .collect(),
473 sandbox_profiles: BTreeSet::from([MCP_STDIO_SANDBOX_PROFILE.to_owned()]),
474 network_targets: sandbox.network_targets.clone(),
475 allow_unrestricted_network: sandbox.allow_unrestricted_network,
476 environment_variables: environment.keys().cloned().collect(),
477 credential_references: BTreeSet::new(),
478 };
479 authority
480 .validate()
481 .map_err(|error| McpToolsAdapterError::InvalidConfig(error.to_string()))?;
482 Ok(Self {
483 program,
484 args,
485 environment,
486 sandbox,
487 authority,
488 })
489 }
490}
491
492#[async_trait]
493impl McpTransportFactory for StdioMcpTransportFactory {
494 fn authority(&self) -> &McpTransportAuthority {
495 &self.authority
496 }
497
498 async fn connect(&self) -> Result<Box<dyn McpTransportConnection>, McpTransportError> {
499 let command = sandbox_command(
500 self.program.to_string_lossy().into_owned(),
501 self.args.clone(),
502 &self.sandbox.cwd,
503 &ShellSandboxPolicy {
504 readable_roots: self.sandbox.readable_roots.iter().cloned().collect(),
505 readable_files: Vec::new(),
506 writable_roots: self.sandbox.writable_roots.iter().cloned().collect(),
507 allow_child_processes: self.sandbox.allow_child_processes,
508 allow_host_ui: self.sandbox.allow_host_ui,
509 launcher_programs: vec![self.program.clone()],
510 network: if self.sandbox.allow_unrestricted_network {
511 SandboxNetworkAccess::Unrestricted
512 } else if self.sandbox.network_targets.is_empty() {
513 SandboxNetworkAccess::Disabled
514 } else {
515 SandboxNetworkAccess::ExactTargets(self.sandbox.network_targets.clone())
516 },
517 linux_bwrap_path: None,
518 },
519 )
520 .map_err(McpTransportError::Transport)?;
521 let mut environment = self
522 .environment
523 .iter()
524 .map(|(key, value)| (key.clone(), value.clone()))
525 .collect::<HashMap<_, _>>();
526 if let Some(home) = &self.sandbox.private_runtime_home {
527 let home = home.to_string_lossy().into_owned();
528 environment
529 .entry("HOME".to_owned())
530 .or_insert_with(|| home.clone());
531 environment
532 .entry("USERPROFILE".to_owned())
533 .or_insert_with(|| home.clone());
534 environment.entry("TMPDIR".to_owned()).or_insert(home);
535 }
536 environment.extend(command.env);
537 StdioMcpTransport::connect(
538 &command.program,
539 &command.args,
540 &environment,
541 &self.sandbox.cwd,
542 command.backend_starts_new_session,
543 )
544 .await
545 .map(|transport| Box::new(transport) as Box<dyn McpTransportConnection>)
546 .map_err(McpTransportError::Transport)
547 }
548}
549
550#[derive(Clone)]
553pub struct GuardedMcpServerConfig {
554 pub server_id: McpServerId,
555 pub required: bool,
556 pub transport: Arc<dyn McpTransportFactory>,
557 pub startup_timeout: Duration,
558 pub tool_timeout: Duration,
559 pub enabled_tools: BTreeSet<String>,
560 pub disabled_tools: BTreeSet<String>,
561}
562
563impl std::fmt::Debug for GuardedMcpServerConfig {
564 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
565 formatter
566 .debug_struct("GuardedMcpServerConfig")
567 .field("server_id", &self.server_id)
568 .field("required", &self.required)
569 .field("transport_authority", self.transport.authority())
570 .field("startup_timeout", &self.startup_timeout)
571 .field("tool_timeout", &self.tool_timeout)
572 .field("enabled_tools", &self.enabled_tools)
573 .field("disabled_tools", &self.disabled_tools)
574 .finish()
575 }
576}
577
578impl GuardedMcpServerConfig {
579 pub fn validate(&self) -> Result<(), McpToolsAdapterError> {
580 if self.server_id.is_empty()
581 || self.startup_timeout.is_zero()
582 || self.tool_timeout.is_zero()
583 || self
584 .enabled_tools
585 .iter()
586 .chain(self.disabled_tools.iter())
587 .any(|name| name.trim().is_empty())
588 {
589 return Err(McpToolsAdapterError::InvalidConfig(format!(
590 "invalid guarded MCP configuration for '{}'",
591 self.server_id
592 )));
593 }
594 self.transport.authority().validate().map_err(|error| {
595 McpToolsAdapterError::InvalidConfig(format!(
596 "invalid MCP transport authority for '{}': {error}",
597 self.server_id
598 ))
599 })?;
600 Ok(())
601 }
602
603 fn allows_tool(&self, name: &str) -> bool {
604 !name.trim().is_empty()
605 && (self.enabled_tools.is_empty() || self.enabled_tools.contains(name))
606 && !self.disabled_tools.contains(name)
607 }
608
609 pub fn effect_scopes(&self) -> BTreeSet<EffectScope> {
610 self.transport.authority().effect_scopes.clone()
611 }
612
613 pub fn allowed_programs(&self) -> BTreeSet<String> {
614 self.transport.authority().process_programs.clone()
615 }
616
617 pub fn allows_child_processes(&self) -> bool {
618 self.transport.authority().allow_child_processes
619 }
620
621 pub fn allowed_network_targets(&self) -> BTreeSet<String> {
622 self.transport.authority().network_targets.clone()
623 }
624
625 pub fn allows_unrestricted_network(&self) -> bool {
626 self.transport.authority().allow_unrestricted_network
627 }
628
629 pub fn filesystem_read_roots(&self) -> BTreeSet<String> {
630 self.transport.authority().filesystem_read_roots.clone()
631 }
632
633 pub fn filesystem_write_roots(&self) -> BTreeSet<String> {
634 self.transport.authority().filesystem_write_roots.clone()
635 }
636
637 pub fn sandbox_profiles(&self) -> BTreeSet<String> {
638 self.transport.authority().sandbox_profiles.clone()
639 }
640
641 pub fn environment_names(&self) -> BTreeSet<String> {
642 self.transport.authority().environment_variables.clone()
643 }
644
645 pub fn credential_references(&self) -> BTreeSet<String> {
646 self.transport.authority().credential_references.clone()
647 }
648
649 fn transport_kind(&self) -> McpTransportKind {
650 self.transport.authority().kind
651 }
652}
653
654#[derive(Debug, Clone, Copy, PartialEq, Eq)]
655#[non_exhaustive]
656pub enum McpServerHealth {
657 Connecting,
658 Ready,
659 Degraded,
660 Closed,
661}
662
663impl McpServerHealth {
664 fn encode(self) -> u64 {
665 match self {
666 Self::Connecting => 0,
667 Self::Ready => 1,
668 Self::Degraded => 2,
669 Self::Closed => 3,
670 }
671 }
672
673 fn decode(value: u64) -> Self {
674 match value {
675 1 => Self::Ready,
676 2 => Self::Degraded,
677 3 => Self::Closed,
678 _ => Self::Connecting,
679 }
680 }
681}
682
683pub struct McpServerConnectionManager {
685 config: GuardedMcpServerConfig,
686 snapshot: McpServerSnapshot,
687 parameter_headers: BTreeMap<String, Vec<McpHeaderBinding>>,
688 session: tokio::sync::Mutex<Option<McpTransportSession>>,
689 connection_generation: AtomicU64,
690 health: AtomicU64,
691}
692
693impl McpServerConnectionManager {
694 async fn connect(
695 config: GuardedMcpServerConfig,
696 cancellation: CancellationToken,
697 ) -> Result<Arc<Self>, McpToolsAdapterError> {
698 config.validate()?;
699 let (session, snapshot, parameter_headers) =
700 discover_guarded_mcp_session(&config, &cancellation).await?;
701 Ok(Arc::new(Self {
702 config,
703 snapshot,
704 parameter_headers,
705 session: tokio::sync::Mutex::new(Some(session)),
706 connection_generation: AtomicU64::new(1),
707 health: AtomicU64::new(McpServerHealth::Ready.encode()),
708 }))
709 }
710
711 pub fn config(&self) -> &GuardedMcpServerConfig {
712 &self.config
713 }
714
715 pub fn snapshot(&self) -> &McpServerSnapshot {
716 &self.snapshot
717 }
718
719 pub fn connection_generation(&self) -> u64 {
720 self.connection_generation.load(Ordering::Acquire)
721 }
722
723 pub fn health(&self) -> McpServerHealth {
724 McpServerHealth::decode(self.health.load(Ordering::Acquire))
725 }
726
727 async fn invoke(
728 &self,
729 tool_name: &str,
730 arguments: Value,
731 cancellation: CancellationToken,
732 ) -> Result<Value, GuardedMcpCallError> {
733 if !self.config.allows_tool(tool_name)
734 || !self
735 .snapshot
736 .tools
737 .iter()
738 .any(|tool| tool.name == tool_name)
739 {
740 return Err(GuardedMcpCallError::Rejected(format!(
741 "MCP Tool '{tool_name}' is outside the pinned Host snapshot/filter"
742 )));
743 }
744 let parameter_headers = build_mcp_parameter_headers(
745 self.parameter_headers
746 .get(tool_name)
747 .map(Vec::as_slice)
748 .unwrap_or_default(),
749 &arguments,
750 )
751 .map_err(GuardedMcpCallError::Rejected)?;
752 let mut guard = tokio::select! {
753 biased;
754 _ = cancellation.cancelled() => return Err(GuardedMcpCallError::Cancelled),
755 guard = self.session.lock() => guard,
756 };
757 if guard.is_none() {
758 self.health
759 .store(McpServerHealth::Connecting.encode(), Ordering::Release);
760 let (mut session, current_snapshot, current_parameter_headers) =
761 discover_guarded_mcp_session(&self.config, &cancellation)
762 .await
763 .map_err(|error| GuardedMcpCallError::Failed(error.to_string()))?;
764 if current_snapshot.revision != self.snapshot.revision
765 || current_parameter_headers != self.parameter_headers
766 {
767 let _ = session.shutdown().await;
768 self.health
769 .store(McpServerHealth::Degraded.encode(), Ordering::Release);
770 return Err(GuardedMcpCallError::Failed(
771 "MCP Tool catalog changed after reconnect; the pinned Host snapshot is stale"
772 .to_owned(),
773 ));
774 }
775 self.connection_generation.fetch_add(1, Ordering::AcqRel);
776 *guard = Some(session);
777 }
778 let request = {
779 let session = guard.as_mut().expect("MCP session was initialized");
780 let request_id = session.next_request_id();
781 let exchange_cancellation = cancellation.child_token();
782 enum Wait {
783 Response(Result<Value, String>),
784 Cancelled,
785 TimedOut,
786 }
787 let wait = tokio::select! {
788 biased;
789 _ = cancellation.cancelled() => Wait::Cancelled,
790 result = timeout(
791 self.config.tool_timeout,
792 session.request(
793 "tools/call",
794 json!({"name": tool_name, "arguments": arguments}),
795 parameter_headers,
796 exchange_cancellation.clone(),
797 ),
798 ) => match result {
799 Ok(response) => Wait::Response(response),
800 Err(_) => Wait::TimedOut,
801 }
802 };
803 match wait {
804 Wait::Response(Ok(value)) => validate_mcp_call_result(
805 session
806 .negotiated_protocol()
807 .map_err(|error| GuardedMcpCallError::Failed(error.to_string()))?,
808 value,
809 ),
810 Wait::Response(Err(error)) => Err(GuardedMcpCallError::UnknownEffect(error)),
811 Wait::Cancelled => {
812 exchange_cancellation.cancel();
813 session
814 .cancel_request(request_id, "Agent Run cancelled")
815 .await;
816 Err(GuardedMcpCallError::UnknownEffect(
817 "MCP call was cancelled after dispatch; remote effect is unknown"
818 .to_owned(),
819 ))
820 }
821 Wait::TimedOut => {
822 exchange_cancellation.cancel();
823 session
824 .cancel_request(request_id, "Host deadline exceeded")
825 .await;
826 Err(GuardedMcpCallError::UnknownEffect(
827 format!(
828 "MCP server '{}' Tool '{}' timed out after {} ms after dispatch; remote effect is unknown",
829 self.config.server_id,
830 tool_name,
831 self.config.tool_timeout.as_millis()
832 ),
833 ))
834 }
835 }
836 };
837 if request
838 .as_ref()
839 .is_err_and(GuardedMcpCallError::invalidates_session)
840 {
841 self.health
842 .store(McpServerHealth::Degraded.encode(), Ordering::Release);
843 if let Some(mut session) = guard.take() {
844 let _ = session.shutdown().await;
845 }
846 } else {
847 self.health
848 .store(McpServerHealth::Ready.encode(), Ordering::Release);
849 }
850 request
851 }
852
853 pub async fn shutdown(&self) {
854 let mut guard = self.session.lock().await;
855 if let Some(mut session) = guard.take() {
856 let _ = session.shutdown().await;
857 }
858 self.health
859 .store(McpServerHealth::Closed.encode(), Ordering::Release);
860 }
861}
862
863async fn connect_guarded_mcp_session(
864 config: &GuardedMcpServerConfig,
865 cancellation: &CancellationToken,
866) -> Result<McpTransportSession, McpToolsAdapterError> {
867 let mut session = spawn_guarded_mcp_session(config, cancellation).await?;
868 let probe = tokio::select! {
869 biased;
870 _ = cancellation.cancelled() => {
871 let _ = session.shutdown().await;
872 return Err(McpToolsAdapterError::Cancelled);
873 }
874 result = timeout(
875 config.startup_timeout,
876 session.probe_stateless(cancellation.child_token()),
877 ) => result,
878 };
879 match probe {
880 Ok(Ok(true)) => return Ok(session),
881 Ok(Ok(false)) if config.transport_kind() == McpTransportKind::Stdio => {}
882 Ok(Ok(false)) => {
883 let _ = session.shutdown().await;
884 return Err(McpToolsAdapterError::Protocol(
885 "Streamable HTTP endpoint does not support MCP 2026-07-28".to_owned(),
886 ));
887 }
888 Ok(Err(error)) if error.is_transport() => {
889 let _ = session.shutdown().await;
890 if config.transport_kind() != McpTransportKind::Stdio {
891 return Err(mcp_request_adapter_error(error));
892 }
893 session = spawn_guarded_mcp_session(config, cancellation).await?;
894 }
895 Ok(Err(error)) => {
896 let _ = session.shutdown().await;
897 return Err(mcp_request_adapter_error(error));
898 }
899 Err(_) => {
900 let _ = session.shutdown().await;
901 if config.transport_kind() != McpTransportKind::Stdio {
902 return Err(McpToolsAdapterError::Transport(
903 "MCP server/discover timed out".to_owned(),
904 ));
905 }
906 session = spawn_guarded_mcp_session(config, cancellation).await?;
909 }
910 }
911 let initialized = tokio::select! {
912 biased;
913 _ = cancellation.cancelled() => {
914 let _ = session.shutdown().await;
915 return Err(McpToolsAdapterError::Cancelled);
916 }
917 result = timeout(
918 config.startup_timeout,
919 session.initialize_guarded_legacy(cancellation.child_token()),
920 ) => result,
921 };
922 match initialized {
923 Ok(Ok(())) => Ok(session),
924 Ok(Err(error)) => {
925 let _ = session.shutdown().await;
926 Err(mcp_request_adapter_error(error))
927 }
928 Err(_) => {
929 let _ = session.shutdown().await;
930 Err(McpToolsAdapterError::Transport(
931 "legacy MCP initialize timed out".to_owned(),
932 ))
933 }
934 }
935}
936
937async fn spawn_guarded_mcp_session(
938 config: &GuardedMcpServerConfig,
939 cancellation: &CancellationToken,
940) -> Result<McpTransportSession, McpToolsAdapterError> {
941 let connection = tokio::select! {
942 biased;
943 _ = cancellation.cancelled() => return Err(McpToolsAdapterError::Cancelled),
944 result = timeout(config.startup_timeout, config.transport.connect()) => match result {
945 Ok(Ok(connection)) => connection,
946 Ok(Err(error)) => return Err(mcp_request_adapter_error(error)),
947 Err(_) => {
948 return Err(McpToolsAdapterError::Transport(
949 "MCP transport connect timed out".to_owned(),
950 ));
951 }
952 },
953 };
954 if connection.kind() != config.transport_kind() {
955 let _ = timeout(MCP_TRANSPORT_CLOSE_TIMEOUT, connection.close()).await;
956 return Err(McpToolsAdapterError::InvalidConfig(format!(
957 "MCP transport factory for '{}' returned a connection of the wrong kind",
958 config.server_id
959 )));
960 }
961 Ok(McpTransportSession::new(connection))
962}
963
964async fn discover_guarded_mcp_session(
965 config: &GuardedMcpServerConfig,
966 cancellation: &CancellationToken,
967) -> Result<
968 (
969 McpTransportSession,
970 McpServerSnapshot,
971 BTreeMap<String, Vec<McpHeaderBinding>>,
972 ),
973 McpToolsAdapterError,
974> {
975 let mut session = connect_guarded_mcp_session(config, cancellation).await?;
976 let listed = tokio::select! {
977 biased;
978 _ = cancellation.cancelled() => {
979 let _ = session.shutdown().await;
980 return Err(McpToolsAdapterError::Cancelled);
981 }
982 result = timeout(
983 config.startup_timeout,
984 list_all_mcp_tools(&mut session, cancellation),
985 ) => result,
986 };
987 let listed = match listed {
988 Ok(Ok(value)) => value,
989 Ok(Err(error)) => {
990 let _ = session.shutdown().await;
991 return Err(McpToolsAdapterError::Transport(error));
992 }
993 Err(_) => {
994 let _ = session.shutdown().await;
995 return Err(McpToolsAdapterError::Transport(
996 "MCP tools/list timed out".to_owned(),
997 ));
998 }
999 };
1000 let negotiated = session
1001 .negotiated_protocol()
1002 .map_err(mcp_request_adapter_error)?
1003 .clone();
1004 match parse_guarded_tool_snapshot(config, &negotiated, &listed) {
1005 Ok((snapshot, parameter_headers)) => Ok((session, snapshot, parameter_headers)),
1006 Err(error) => {
1007 let _ = session.shutdown().await;
1008 Err(error)
1009 }
1010 }
1011}
1012
1013async fn list_all_mcp_tools(
1014 session: &mut McpTransportSession,
1015 cancellation: &CancellationToken,
1016) -> Result<Value, String> {
1017 const MAX_PAGES: usize = 256;
1018 let stateless = session
1019 .negotiated_protocol()
1020 .map_err(|error| error.to_string())?
1021 .era
1022 == McpProtocolEra::Stateless;
1023 let mut tools = Vec::new();
1024 let mut cursor: Option<String> = None;
1025 let mut seen_cursors = BTreeSet::new();
1026 for _ in 0..MAX_PAGES {
1027 let params = cursor
1028 .as_ref()
1029 .map(|cursor| json!({"cursor": cursor}))
1030 .unwrap_or_else(|| json!({}));
1031 let result = session
1032 .request(
1033 "tools/list",
1034 params,
1035 BTreeMap::new(),
1036 cancellation.child_token(),
1037 )
1038 .await?;
1039 match result.get("resultType").and_then(Value::as_str) {
1040 Some("input_required") => {
1041 return Err(
1042 "MCP tools/list requested unsupported multi-round-trip input".to_owned(),
1043 )
1044 }
1045 Some("complete") => {}
1046 None if !stateless => {}
1047 Some(other) => {
1048 return Err(format!(
1049 "MCP tools/list returned unknown resultType '{other}'"
1050 ))
1051 }
1052 None => return Err("stateless MCP tools/list omitted resultType".to_owned()),
1053 }
1054 if stateless {
1055 validate_stateless_cache_hint(&result)?;
1056 }
1057 let page = result
1058 .get("tools")
1059 .and_then(Value::as_array)
1060 .ok_or_else(|| "MCP tools/list omitted a tools array".to_owned())?;
1061 tools.extend(page.iter().cloned());
1062 cursor = result
1063 .get("nextCursor")
1064 .and_then(Value::as_str)
1065 .filter(|cursor| !cursor.is_empty())
1066 .map(str::to_owned);
1067 let Some(next) = cursor.as_ref() else {
1068 return Ok(json!({"tools": tools}));
1069 };
1070 if !seen_cursors.insert(next.clone()) {
1071 return Err("MCP tools/list repeated a pagination cursor".to_owned());
1072 }
1073 }
1074 Err(format!(
1075 "MCP tools/list exceeded the {MAX_PAGES}-page safety limit"
1076 ))
1077}
1078
1079fn validate_stateless_cache_hint(result: &Value) -> Result<(), String> {
1080 if result.get("ttlMs").and_then(Value::as_u64).is_none()
1081 || !matches!(
1082 result.get("cacheScope").and_then(Value::as_str),
1083 Some("public" | "private")
1084 )
1085 {
1086 return Err(
1087 "stateless MCP tools/list requires non-negative ttlMs and public/private cacheScope"
1088 .to_owned(),
1089 );
1090 }
1091 Ok(())
1092}
1093
1094fn parse_guarded_tool_snapshot(
1095 config: &GuardedMcpServerConfig,
1096 negotiated: &NegotiatedMcpProtocol,
1097 result: &Value,
1098) -> Result<(McpServerSnapshot, BTreeMap<String, Vec<McpHeaderBinding>>), McpToolsAdapterError> {
1099 let tools = result
1100 .get("tools")
1101 .and_then(Value::as_array)
1102 .ok_or_else(|| {
1103 McpToolsAdapterError::Protocol("MCP tools/list omitted a tools array".to_owned())
1104 })?;
1105 let mut names = BTreeSet::new();
1106 let mut snapshots = Vec::new();
1107 let mut parameter_headers = BTreeMap::new();
1108 for raw in tools {
1109 let name = raw
1110 .get("name")
1111 .and_then(Value::as_str)
1112 .map(str::trim)
1113 .filter(|name| !name.is_empty())
1114 .ok_or_else(|| {
1115 McpToolsAdapterError::Protocol(
1116 "MCP tools/list returned an invalid Tool name".to_owned(),
1117 )
1118 })?;
1119 if !config.allows_tool(name) {
1120 continue;
1121 }
1122 if !names.insert(name.to_owned()) {
1123 return Err(McpToolsAdapterError::Protocol(format!(
1124 "MCP server '{}' returned duplicate Tool '{name}'",
1125 config.server_id
1126 )));
1127 }
1128 let input_schema = raw
1129 .get("inputSchema")
1130 .cloned()
1131 .unwrap_or_else(|| json!({"type": "object"}));
1132 if config.transport_kind() == McpTransportKind::StreamableHttp {
1133 match compile_mcp_header_bindings(&input_schema) {
1134 Ok(bindings) => {
1135 parameter_headers.insert(name.to_owned(), bindings);
1136 }
1137 Err(error) => {
1138 tracing::warn!(
1139 server = %config.server_id,
1140 tool = name,
1141 %error,
1142 "excluding MCP Tool with invalid x-mcp-header annotation"
1143 );
1144 continue;
1145 }
1146 }
1147 }
1148 let annotations = raw
1149 .get("annotations")
1150 .cloned()
1151 .map(serde_json::from_value::<McpToolAnnotations>)
1152 .transpose()
1153 .map_err(|error| {
1154 McpToolsAdapterError::Protocol(format!(
1155 "MCP Tool '{name}' returned invalid annotations: {error}"
1156 ))
1157 })?
1158 .unwrap_or_default();
1159 snapshots.push(
1160 McpToolSnapshot::seal_with_annotations(
1161 config.server_id.clone(),
1162 name,
1163 raw.get("description").and_then(Value::as_str).unwrap_or(""),
1164 input_schema,
1165 raw.get("outputSchema").cloned(),
1166 annotations,
1167 )
1168 .map_err(|error| McpToolsAdapterError::Protocol(error.to_string()))?,
1169 );
1170 }
1171 let snapshot = McpServerSnapshot::seal(
1172 config.server_id.clone(),
1173 config.transport_kind(),
1174 config.transport.authority().binding_digest.clone(),
1175 negotiated.version.clone(),
1176 negotiated.era,
1177 snapshots,
1178 )
1179 .map_err(|error| McpToolsAdapterError::Protocol(error.to_string()))?;
1180 Ok((snapshot, parameter_headers))
1181}
1182
1183fn mcp_request_adapter_error(error: McpRequestError) -> McpToolsAdapterError {
1184 match error {
1185 McpRequestError::Transport(message) => McpToolsAdapterError::Transport(message),
1186 error => McpToolsAdapterError::Protocol(error.to_string()),
1187 }
1188}
1189
1190pub struct McpToolsAdapterRegistry {
1192 managers: BTreeMap<McpServerId, Arc<McpServerConnectionManager>>,
1193 skipped_optional_servers: BTreeMap<McpServerId, String>,
1194 tool_count: usize,
1195}
1196
1197impl McpToolsAdapterRegistry {
1198 pub async fn register<S: ApprovalCapabilityStore>(
1199 runtime: &GuardedToolRuntime<S>,
1200 mut configs: Vec<GuardedMcpServerConfig>,
1201 restriction: ToolRestriction,
1202 cancellation: CancellationToken,
1203 ) -> Result<Self, McpToolsAdapterError> {
1204 configs.sort_by(|left, right| left.server_id.cmp(&right.server_id));
1205 restriction
1206 .bounds
1207 .validate()
1208 .map_err(|error| McpToolsAdapterError::InvalidConfig(error.to_string()))?;
1209 let mut configured_ids = BTreeSet::new();
1210 for config in &configs {
1211 config.validate()?;
1212 if !configured_ids.insert(config.server_id.clone()) {
1213 return Err(McpToolsAdapterError::Conflict(format!(
1214 "duplicate MCP server id: {}",
1215 config.server_id
1216 )));
1217 }
1218 let programs = config.allowed_programs();
1219 let allow_child_processes = config.allows_child_processes();
1220 let read_roots = config.filesystem_read_roots();
1221 let write_roots = config.filesystem_write_roots();
1222 let sandbox_profiles = config.sandbox_profiles();
1223 let network_targets = config.allowed_network_targets();
1224 let allow_unrestricted_network = config.allows_unrestricted_network();
1225 let environment = config.environment_names();
1226 let credentials = config.credential_references();
1227 if !config
1228 .effect_scopes()
1229 .is_subset(&restriction.bounds.allowed_effects)
1230 || !programs.is_subset(&restriction.bounds.process.transport.allowed_programs)
1231 || (allow_child_processes
1232 && !restriction.bounds.process.transport.allow_child_processes)
1233 || !read_roots.is_subset(&restriction.bounds.filesystem.readable_roots)
1234 || !write_roots.is_subset(&restriction.bounds.filesystem.writable_roots)
1235 || !sandbox_profiles.is_subset(&restriction.bounds.sandbox.allowed_profiles)
1236 || (!sandbox_profiles.is_empty() && !restriction.bounds.sandbox.required)
1237 || !network_targets.is_subset(&restriction.bounds.network.allowed_targets)
1238 || (allow_unrestricted_network && !restriction.bounds.network.allow_unrestricted)
1239 || !environment.is_subset(&restriction.bounds.environment.allowed_variables)
1240 || !credentials.is_subset(&restriction.bounds.allowed_credentials)
1241 {
1242 return Err(McpToolsAdapterError::InvalidConfig(format!(
1243 "MCP server '{}' exceeds its Host Tool restriction",
1244 config.server_id
1245 )));
1246 }
1247 }
1248 let mut managers = BTreeMap::new();
1249 let mut skipped_optional_servers = BTreeMap::new();
1250 for config in configs {
1251 match McpServerConnectionManager::connect(config.clone(), cancellation.child_token())
1252 .await
1253 {
1254 Ok(manager) => {
1255 managers.insert(config.server_id.clone(), manager);
1256 }
1257 Err(error) if !config.required => {
1258 skipped_optional_servers.insert(config.server_id.clone(), error.to_string());
1259 }
1260 Err(error) => {
1261 shutdown_mcp_managers(&managers).await;
1262 return Err(error);
1263 }
1264 }
1265 }
1266
1267 let existing_names = match runtime.model_tool_schemas() {
1268 Ok(schemas) => schemas
1269 .into_iter()
1270 .map(|schema| schema.name)
1271 .collect::<BTreeSet<_>>(),
1272 Err(error) => {
1273 shutdown_mcp_managers(&managers).await;
1274 return Err(McpToolsAdapterError::ToolRuntime(error));
1275 }
1276 };
1277 let mut registrations = Vec::new();
1278 let mut model_names = BTreeSet::new();
1279 for manager in managers.values() {
1280 let server_restriction = mcp_server_restriction(&restriction, &manager.config);
1281 let mut sanitized_names = BTreeSet::new();
1282 for tool in &manager.snapshot.tools {
1283 let server = sanitize_mcp_identifier(manager.config.server_id.as_str());
1284 let tool_name = sanitize_mcp_identifier(&tool.name);
1285 if !sanitized_names.insert(tool_name.clone()) {
1286 let error = McpToolsAdapterError::Conflict(format!(
1287 "MCP server '{}' has Tool names that collide after namespacing",
1288 manager.config.server_id
1289 ));
1290 shutdown_mcp_managers(&managers).await;
1291 return Err(error);
1292 }
1293 let model_name = format!("mcp__{server}__{tool_name}");
1294 if existing_names.contains(&model_name) || !model_names.insert(model_name.clone()) {
1295 let error = McpToolsAdapterError::Conflict(format!(
1296 "MCP model Tool name collides: {model_name}"
1297 ));
1298 shutdown_mcp_managers(&managers).await;
1299 return Err(error);
1300 }
1301 let descriptor = ToolDescriptor {
1302 tool_id: ToolId::new(format!("mcp/{server}/{tool_name}/v1")),
1303 model_schema: ModelToolSchema {
1304 name: model_name,
1305 description: if tool.description.trim().is_empty() {
1306 format!(
1307 "MCP Tool '{}' on server '{}'",
1308 tool.name, manager.config.server_id
1309 )
1310 } else {
1311 tool.description.clone()
1312 },
1313 input_schema: tool.input_schema.clone(),
1314 },
1315 output_schema: json!({
1316 "type": "object",
1317 "required": ["server", "tool", "result"],
1318 "properties": {
1319 "server": {"type": "string"},
1320 "tool": {"type": "string"},
1321 "result": {}
1322 },
1323 "additionalProperties": false
1324 }),
1325 effect_scopes: manager.config.effect_scopes(),
1326 restriction: server_restriction.clone(),
1327 idempotency: mcp_tool_idempotency(&tool.annotations),
1328 concurrency: ToolConcurrency::GlobalSerial,
1329 };
1330 if let Err(error) = descriptor.validate() {
1331 shutdown_mcp_managers(&managers).await;
1332 return Err(McpToolsAdapterError::Protocol(error.to_string()));
1333 }
1334 registrations.push((
1335 descriptor,
1336 Arc::new(GuardedMcpToolExecutor {
1337 manager: manager.clone(),
1338 tool_name: tool.name.clone(),
1339 annotations: tool.annotations.clone(),
1340 session_approval_scope: tool.schema_digest.clone(),
1341 }) as Arc<dyn GuardedToolExecutor>,
1342 ));
1343 }
1344 }
1345 let tool_count = registrations.len();
1346 for (descriptor, executor) in registrations {
1347 if let Err(error) = runtime.register(descriptor, executor) {
1348 shutdown_mcp_managers(&managers).await;
1349 return Err(McpToolsAdapterError::ToolRuntime(error));
1350 }
1351 }
1352 Ok(Self {
1353 managers,
1354 skipped_optional_servers,
1355 tool_count,
1356 })
1357 }
1358
1359 pub fn tool_count(&self) -> usize {
1360 self.tool_count
1361 }
1362
1363 pub fn server_names(&self) -> BTreeSet<String> {
1364 self.managers
1365 .keys()
1366 .map(|server| server.as_str().to_owned())
1367 .collect()
1368 }
1369
1370 pub fn skipped_optional_servers(&self) -> &BTreeMap<McpServerId, String> {
1371 &self.skipped_optional_servers
1372 }
1373
1374 pub fn manager(&self, server: &McpServerId) -> Option<Arc<McpServerConnectionManager>> {
1375 self.managers.get(server).cloned()
1376 }
1377
1378 pub async fn shutdown(&self) {
1379 for manager in self.managers.values() {
1380 manager.shutdown().await;
1381 }
1382 }
1383}
1384
1385fn mcp_server_restriction(
1389 ceiling: &ToolRestriction,
1390 config: &GuardedMcpServerConfig,
1391) -> ToolRestriction {
1392 let mut bounds = ceiling.bounds.clone();
1393 bounds.allowed_effects = config.effect_scopes();
1394 bounds.sandbox.allowed_profiles = config.sandbox_profiles();
1395 bounds.sandbox.required = !bounds.sandbox.allowed_profiles.is_empty();
1396 bounds.process.interactive = Default::default();
1397 bounds.process.transport.allowed_programs = config.allowed_programs();
1398 bounds.process.transport.allow_child_processes = config.allows_child_processes();
1399 bounds.filesystem.readable_roots = config.filesystem_read_roots();
1400 bounds.filesystem.writable_roots = config.filesystem_write_roots();
1401 bounds.network.allowed_targets = config.allowed_network_targets();
1402 bounds.network.allow_unrestricted = config.allows_unrestricted_network();
1403 bounds.environment.allowed_variables = config.environment_names();
1404 bounds.environment.inherit_host_environment = false;
1405 bounds.allowed_credentials = config.credential_references();
1406 ToolRestriction { bounds }
1407}
1408
1409async fn shutdown_mcp_managers(managers: &BTreeMap<McpServerId, Arc<McpServerConnectionManager>>) {
1410 for manager in managers.values() {
1411 manager.shutdown().await;
1412 }
1413}
1414
1415struct GuardedMcpToolExecutor {
1416 manager: Arc<McpServerConnectionManager>,
1417 tool_name: String,
1418 annotations: McpToolAnnotations,
1419 session_approval_scope: Digest,
1420}
1421
1422fn mcp_tool_idempotency(annotations: &McpToolAnnotations) -> ToolIdempotency {
1423 if annotations.read_only_hint == Some(true) || annotations.idempotent_hint == Some(true) {
1424 ToolIdempotency::Idempotent
1425 } else {
1426 ToolIdempotency::NonIdempotent
1427 }
1428}
1429
1430fn mcp_unknown_effect_outcome(annotations: &McpToolAnnotations, message: String) -> ToolOutcome {
1431 if matches!(
1432 mcp_tool_idempotency(annotations),
1433 ToolIdempotency::Idempotent
1434 ) {
1435 ToolOutcome::Failed {
1436 code: "mcp_call_interrupted".to_owned(),
1437 message: format!(
1438 "{message}; the MCP Tool is annotated read-only or idempotent and may be retried"
1439 ),
1440 retryable: true,
1441 }
1442 } else {
1443 ToolOutcome::UnknownEffect { message }
1444 }
1445}
1446
1447fn mcp_operation_capabilities(
1452 config: &GuardedMcpServerConfig,
1453 tool_name: &str,
1454 annotations: &McpToolAnnotations,
1455) -> CapabilityRequest {
1456 let transport_effects = config.effect_scopes();
1457 let mut required = CapabilityRequest::default();
1458
1459 if transport_effects.contains(&EffectScope::Network) {
1460 if config.allows_unrestricted_network() {
1461 required.insert_resource(EffectScope::Network, CapabilitySelector::Unrestricted);
1462 } else {
1463 for target in config.allowed_network_targets() {
1464 required.insert_resource(EffectScope::Network, CapabilitySelector::Exact(target));
1465 }
1466 }
1467 }
1468
1469 let can_change_external_state =
1470 annotations.read_only_hint != Some(true) || annotations.destructive_hint == Some(true);
1471 if can_change_external_state && transport_effects.contains(&EffectScope::ExternalSideEffect) {
1472 required.insert_resource(
1473 EffectScope::ExternalSideEffect,
1474 CapabilitySelector::Exact(format!("mcp:{}:{tool_name}", config.server_id)),
1475 );
1476 }
1477
1478 required
1479}
1480
1481#[async_trait]
1482impl GuardedToolExecutor for GuardedMcpToolExecutor {
1483 fn planning_contract(&self) -> Value {
1484 json!({
1485 "contract": "orchestral.mcp-operation-planner/v1",
1486 "transportBinding": self.manager.config.transport.authority().binding_digest,
1487 "tool": self.tool_name,
1488 })
1489 }
1490
1491 fn activity_evidence(
1492 &self,
1493 _invocation: &ToolInvocation,
1494 _outcome: Option<&ToolOutcome>,
1495 ) -> Vec<ToolActivityEvidence> {
1496 vec![ToolActivityEvidence::Note {
1497 text: format!("{}/{}", self.manager.config.server_id, self.tool_name),
1498 }]
1499 }
1500
1501 fn plan_operation(
1502 &self,
1503 invocation: &ToolInvocation,
1504 descriptor: &ToolDescriptor,
1505 _effective_policy: &orchestral_core::tool_protocol::EffectiveToolPolicy,
1506 ) -> Result<ToolOperationPlan, ToolOutcome> {
1507 let required_capabilities =
1508 mcp_operation_capabilities(&self.manager.config, &self.tool_name, &self.annotations);
1509 let operation = ToolOperationPlan {
1510 required_capabilities,
1511 risk: if self.annotations.destructive_hint == Some(true) {
1512 ToolOperationRisk::Destructive
1513 } else if self.annotations.requires_approval() {
1514 ToolOperationRisk::Elevated
1515 } else {
1516 ToolOperationRisk::Routine
1517 },
1518 session_approval_scope: Some(self.session_approval_scope.clone()),
1522 summary: self.approval_summary(invocation),
1523 };
1524 operation
1525 .validate_envelope(&descriptor.effect_scopes)
1526 .map_err(|error| ToolOutcome::Rejected {
1527 code: "mcp_operation_invalid".to_owned(),
1528 message: error.message,
1529 })?;
1530 Ok(operation)
1531 }
1532
1533 fn approval_summary(
1534 &self,
1535 invocation: &orchestral_core::tool_protocol::ToolInvocation,
1536 ) -> String {
1537 let digest = invocation
1538 .args_digest()
1539 .map(|digest| digest.to_string())
1540 .unwrap_or_else(|_| "invalid-arguments".to_owned());
1541 format!(
1542 "Call MCP server '{}' Tool '{}' with arguments {}",
1543 self.manager.config.server_id, self.tool_name, digest
1544 )
1545 }
1546
1547 async fn execute(&self, execution: GuardedToolExecution) -> ToolOutcome {
1548 let bounds = execution.effective_policy.bounds();
1549 let programs = self.manager.config.allowed_programs();
1550 let allow_child_processes = self.manager.config.allows_child_processes();
1551 let read_roots = self.manager.config.filesystem_read_roots();
1552 let write_roots = self.manager.config.filesystem_write_roots();
1553 let sandbox_profiles = self.manager.config.sandbox_profiles();
1554 let network_targets = self.manager.config.allowed_network_targets();
1555 let allow_unrestricted_network = self.manager.config.allows_unrestricted_network();
1556 let environment = self.manager.config.environment_names();
1557 let credentials = self.manager.config.credential_references();
1558 if !programs.is_subset(&bounds.process.transport.allowed_programs)
1559 || (allow_child_processes && !bounds.process.transport.allow_child_processes)
1560 || !read_roots.is_subset(&bounds.filesystem.readable_roots)
1561 || !write_roots.is_subset(&bounds.filesystem.writable_roots)
1562 || !sandbox_profiles.is_subset(&bounds.sandbox.allowed_profiles)
1563 || (!sandbox_profiles.is_empty() && !bounds.sandbox.required)
1564 || !network_targets.is_subset(&bounds.network.allowed_targets)
1565 || (allow_unrestricted_network && !bounds.network.allow_unrestricted)
1566 || !environment.is_subset(&bounds.environment.allowed_variables)
1567 || !credentials.is_subset(&bounds.allowed_credentials)
1568 {
1569 return ToolOutcome::Rejected {
1570 code: "mcp_policy_rejected".to_owned(),
1571 message: "MCP transport exceeds the effective Host policy".to_owned(),
1572 };
1573 }
1574 let cleanup_manager = self.manager.clone();
1578 let cleanup_cancellation = execution.cancellation.clone();
1579 let cleanup = tokio::spawn(async move {
1580 cleanup_cancellation.cancelled().await;
1581 cleanup_manager.shutdown().await;
1582 });
1583 let result = self
1584 .manager
1585 .invoke(
1586 &self.tool_name,
1587 execution.invocation.arguments,
1588 execution.cancellation,
1589 )
1590 .await;
1591 cleanup.abort();
1592 match result {
1593 Ok(result) => ToolOutcome::Completed {
1594 output: json!({
1595 "server": self.manager.config.server_id,
1596 "tool": self.tool_name,
1597 "result": result,
1598 })
1599 .into(),
1600 },
1601 Err(GuardedMcpCallError::Rejected(message)) => ToolOutcome::Rejected {
1602 code: "mcp_tool_rejected".to_owned(),
1603 message,
1604 },
1605 Err(GuardedMcpCallError::Failed(message)) => ToolOutcome::Failed {
1606 code: "mcp_call_failed".to_owned(),
1607 message,
1608 retryable: true,
1609 },
1610 Err(GuardedMcpCallError::ToolError(message)) => ToolOutcome::Failed {
1611 code: "mcp_tool_error".to_owned(),
1612 message,
1613 retryable: false,
1614 },
1615 Err(GuardedMcpCallError::Unsupported(message)) => ToolOutcome::Failed {
1616 code: "mcp_feature_unsupported".to_owned(),
1617 message,
1618 retryable: false,
1619 },
1620 Err(GuardedMcpCallError::UnknownEffect(message)) => {
1621 mcp_unknown_effect_outcome(&self.annotations, message)
1622 }
1623 Err(GuardedMcpCallError::Cancelled) => ToolOutcome::Cancelled,
1624 }
1625 }
1626}
1627
1628enum GuardedMcpCallError {
1629 Rejected(String),
1630 Failed(String),
1631 ToolError(String),
1632 Unsupported(String),
1633 UnknownEffect(String),
1634 Cancelled,
1635}
1636
1637impl GuardedMcpCallError {
1638 fn invalidates_session(&self) -> bool {
1639 matches!(self, Self::Failed(_) | Self::UnknownEffect(_))
1640 }
1641}
1642
1643fn validate_mcp_call_result(
1644 negotiated: &NegotiatedMcpProtocol,
1645 result: Value,
1646) -> Result<Value, GuardedMcpCallError> {
1647 if negotiated.era == McpProtocolEra::Stateless {
1648 match result.get("resultType").and_then(Value::as_str) {
1649 Some("complete") => {}
1650 Some("input_required") => {
1651 return Err(GuardedMcpCallError::Unsupported(
1652 "MCP multi-round-trip input is not implemented by Tools Adapter v1".to_owned(),
1653 ))
1654 }
1655 Some(other) => {
1656 return Err(GuardedMcpCallError::Failed(format!(
1657 "MCP tools/call returned unknown resultType '{other}'"
1658 )))
1659 }
1660 None => {
1661 return Err(GuardedMcpCallError::Failed(
1662 "stateless MCP tools/call omitted resultType".to_owned(),
1663 ))
1664 }
1665 }
1666 }
1667 if result.get("isError").and_then(Value::as_bool) == Some(true) {
1668 return Err(GuardedMcpCallError::ToolError(mcp_tool_error_message(
1669 &result,
1670 )));
1671 }
1672 Ok(result)
1673}
1674
1675fn mcp_tool_error_message(result: &Value) -> String {
1679 if let Some(structured) = result
1680 .get("structuredContent")
1681 .filter(|value| !value.is_null())
1682 {
1683 return actionable_mcp_error(structured);
1684 }
1685
1686 let text = result
1687 .get("content")
1688 .and_then(Value::as_array)
1689 .into_iter()
1690 .flatten()
1691 .filter_map(|item| {
1692 (item.get("type").and_then(Value::as_str) == Some("text"))
1693 .then(|| item.get("text").and_then(Value::as_str))
1694 .flatten()
1695 })
1696 .map(|text| {
1697 embedded_json(text)
1698 .as_ref()
1699 .map(actionable_mcp_error)
1700 .unwrap_or_else(|| text.to_owned())
1701 })
1702 .collect::<Vec<_>>()
1703 .join("\n");
1704 if !text.trim().is_empty() {
1705 text
1706 } else {
1707 result.to_string()
1708 }
1709}
1710
1711fn embedded_json(text: &str) -> Option<Value> {
1712 let mut value = serde_json::from_str::<Value>(text.trim()).ok()?;
1713 for _ in 0..2 {
1714 let Value::String(nested) = &value else {
1715 break;
1716 };
1717 let Ok(parsed) = serde_json::from_str::<Value>(nested.trim()) else {
1718 break;
1719 };
1720 value = parsed;
1721 }
1722 Some(value)
1723}
1724
1725fn actionable_mcp_error(value: &Value) -> String {
1726 const KEYS: &[&str] = &["code", "message", "reason", "how_to_get", "hint"];
1727 const MAX_ITEMS: usize = 8;
1728 const MAX_CHARS: usize = 4_096;
1729
1730 fn collect(value: &Value, output: &mut Vec<(String, String)>) {
1731 if output.len() >= MAX_ITEMS {
1732 return;
1733 }
1734 match value {
1735 Value::Object(fields) => {
1736 for key in KEYS {
1737 let Some(value) = fields.get(*key) else {
1738 continue;
1739 };
1740 let rendered = match value {
1741 Value::String(value) => value.trim().to_owned(),
1742 Value::Null => continue,
1743 value => value.to_string(),
1744 };
1745 if !rendered.is_empty()
1746 && !output.iter().any(|(existing_key, existing)| {
1747 existing_key == key && existing == &rendered
1748 })
1749 {
1750 output.push(((*key).to_owned(), rendered));
1751 if output.len() >= MAX_ITEMS {
1752 return;
1753 }
1754 }
1755 }
1756 for value in fields.values() {
1757 collect(value, output);
1758 if output.len() >= MAX_ITEMS {
1759 return;
1760 }
1761 }
1762 }
1763 Value::Array(values) => {
1764 for value in values {
1765 collect(value, output);
1766 if output.len() >= MAX_ITEMS {
1767 return;
1768 }
1769 }
1770 }
1771 _ => {}
1772 }
1773 }
1774
1775 let mut fields = Vec::new();
1776 collect(value, &mut fields);
1777 let message = if fields.is_empty() {
1778 serde_json::to_string_pretty(value).unwrap_or_else(|_| value.to_string())
1779 } else {
1780 fields
1781 .into_iter()
1782 .map(|(key, value)| format!("{key}: {value}"))
1783 .collect::<Vec<_>>()
1784 .join("; ")
1785 };
1786 let mut chars = message.chars();
1787 let mut bounded = chars.by_ref().take(MAX_CHARS).collect::<String>();
1788 if chars.next().is_some() {
1789 bounded.push('…');
1790 }
1791 bounded
1792}
1793
1794#[derive(Debug, thiserror::Error)]
1795#[non_exhaustive]
1796pub enum McpToolsAdapterError {
1797 #[error("invalid MCP Tools Adapter configuration: {0}")]
1798 InvalidConfig(String),
1799 #[error("MCP Tools transport failed: {0}")]
1800 Transport(String),
1801 #[error("MCP Tools protocol failed: {0}")]
1802 Protocol(String),
1803 #[error("MCP Tools registry conflict: {0}")]
1804 Conflict(String),
1805 #[error("MCP Tools startup was cancelled")]
1806 Cancelled,
1807 #[error(transparent)]
1808 ToolRuntime(#[from] ToolRuntimeError),
1809}
1810
1811fn sanitize_mcp_identifier(value: &str) -> String {
1812 let normalized = value
1813 .chars()
1814 .map(|character| {
1815 if character.is_ascii_alphanumeric() || matches!(character, '_' | '-') {
1816 character.to_ascii_lowercase()
1817 } else {
1818 '_'
1819 }
1820 })
1821 .collect::<String>();
1822 let normalized = normalized.trim_matches('_');
1823 if normalized.is_empty() {
1824 "unnamed".to_owned()
1825 } else {
1826 normalized.to_owned()
1827 }
1828}
1829
1830#[derive(Debug, Clone, PartialEq, Eq)]
1831struct NegotiatedMcpProtocol {
1832 version: String,
1833 era: McpProtocolEra,
1834}
1835
1836type McpRequestError = McpTransportError;
1837
1838struct McpTransportSession {
1839 connection: Box<dyn McpTransportConnection>,
1840 next_id: u64,
1841 negotiated: Option<NegotiatedMcpProtocol>,
1842}
1843
1844impl McpTransportSession {
1845 fn new(connection: Box<dyn McpTransportConnection>) -> Self {
1846 Self {
1847 connection,
1848 next_id: 1,
1849 negotiated: None,
1850 }
1851 }
1852
1853 fn next_request_id(&self) -> u64 {
1854 self.next_id
1855 }
1856
1857 async fn probe_stateless(
1858 &mut self,
1859 cancellation: CancellationToken,
1860 ) -> Result<bool, McpRequestError> {
1861 let params = attach_stateless_request_metadata(json!({}))?;
1862 let result = match self
1863 .request_raw("server/discover", params, BTreeMap::new(), cancellation)
1864 .await
1865 {
1866 Ok(result) => result,
1867 Err(error) if error.permits_legacy_fallback() => return Ok(false),
1868 Err(error) => return Err(error),
1869 };
1870 let versions = result
1871 .get("supportedVersions")
1872 .and_then(Value::as_array)
1873 .ok_or_else(|| {
1874 McpRequestError::Protocol("server/discover omitted supportedVersions".to_owned())
1875 })?;
1876 if !versions
1877 .iter()
1878 .any(|version| version.as_str() == Some(MCP_STATELESS_PROTOCOL_2026_07_28))
1879 {
1880 return Ok(false);
1881 }
1882 if result
1883 .get("capabilities")
1884 .and_then(|value| value.get("tools"))
1885 .and_then(Value::as_object)
1886 .is_none()
1887 {
1888 return Err(McpRequestError::Protocol(
1889 "MCP server does not advertise the tools capability".to_owned(),
1890 ));
1891 }
1892 self.negotiated = Some(NegotiatedMcpProtocol {
1893 version: MCP_STATELESS_PROTOCOL_2026_07_28.to_owned(),
1894 era: McpProtocolEra::Stateless,
1895 });
1896 Ok(true)
1897 }
1898
1899 async fn initialize_guarded_legacy(
1900 &mut self,
1901 cancellation: CancellationToken,
1902 ) -> Result<(), McpRequestError> {
1903 let params = json!({
1904 "protocolVersion": MCP_LATEST_LEGACY_PROTOCOL,
1905 "capabilities": {},
1906 "clientInfo": {
1907 "name": "orchestral",
1908 "version": env!("CARGO_PKG_VERSION")
1909 }
1910 });
1911 let result = self
1912 .request_raw(
1913 "initialize",
1914 params,
1915 BTreeMap::new(),
1916 cancellation.child_token(),
1917 )
1918 .await?;
1919 let version = result
1920 .get("protocolVersion")
1921 .and_then(Value::as_str)
1922 .ok_or_else(|| {
1923 McpRequestError::Protocol(
1924 "legacy initialize response omitted protocolVersion".to_owned(),
1925 )
1926 })?;
1927 const SUPPORTED_LEGACY: &[&str] = &["2025-11-25", "2025-06-18", "2025-03-26", "2024-11-05"];
1928 if !SUPPORTED_LEGACY.contains(&version) {
1929 return Err(McpRequestError::Protocol(format!(
1930 "server selected unsupported legacy MCP version '{version}'"
1931 )));
1932 }
1933 if result
1934 .get("capabilities")
1935 .and_then(|value| value.get("tools"))
1936 .and_then(Value::as_object)
1937 .is_none()
1938 {
1939 return Err(McpRequestError::Protocol(
1940 "legacy MCP server does not advertise the tools capability".to_owned(),
1941 ));
1942 }
1943 self.negotiated = Some(NegotiatedMcpProtocol {
1944 version: version.to_owned(),
1945 era: McpProtocolEra::LegacyHandshake,
1946 });
1947 self.connection
1948 .notification(
1949 "notifications/initialized",
1950 json!({}),
1951 cancellation.child_token(),
1952 )
1953 .await
1954 }
1955
1956 fn negotiated_protocol(&self) -> Result<&NegotiatedMcpProtocol, McpRequestError> {
1957 self.negotiated
1958 .as_ref()
1959 .ok_or_else(|| McpRequestError::Protocol("MCP transport was not negotiated".to_owned()))
1960 }
1961
1962 async fn shutdown(&mut self) -> Result<(), String> {
1963 match timeout(MCP_TRANSPORT_CLOSE_TIMEOUT, self.connection.close()).await {
1964 Ok(result) => result.map_err(|error| error.to_string()),
1965 Err(_) => Err("MCP transport close timed out".to_owned()),
1966 }
1967 }
1968
1969 async fn cancel_request(&self, request_id: u64, reason: &str) {
1970 if self.connection.cancellation() == McpTransportCancellation::ProtocolNotification {
1971 let _ = timeout(
1972 Duration::from_millis(100),
1973 self.connection.notification(
1974 "notifications/cancelled",
1975 json!({"requestId": request_id, "reason": reason}),
1976 CancellationToken::new(),
1977 ),
1978 )
1979 .await;
1980 }
1981 }
1982
1983 async fn request(
1984 &mut self,
1985 method: &str,
1986 params: Value,
1987 parameter_headers: BTreeMap<String, String>,
1988 cancellation: CancellationToken,
1989 ) -> Result<Value, String> {
1990 let params = match self.negotiated.as_ref().map(|value| value.era) {
1991 Some(McpProtocolEra::Stateless) => {
1992 attach_stateless_request_metadata(params).map_err(|error| error.to_string())?
1993 }
1994 _ => params,
1995 };
1996 self.request_raw(method, params, parameter_headers, cancellation)
1997 .await
1998 .map_err(|error| error.to_string())
1999 }
2000
2001 async fn request_raw(
2002 &mut self,
2003 method: &str,
2004 params: Value,
2005 parameter_headers: BTreeMap<String, String>,
2006 cancellation: CancellationToken,
2007 ) -> Result<Value, McpRequestError> {
2008 let id = self.next_id;
2009 self.next_id = self.next_id.saturating_add(1);
2010 self.connection
2011 .request(
2012 McpTransportRequest {
2013 id,
2014 method: method.to_owned(),
2015 params,
2016 parameter_headers,
2017 },
2018 cancellation,
2019 )
2020 .await
2021 }
2022}
2023
2024struct StdioMcpTransport {
2025 state: tokio::sync::Mutex<StdioMcpTransportState>,
2026}
2027
2028struct StdioMcpTransportState {
2029 child: Child,
2030 stdin: ChildStdin,
2031 stdout: BufReader<ChildStdout>,
2032 stderr_tail: Arc<StdMutex<VecDeque<u8>>>,
2033 stderr_task: tokio::task::JoinHandle<()>,
2034 process_group_id: Option<u32>,
2035 max_frame_bytes: usize,
2036}
2037
2038impl StdioMcpTransport {
2039 async fn connect(
2040 command: &str,
2041 args: &[String],
2042 env: &HashMap<String, String>,
2043 cwd: &std::path::Path,
2044 backend_starts_new_session: bool,
2045 ) -> Result<Self, String> {
2046 let mut cmd = Command::new(command);
2047 cmd.args(args);
2048 cmd.env_clear();
2050 cmd.envs(env);
2051 cmd.current_dir(cwd);
2052 isolate_mcp_process_group(&mut cmd, backend_starts_new_session);
2053 cmd.kill_on_drop(true)
2054 .stdin(std::process::Stdio::piped())
2055 .stdout(std::process::Stdio::piped())
2056 .stderr(std::process::Stdio::piped());
2057
2058 let mut child = cmd
2059 .spawn()
2060 .map_err(|err| format!("spawn mcp process failed: {}", err))?;
2061
2062 let process_group_id = child.id();
2063 let stdin = child
2064 .stdin
2065 .take()
2066 .ok_or_else(|| "mcp stdio missing stdin pipe".to_string())?;
2067 let stdout = child
2068 .stdout
2069 .take()
2070 .ok_or_else(|| "mcp stdio missing stdout pipe".to_string())?;
2071 let stderr = child
2072 .stderr
2073 .take()
2074 .ok_or_else(|| "mcp stdio missing stderr pipe".to_string())?;
2075 let (stderr_tail, stderr_task) = capture_mcp_stderr(stderr);
2076
2077 Ok(Self {
2078 state: tokio::sync::Mutex::new(StdioMcpTransportState {
2079 child,
2080 stdin,
2081 stdout: BufReader::new(stdout),
2082 stderr_tail,
2083 stderr_task,
2084 process_group_id,
2085 max_frame_bytes: DEFAULT_MCP_MAX_FRAME_BYTES,
2086 }),
2087 })
2088 }
2089}
2090
2091impl StdioMcpTransportState {
2092 async fn request(&mut self, request: McpTransportRequest) -> Result<Value, McpTransportError> {
2093 let id = request.id;
2094 let payload = json!({
2095 "jsonrpc": "2.0",
2096 "id": id,
2097 "method": request.method,
2098 "params": request.params,
2099 });
2100 self.write_frame(&payload)
2101 .await
2102 .map_err(McpTransportError::Transport)?;
2103
2104 loop {
2105 let msg = self
2106 .read_frame()
2107 .await
2108 .map_err(McpTransportError::Transport)?;
2109 let matched = msg
2110 .get("id")
2111 .and_then(Value::as_u64)
2112 .map(|value| value == id)
2113 .unwrap_or(false);
2114 if !matched {
2115 if msg.get("id").is_some() && msg.get("method").is_some() {
2116 return Err(McpTransportError::Protocol(
2117 "MCP server initiated a request without a negotiated client capability"
2118 .to_owned(),
2119 ));
2120 }
2121 if msg.get("id").is_some() {
2122 return Err(McpTransportError::Protocol(format!(
2123 "MCP server returned an unexpected response id while waiting for {id}"
2124 )));
2125 }
2126 continue;
2127 }
2128
2129 if let Some(error) = msg.get("error") {
2130 return Err(McpTransportError::Rpc {
2131 code: error.get("code").and_then(Value::as_i64).unwrap_or(-32000),
2132 message: error
2133 .get("message")
2134 .and_then(Value::as_str)
2135 .map(str::to_owned)
2136 .unwrap_or_else(|| error.to_string()),
2137 });
2138 }
2139 return Ok(msg.get("result").cloned().unwrap_or(Value::Null));
2140 }
2141 }
2142
2143 async fn write_frame(&mut self, payload: &Value) -> Result<(), String> {
2144 let body = serde_json::to_vec(payload)
2146 .map_err(|err| format!("serialize mcp payload failed: {}", err))?;
2147 if body.len() > self.max_frame_bytes {
2148 return Err(format!(
2149 "MCP request frame exceeds the {} byte limit",
2150 self.max_frame_bytes
2151 ));
2152 }
2153 self.stdin
2154 .write_all(&body)
2155 .await
2156 .map_err(|err| format!("write mcp payload failed: {}", err))?;
2157 self.stdin
2158 .write_all(b"\n")
2159 .await
2160 .map_err(|err| format!("write mcp newline failed: {}", err))?;
2161 self.stdin
2162 .flush()
2163 .await
2164 .map_err(|err| format!("flush mcp payload failed: {}", err))
2165 }
2166
2167 async fn read_frame(&mut self) -> Result<Value, String> {
2168 loop {
2170 let line = self.read_bounded_line().await?;
2171 let trimmed = line.trim();
2172 if trimmed.is_empty() {
2173 continue;
2174 }
2175 if trimmed.starts_with('{') {
2177 return serde_json::from_str::<Value>(trimmed)
2178 .map_err(|err| format!("parse mcp NDJSON failed: {}", err));
2179 }
2180 if let Some((key, value)) = trimmed.split_once(':') {
2182 if key.trim().eq_ignore_ascii_case("content-length") {
2183 if let Ok(len) = value.trim().parse::<usize>() {
2184 if len > self.max_frame_bytes {
2185 return Err(format!(
2186 "MCP response frame exceeds the {} byte limit",
2187 self.max_frame_bytes
2188 ));
2189 }
2190 let _ = self.read_bounded_line().await?;
2192 let mut body = vec![0_u8; len];
2194 self.stdout
2195 .read_exact(&mut body)
2196 .await
2197 .map_err(|err| format!("read mcp payload failed: {}", err))?;
2198 return serde_json::from_slice::<Value>(&body)
2199 .map_err(|err| format!("parse mcp payload failed: {}", err));
2200 }
2201 }
2202 }
2203 }
2204 }
2205
2206 async fn read_bounded_line(&mut self) -> Result<String, String> {
2207 let mut bytes = Vec::new();
2208 loop {
2209 let (chunk, consumed, complete) = {
2210 let available = self
2211 .stdout
2212 .fill_buf()
2213 .await
2214 .map_err(|error| format!("read MCP frame failed: {error}"))?;
2215 if available.is_empty() {
2216 if bytes.is_empty() {
2217 return Err(self.closed_process_error().await);
2218 }
2219 (Vec::new(), 0, true)
2220 } else if let Some(index) = available.iter().position(|byte| *byte == b'\n') {
2221 (available[..=index].to_vec(), index + 1, true)
2222 } else {
2223 (available.to_vec(), available.len(), false)
2224 }
2225 };
2226 if bytes.len().saturating_add(chunk.len()) > self.max_frame_bytes {
2227 return Err(format!(
2228 "MCP response line exceeds the {} byte limit",
2229 self.max_frame_bytes
2230 ));
2231 }
2232 bytes.extend_from_slice(&chunk);
2233 self.stdout.consume(consumed);
2234 if complete {
2235 return String::from_utf8(bytes)
2236 .map_err(|error| format!("MCP response is not UTF-8: {error}"));
2237 }
2238 }
2239 }
2240
2241 async fn closed_process_error(&mut self) -> String {
2242 let status = self
2243 .child
2244 .try_wait()
2245 .ok()
2246 .flatten()
2247 .map(|status| status.to_string())
2248 .unwrap_or_else(|| "unknown status".to_owned());
2249 tokio::task::yield_now().await;
2250 let stderr = self
2251 .stderr_tail
2252 .lock()
2253 .map(|tail| tail.iter().copied().collect::<Vec<_>>())
2254 .unwrap_or_default();
2255 let stderr = String::from_utf8_lossy(&stderr)
2256 .split_whitespace()
2257 .collect::<Vec<_>>()
2258 .join(" ");
2259 if stderr.is_empty() {
2260 format!("mcp process closed stdout ({status})")
2261 } else {
2262 format!("mcp process closed stdout ({status}): {stderr}")
2263 }
2264 }
2265}
2266
2267#[async_trait]
2268impl McpTransportConnection for StdioMcpTransport {
2269 fn kind(&self) -> McpTransportKind {
2270 McpTransportKind::Stdio
2271 }
2272
2273 fn cancellation(&self) -> McpTransportCancellation {
2274 McpTransportCancellation::ProtocolNotification
2275 }
2276
2277 async fn request(
2278 &self,
2279 request: McpTransportRequest,
2280 cancellation: CancellationToken,
2281 ) -> Result<Value, McpTransportError> {
2282 let mut state = tokio::select! {
2283 biased;
2284 _ = cancellation.cancelled() => {
2285 return Err(McpTransportError::Transport(
2286 "stdio MCP request was cancelled before dispatch".to_owned(),
2287 ));
2288 }
2289 state = self.state.lock() => state,
2290 };
2291 tokio::select! {
2292 biased;
2293 _ = cancellation.cancelled() => Err(McpTransportError::Transport(
2294 "stdio MCP request was cancelled after dispatch".to_owned(),
2295 )),
2296 result = state.request(request) => result,
2297 }
2298 }
2299
2300 async fn notification(
2301 &self,
2302 method: &str,
2303 params: Value,
2304 cancellation: CancellationToken,
2305 ) -> Result<(), McpTransportError> {
2306 let payload = json!({
2307 "jsonrpc": "2.0",
2308 "method": method,
2309 "params": params,
2310 });
2311 let mut state = tokio::select! {
2312 biased;
2313 _ = cancellation.cancelled() => {
2314 return Err(McpTransportError::Transport(
2315 "stdio MCP notification was cancelled".to_owned(),
2316 ));
2317 }
2318 state = self.state.lock() => state,
2319 };
2320 state
2321 .write_frame(&payload)
2322 .await
2323 .map_err(McpTransportError::Transport)
2324 }
2325
2326 async fn close(&self) -> Result<(), McpTransportError> {
2327 let mut state = self.state.lock().await;
2328 let process_group_id = state.process_group_id;
2329 terminate_mcp_process_tree(&mut state.child, process_group_id).await;
2330 state.stderr_task.abort();
2331 Ok(())
2332 }
2333}
2334
2335fn capture_mcp_stderr(
2336 stderr: ChildStderr,
2337) -> (Arc<StdMutex<VecDeque<u8>>>, tokio::task::JoinHandle<()>) {
2338 let tail = Arc::new(StdMutex::new(VecDeque::with_capacity(
2339 MCP_STDERR_CAPTURE_BYTES,
2340 )));
2341 let captured = tail.clone();
2342 let task = tokio::spawn(async move {
2343 let mut stderr = BufReader::new(stderr);
2344 let mut buffer = [0_u8; 4096];
2345 loop {
2346 let Ok(count) = stderr.read(&mut buffer).await else {
2347 break;
2348 };
2349 if count == 0 {
2350 break;
2351 }
2352 let Ok(mut tail) = captured.lock() else {
2353 break;
2354 };
2355 tail.extend(&buffer[..count]);
2356 while tail.len() > MCP_STDERR_CAPTURE_BYTES {
2357 tail.pop_front();
2358 }
2359 }
2360 });
2361 (tail, task)
2362}
2363
2364fn attach_stateless_request_metadata(mut params: Value) -> Result<Value, McpRequestError> {
2365 let object = params.as_object_mut().ok_or_else(|| {
2366 McpRequestError::Protocol("MCP request params must be an object".to_owned())
2367 })?;
2368 if object.contains_key("_meta") {
2369 return Err(McpRequestError::Protocol(
2370 "MCP caller cannot override Host request metadata".to_owned(),
2371 ));
2372 }
2373 object.insert(
2374 "_meta".to_owned(),
2375 json!({
2376 "io.modelcontextprotocol/protocolVersion": MCP_STATELESS_PROTOCOL_2026_07_28,
2377 "io.modelcontextprotocol/clientInfo": {
2378 "name": "orchestral",
2379 "version": env!("CARGO_PKG_VERSION")
2380 },
2381 "io.modelcontextprotocol/clientCapabilities": {}
2382 }),
2383 );
2384 Ok(params)
2385}
2386
2387#[cfg(unix)]
2388fn isolate_mcp_process_group(command: &mut Command, backend_starts_new_session: bool) {
2389 if !backend_starts_new_session {
2390 command.process_group(0);
2391 }
2392}
2393
2394#[cfg(not(unix))]
2395fn isolate_mcp_process_group(_command: &mut Command, _backend_starts_new_session: bool) {}
2396
2397async fn terminate_mcp_process_tree(child: &mut Child, process_group_id: Option<u32>) {
2398 #[cfg(unix)]
2399 if let Some(process_group_id) = process_group_id.filter(|id| *id <= i32::MAX as u32) {
2400 unsafe {
2402 libc::kill(-(process_group_id as i32), libc::SIGKILL);
2403 }
2404 }
2405 #[cfg(not(unix))]
2406 let _ = process_group_id;
2407 let _ = child.kill().await;
2408 let _ = timeout(MCP_PROCESS_REAP_TIMEOUT, child.wait()).await;
2409}
2410
2411#[cfg(test)]
2412mod mcp_header_tests {
2413 use super::*;
2414
2415 #[test]
2416 fn nested_static_property_headers_compile_and_extract_canonically() {
2417 let bindings = compile_mcp_header_bindings(&json!({
2418 "type": "object",
2419 "properties": {
2420 "region": {"type": "string", "x-mcp-header": "Region"},
2421 "routing": {
2422 "type": "object",
2423 "properties": {
2424 "priority": {"type": "integer", "x-mcp-header": "Priority"},
2425 "dryRun": {"type": "boolean", "x-mcp-header": "Dry-Run"}
2426 }
2427 }
2428 }
2429 }))
2430 .unwrap();
2431 let headers = build_mcp_parameter_headers(
2432 &bindings,
2433 &json!({
2434 "region": "us-west1",
2435 "routing": {"priority": -7, "dryRun": true}
2436 }),
2437 )
2438 .unwrap();
2439 assert_eq!(headers.get("Region").map(String::as_str), Some("us-west1"));
2440 assert_eq!(headers.get("Priority").map(String::as_str), Some("-7"));
2441 assert_eq!(headers.get("Dry-Run").map(String::as_str), Some("true"));
2442 }
2443
2444 #[test]
2445 fn non_static_duplicate_and_non_primitive_annotations_are_rejected() {
2446 let invalid = [
2447 json!({
2448 "type": "object",
2449 "properties": {"values": {"type": "array", "items": {
2450 "type": "string", "x-mcp-header": "Item"
2451 }}}
2452 }),
2453 json!({
2454 "type": "object",
2455 "allOf": [{"properties": {"value": {
2456 "type": "string", "x-mcp-header": "Value"
2457 }}}]
2458 }),
2459 json!({
2460 "type": "object",
2461 "properties": {
2462 "left": {"type": "string", "x-mcp-header": "Route"},
2463 "right": {"type": "string", "x-mcp-header": "route"}
2464 }
2465 }),
2466 json!({
2467 "type": "object",
2468 "properties": {"ratio": {"type": "number", "x-mcp-header": "Ratio"}}
2469 }),
2470 json!({
2471 "type": "object",
2472 "properties": {"value": {"type": "string", "x-mcp-header": "bad header"}}
2473 }),
2474 ];
2475 for schema in invalid {
2476 assert!(compile_mcp_header_bindings(&schema).is_err());
2477 }
2478 }
2479
2480 #[test]
2481 fn header_integer_extraction_enforces_the_javascript_safe_range() {
2482 let bindings = compile_mcp_header_bindings(&json!({
2483 "type": "object",
2484 "properties": {"value": {"type": "integer", "x-mcp-header": "Value"}}
2485 }))
2486 .unwrap();
2487 assert!(build_mcp_parameter_headers(
2488 &bindings,
2489 &json!({"value": MAX_SAFE_MCP_HEADER_INTEGER})
2490 )
2491 .is_ok());
2492 assert!(build_mcp_parameter_headers(
2493 &bindings,
2494 &json!({"value": MAX_SAFE_MCP_HEADER_INTEGER + 1})
2495 )
2496 .is_err());
2497 }
2498}
2499
2500#[cfg(test)]
2501mod mcp_lifecycle_gate_tests {
2502 use super::*;
2503 use crate::tool_runtime::{GuardedToolResult, ToolArtifactStore};
2504 use crate::InMemoryBlobStore;
2505 use orchestral_core::agent_protocol::wire::RunId;
2506 use orchestral_core::tool_effect::{
2507 replay_tool_effect, InMemoryToolEffectJournalStore, ToolEffectJournalStore, ToolEffectKey,
2508 ToolEffectPhase,
2509 };
2510 use orchestral_core::tool_protocol::ApprovalPolicy;
2511 use orchestral_core::tool_protocol::{
2512 HostApprovalIssuer, HostApprovalVerifier, HostToolPolicy, InMemoryApprovalCapabilityStore,
2513 RunToolGrant, ToolCallId, ToolInvocation, ToolOutput, ToolPolicyBounds,
2514 };
2515 use std::future::pending;
2516 use std::sync::atomic::{AtomicBool, AtomicUsize};
2517 use std::time::Instant;
2518 use tokio::sync::Notify;
2519
2520 const GATE_CASES: usize = 1_000;
2521 const GATE_SIGNING_KEY: &[u8] = b"0123456789abcdef0123456789abcdef";
2522
2523 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
2524 enum FaultStage {
2525 Healthy,
2526 Connect,
2527 Discover,
2528 Initialize,
2529 List,
2530 Call,
2531 BodyDecode,
2532 ReconnectStable,
2533 SchemaChanged,
2534 NameConflict,
2535 LargeResult,
2536 }
2537
2538 struct FaultState {
2539 stage: FaultStage,
2540 connects: AtomicUsize,
2541 active_connects: AtomicUsize,
2542 active_connections: AtomicUsize,
2543 active_requests: AtomicUsize,
2544 closes: AtomicUsize,
2545 explicit_closes: AtomicUsize,
2546 tool_call_dispatches: AtomicUsize,
2547 completed_call_responses: AtomicUsize,
2548 entered: Notify,
2549 }
2550
2551 impl FaultState {
2552 fn new(stage: FaultStage) -> Arc<Self> {
2553 Arc::new(Self {
2554 stage,
2555 connects: AtomicUsize::new(0),
2556 active_connects: AtomicUsize::new(0),
2557 active_connections: AtomicUsize::new(0),
2558 active_requests: AtomicUsize::new(0),
2559 closes: AtomicUsize::new(0),
2560 explicit_closes: AtomicUsize::new(0),
2561 tool_call_dispatches: AtomicUsize::new(0),
2562 completed_call_responses: AtomicUsize::new(0),
2563 entered: Notify::new(),
2564 })
2565 }
2566
2567 fn assert_released(&self) {
2568 assert_eq!(self.active_connects.load(Ordering::SeqCst), 0);
2569 assert_eq!(self.active_connections.load(Ordering::SeqCst), 0);
2570 assert_eq!(self.active_requests.load(Ordering::SeqCst), 0);
2571 }
2572 }
2573
2574 enum ActiveCounter {
2575 Connect(Arc<FaultState>),
2576 Request(Arc<FaultState>),
2577 }
2578
2579 impl Drop for ActiveCounter {
2580 fn drop(&mut self) {
2581 match self {
2582 Self::Connect(state) => {
2583 state.active_connects.fetch_sub(1, Ordering::SeqCst);
2584 }
2585 Self::Request(state) => {
2586 state.active_requests.fetch_sub(1, Ordering::SeqCst);
2587 }
2588 }
2589 }
2590 }
2591
2592 struct FaultFactory {
2593 authority: McpTransportAuthority,
2594 state: Arc<FaultState>,
2595 }
2596
2597 impl FaultFactory {
2598 fn new(stage: FaultStage) -> Arc<Self> {
2599 let kind = match stage {
2600 FaultStage::Discover | FaultStage::BodyDecode => McpTransportKind::StreamableHttp,
2601 FaultStage::Healthy
2602 | FaultStage::Connect
2603 | FaultStage::Initialize
2604 | FaultStage::List
2605 | FaultStage::Call
2606 | FaultStage::ReconnectStable
2607 | FaultStage::SchemaChanged
2608 | FaultStage::NameConflict
2609 | FaultStage::LargeResult => McpTransportKind::Stdio,
2610 };
2611 let (effect_scopes, process_programs, network_targets) = match kind {
2612 McpTransportKind::Stdio => (
2613 BTreeSet::from([
2614 EffectScope::Process,
2615 EffectScope::FilesystemRead,
2616 EffectScope::FilesystemWrite,
2617 EffectScope::ExternalSideEffect,
2618 ]),
2619 BTreeSet::from(["/fault/mcp".to_owned()]),
2620 BTreeSet::new(),
2621 ),
2622 McpTransportKind::StreamableHttp => (
2623 BTreeSet::from([EffectScope::Network, EffectScope::ExternalSideEffect]),
2624 BTreeSet::new(),
2625 BTreeSet::from(["http://127.0.0.1/fault-mcp".to_owned()]),
2626 ),
2627 _ => unreachable!("fault factory only selects v1 transport kinds"),
2628 };
2629 let (filesystem_read_roots, filesystem_write_roots, sandbox_profiles) = match kind {
2630 McpTransportKind::Stdio => (
2631 BTreeSet::from(["/fault/read".to_owned()]),
2632 BTreeSet::from(["/fault/write".to_owned()]),
2633 BTreeSet::from([MCP_STDIO_SANDBOX_PROFILE.to_owned()]),
2634 ),
2635 McpTransportKind::StreamableHttp => {
2636 (BTreeSet::new(), BTreeSet::new(), BTreeSet::new())
2637 }
2638 _ => unreachable!("fault factory only selects v1 transport kinds"),
2639 };
2640 Arc::new(Self {
2641 authority: McpTransportAuthority {
2642 kind,
2643 binding_digest: Digest::sha256(format!("fault-stage-{stage:?}")),
2644 effect_scopes,
2645 process_programs,
2646 allow_child_processes: false,
2647 allow_host_ui: false,
2648 filesystem_read_roots,
2649 filesystem_write_roots,
2650 sandbox_profiles,
2651 network_targets,
2652 allow_unrestricted_network: false,
2653 environment_variables: BTreeSet::new(),
2654 credential_references: BTreeSet::new(),
2655 },
2656 state: FaultState::new(stage),
2657 })
2658 }
2659
2660 fn with_authority(authority: McpTransportAuthority) -> Arc<Self> {
2661 Arc::new(Self {
2662 authority,
2663 state: FaultState::new(FaultStage::Healthy),
2664 })
2665 }
2666 }
2667
2668 #[async_trait]
2669 impl McpTransportFactory for FaultFactory {
2670 fn authority(&self) -> &McpTransportAuthority {
2671 &self.authority
2672 }
2673
2674 async fn connect(&self) -> Result<Box<dyn McpTransportConnection>, McpTransportError> {
2675 let generation = self.state.connects.fetch_add(1, Ordering::SeqCst) + 1;
2676 if self.state.stage == FaultStage::Connect {
2677 self.state.active_connects.fetch_add(1, Ordering::SeqCst);
2678 let _guard = ActiveCounter::Connect(self.state.clone());
2679 self.state.entered.notify_one();
2680 return pending().await;
2681 }
2682 self.state.active_connections.fetch_add(1, Ordering::SeqCst);
2683 Ok(Box::new(FaultConnection {
2684 kind: self.authority.kind,
2685 state: self.state.clone(),
2686 closed: AtomicBool::new(false),
2687 generation,
2688 }))
2689 }
2690 }
2691
2692 struct FaultConnection {
2693 kind: McpTransportKind,
2694 state: Arc<FaultState>,
2695 closed: AtomicBool,
2696 generation: usize,
2697 }
2698
2699 impl FaultConnection {
2700 async fn block_request(&self) -> Result<Value, McpTransportError> {
2701 self.state.active_requests.fetch_add(1, Ordering::SeqCst);
2702 let _guard = ActiveCounter::Request(self.state.clone());
2703 self.state.entered.notify_one();
2704 pending().await
2705 }
2706
2707 fn release(&self, explicit: bool) {
2708 if !self.closed.swap(true, Ordering::SeqCst) {
2709 self.state.active_connections.fetch_sub(1, Ordering::SeqCst);
2710 self.state.closes.fetch_add(1, Ordering::SeqCst);
2711 if explicit {
2712 self.state.explicit_closes.fetch_add(1, Ordering::SeqCst);
2713 }
2714 }
2715 }
2716 }
2717
2718 impl Drop for FaultConnection {
2719 fn drop(&mut self) {
2720 self.release(false);
2721 }
2722 }
2723
2724 #[async_trait]
2725 impl McpTransportConnection for FaultConnection {
2726 fn kind(&self) -> McpTransportKind {
2727 self.kind
2728 }
2729
2730 fn cancellation(&self) -> McpTransportCancellation {
2731 McpTransportCancellation::DropExchange
2732 }
2733
2734 async fn request(
2735 &self,
2736 request: McpTransportRequest,
2737 _cancellation: CancellationToken,
2738 ) -> Result<Value, McpTransportError> {
2739 let is_tool_call = request.method == "tools/call";
2740 let result = match request.method.as_str() {
2741 "server/discover" if self.state.stage == FaultStage::Discover => {
2742 self.block_request().await
2743 }
2744 "server/discover" if self.state.stage == FaultStage::Initialize => {
2745 Err(McpTransportError::Rpc {
2746 code: -32601,
2747 message: "legacy server".to_owned(),
2748 })
2749 }
2750 "server/discover" => Ok(json!({
2751 "supportedVersions": [MCP_STATELESS_PROTOCOL_2026_07_28],
2752 "capabilities": {"tools": {}},
2753 "serverInfo": {"name": "fault", "version": "1"}
2754 })),
2755 "initialize" if self.state.stage == FaultStage::Initialize => {
2756 self.block_request().await
2757 }
2758 "initialize" => Ok(json!({
2759 "protocolVersion": MCP_LATEST_LEGACY_PROTOCOL,
2760 "capabilities": {"tools": {}},
2761 "serverInfo": {"name": "fault", "version": "1"}
2762 })),
2763 "tools/list" if self.state.stage == FaultStage::List => self.block_request().await,
2764 "tools/list" if self.state.stage == FaultStage::NameConflict => Ok(json!({
2765 "resultType": "complete",
2766 "ttlMs": 1_000,
2767 "cacheScope": "private",
2768 "tools": [
2769 {"name": "a.b", "inputSchema": {"type": "object"}},
2770 {"name": "a b", "inputSchema": {"type": "object"}}
2771 ]
2772 })),
2773 "tools/list" => {
2774 let echo_schema =
2775 if self.state.stage == FaultStage::SchemaChanged && self.generation > 1 {
2776 json!({
2777 "type": "object",
2778 "properties": {"changed": {"type": "boolean"}},
2779 "additionalProperties": false
2780 })
2781 } else {
2782 json!({"type": "object", "additionalProperties": false})
2783 };
2784 Ok(json!({
2785 "resultType": "complete",
2786 "ttlMs": 1_000,
2787 "cacheScope": "private",
2788 "tools": [
2789 {"name": "echo", "description": "echo", "inputSchema": echo_schema},
2790 {"name": "beta", "description": "beta", "inputSchema": {
2791 "type": "object", "additionalProperties": false
2792 }},
2793 {"name": "hidden", "description": "hidden", "inputSchema": {
2794 "type": "object", "additionalProperties": false
2795 }}
2796 ]
2797 }))
2798 }
2799 "tools/call"
2800 if matches!(self.state.stage, FaultStage::Call | FaultStage::BodyDecode) =>
2801 {
2802 self.block_request().await
2803 }
2804 "tools/call" => {
2805 self.state
2806 .tool_call_dispatches
2807 .fetch_add(1, Ordering::SeqCst);
2808 if matches!(
2809 self.state.stage,
2810 FaultStage::ReconnectStable | FaultStage::SchemaChanged
2811 ) && self.generation == 1
2812 {
2813 Err(McpTransportError::Transport(
2814 "injected disconnect".to_owned(),
2815 ))
2816 } else {
2817 let text = if self.state.stage == FaultStage::LargeResult {
2818 "large-mcp-result/".repeat(256)
2819 } else {
2820 "ok".to_owned()
2821 };
2822 Ok(json!({
2823 "resultType": "complete",
2824 "content": [{"type": "text", "text": text}],
2825 "isError": false
2826 }))
2827 }
2828 }
2829 other => Err(McpTransportError::Protocol(format!(
2830 "unexpected fault transport method: {other}"
2831 ))),
2832 };
2833 if result.is_ok() && is_tool_call {
2834 self.state
2835 .completed_call_responses
2836 .fetch_add(1, Ordering::SeqCst);
2837 }
2838 result
2839 }
2840
2841 async fn notification(
2842 &self,
2843 _method: &str,
2844 _params: Value,
2845 _cancellation: CancellationToken,
2846 ) -> Result<(), McpTransportError> {
2847 Ok(())
2848 }
2849
2850 async fn close(&self) -> Result<(), McpTransportError> {
2851 self.release(true);
2852 Ok(())
2853 }
2854 }
2855
2856 fn fault_config(factory: Arc<FaultFactory>, deadline: Duration) -> GuardedMcpServerConfig {
2857 GuardedMcpServerConfig {
2858 server_id: McpServerId::new(format!("fault-{:?}", factory.state.stage)),
2859 required: true,
2860 transport: factory,
2861 startup_timeout: deadline,
2862 tool_timeout: deadline,
2863 enabled_tools: BTreeSet::new(),
2864 disabled_tools: BTreeSet::new(),
2865 }
2866 }
2867
2868 async fn connect_result(
2869 config: GuardedMcpServerConfig,
2870 cancellation: CancellationToken,
2871 ) -> Result<Arc<McpServerConnectionManager>, McpToolsAdapterError> {
2872 McpServerConnectionManager::connect(config, cancellation).await
2873 }
2874
2875 fn assert_unknown_effect(result: Result<Value, GuardedMcpCallError>) {
2876 assert!(matches!(result, Err(GuardedMcpCallError::UnknownEffect(_))));
2877 }
2878
2879 fn authority_bounds(authority: &McpTransportAuthority) -> ToolPolicyBounds {
2880 let mut bounds = ToolPolicyBounds {
2881 allowed_effects: authority.effect_scopes.clone(),
2882 approval: ApprovalPolicy::Required,
2883 max_timeout_ms: Some(30_000),
2884 max_output_bytes: Some(64 * 1024),
2885 ..ToolPolicyBounds::default()
2886 };
2887 bounds.process.transport.allowed_programs = authority.process_programs.clone();
2888 bounds.process.transport.allow_child_processes = authority.allow_child_processes;
2889 bounds.filesystem.readable_roots = authority.filesystem_read_roots.clone();
2890 bounds.filesystem.writable_roots = authority.filesystem_write_roots.clone();
2891 bounds.sandbox.allowed_profiles = authority.sandbox_profiles.clone();
2892 bounds.sandbox.required = !authority.sandbox_profiles.is_empty();
2893 bounds.network.allowed_targets = authority.network_targets.clone();
2894 bounds.network.allow_unrestricted = authority.allow_unrestricted_network;
2895 bounds.environment.allowed_variables = authority.environment_variables.clone();
2896 bounds.allowed_credentials = authority.credential_references.clone();
2897 bounds
2898 }
2899
2900 #[test]
2901 fn mcp_operation_does_not_reapprove_transport_authority() {
2902 let factory = FaultFactory::new(FaultStage::Healthy);
2903 let config = fault_config(factory, Duration::from_secs(1));
2904
2905 let unknown = mcp_operation_capabilities(&config, "run", &McpToolAnnotations::default());
2906 assert_eq!(
2907 unknown.effects,
2908 BTreeSet::from([EffectScope::ExternalSideEffect])
2909 );
2910 assert!(!unknown.effects.contains(&EffectScope::Process));
2911 assert!(!unknown.effects.contains(&EffectScope::FilesystemRead));
2912 assert!(!unknown.effects.contains(&EffectScope::FilesystemWrite));
2913 assert!(!unknown.effects.contains(&EffectScope::SecretRead));
2914
2915 let read_only = mcp_operation_capabilities(
2916 &config,
2917 "lookup",
2918 &McpToolAnnotations {
2919 read_only_hint: Some(true),
2920 ..McpToolAnnotations::default()
2921 },
2922 );
2923 assert!(read_only.effects.is_empty());
2924 assert!(read_only.resources.is_empty());
2925 }
2926
2927 #[test]
2928 fn mcp_operation_keeps_logical_network_effect_without_transport_secrets() {
2929 let factory = FaultFactory::new(FaultStage::Discover);
2930 let config = fault_config(factory, Duration::from_secs(1));
2931 let request = mcp_operation_capabilities(
2932 &config,
2933 "lookup",
2934 &McpToolAnnotations {
2935 read_only_hint: Some(true),
2936 ..McpToolAnnotations::default()
2937 },
2938 );
2939
2940 assert_eq!(request.effects, BTreeSet::from([EffectScope::Network]));
2941 assert_eq!(
2942 request
2943 .resources_for(EffectScope::Network)
2944 .cloned()
2945 .collect::<BTreeSet<_>>(),
2946 BTreeSet::from([CapabilitySelector::Exact(
2947 "http://127.0.0.1/fault-mcp".to_owned()
2948 )])
2949 );
2950 }
2951
2952 #[test]
2953 fn mcp_tool_error_prefers_structured_content_then_plain_text() {
2954 assert_eq!(
2955 mcp_tool_error_message(&json!({
2956 "structuredContent": {"error": {"code": "bad_args"}},
2957 "content": [{"type": "text", "text": "fallback"}],
2958 "isError": true
2959 })),
2960 "code: bad_args"
2961 );
2962 assert_eq!(
2963 mcp_tool_error_message(&json!({
2964 "content": [
2965 {"type": "text", "text": "first"},
2966 {"type": "image", "data": "ignored"},
2967 {"type": "text", "text": "second"}
2968 ],
2969 "isError": true
2970 })),
2971 "first\nsecond"
2972 );
2973 assert_eq!(
2974 mcp_tool_error_message(&json!({
2975 "content": [{
2976 "type": "text",
2977 "text": serde_json::to_string(&json!({
2978 "capability_schema": {
2979 "how_to_get": "Call search_capabilities first"
2980 },
2981 "reason": "schema omitted"
2982 })).unwrap()
2983 }],
2984 "isError": true
2985 })),
2986 "reason: schema omitted; how_to_get: Call search_capabilities first"
2987 );
2988 }
2989
2990 #[test]
2991 fn mcp_idempotency_follows_standard_annotations() {
2992 assert_eq!(
2993 mcp_tool_idempotency(&McpToolAnnotations {
2994 read_only_hint: Some(true),
2995 ..McpToolAnnotations::default()
2996 }),
2997 ToolIdempotency::Idempotent
2998 );
2999 assert_eq!(
3000 mcp_tool_idempotency(&McpToolAnnotations {
3001 idempotent_hint: Some(true),
3002 ..McpToolAnnotations::default()
3003 }),
3004 ToolIdempotency::Idempotent
3005 );
3006 assert_eq!(
3007 mcp_tool_idempotency(&McpToolAnnotations::default()),
3008 ToolIdempotency::NonIdempotent
3009 );
3010 assert!(matches!(
3011 mcp_unknown_effect_outcome(
3012 &McpToolAnnotations {
3013 read_only_hint: Some(true),
3014 ..McpToolAnnotations::default()
3015 },
3016 "timed out".to_owned()
3017 ),
3018 ToolOutcome::Failed {
3019 ref code,
3020 retryable: true,
3021 ..
3022 } if code == "mcp_call_interrupted"
3023 ));
3024 assert!(matches!(
3025 mcp_unknown_effect_outcome(&McpToolAnnotations::default(), "timed out".to_owned()),
3026 ToolOutcome::UnknownEffect { .. }
3027 ));
3028 }
3029
3030 #[tokio::test]
3031 async fn workspace_auto_policy_runs_unannotated_registered_mcp_without_prompt() {
3032 let factory = FaultFactory::new(FaultStage::Healthy);
3033 let config = fault_config(factory.clone(), Duration::from_secs(1));
3034 let mut bounds = authority_bounds(factory.authority());
3035 bounds.approval = ApprovalPolicy::NotRequired;
3036 let journal = Arc::new(InMemoryToolEffectJournalStore::default());
3037 let verifier =
3038 HostApprovalVerifier::new(GATE_SIGNING_KEY, InMemoryApprovalCapabilityStore::default())
3039 .unwrap();
3040 let runtime = Arc::new(
3041 GuardedToolRuntime::new_with_effect_journal(
3042 HostToolPolicy {
3043 bounds: bounds.clone(),
3044 },
3045 verifier,
3046 journal.clone(),
3047 )
3048 .unwrap()
3049 .with_permission_policy(Arc::new(crate::tool_runtime::WorkspacePermissionPolicy)),
3050 );
3051 let registry = McpToolsAdapterRegistry::register(
3052 runtime.as_ref(),
3053 vec![config],
3054 ToolRestriction {
3055 bounds: bounds.clone(),
3056 },
3057 CancellationToken::new(),
3058 )
3059 .await
3060 .unwrap();
3061 let run_id = RunId::new("auto-mcp");
3062 let call_id = ToolCallId::new("call-1");
3063 let result = runtime
3064 .invoke(
3065 ToolInvocation {
3066 run_id: run_id.clone(),
3067 call_id: call_id.clone(),
3068 tool_id: ToolId::new("mcp/fault-healthy/echo/v1"),
3069 arguments: json!({}),
3070 },
3071 RunToolGrant {
3072 bounds: bounds.clone(),
3073 },
3074 None,
3075 CancellationToken::new(),
3076 )
3077 .await;
3078
3079 assert!(matches!(
3080 result,
3081 GuardedToolResult::Outcome {
3082 outcome: ToolOutcome::Completed { .. },
3083 cached: false
3084 }
3085 ));
3086 let key = ToolEffectKey::new(run_id, call_id);
3087 let records = journal.load_effect(&key).await.unwrap();
3088 let projection = replay_tool_effect(&key, &records).unwrap().unwrap();
3089 assert_eq!(
3090 projection.prepared.effect_scopes,
3091 BTreeSet::from([EffectScope::ExternalSideEffect])
3092 );
3093 registry.shutdown().await;
3094 factory.state.assert_released();
3095 }
3096
3097 fn gate_runtime(
3098 bounds: ToolPolicyBounds,
3099 journal: Arc<InMemoryToolEffectJournalStore>,
3100 ) -> Arc<GuardedToolRuntime<InMemoryApprovalCapabilityStore>> {
3101 let verifier =
3102 HostApprovalVerifier::new(GATE_SIGNING_KEY, InMemoryApprovalCapabilityStore::default())
3103 .unwrap();
3104 Arc::new(
3105 GuardedToolRuntime::new_with_effect_journal(
3106 HostToolPolicy { bounds },
3107 verifier,
3108 journal,
3109 )
3110 .unwrap(),
3111 )
3112 }
3113
3114 fn gate_runtime_with_artifacts(
3115 bounds: ToolPolicyBounds,
3116 journal: Arc<InMemoryToolEffectJournalStore>,
3117 artifacts: ToolArtifactStore,
3118 ) -> Arc<GuardedToolRuntime<InMemoryApprovalCapabilityStore>> {
3119 let verifier =
3120 HostApprovalVerifier::new(GATE_SIGNING_KEY, InMemoryApprovalCapabilityStore::default())
3121 .unwrap();
3122 Arc::new(
3123 GuardedToolRuntime::new_with_effect_journal_and_artifacts(
3124 HostToolPolicy { bounds },
3125 verifier,
3126 journal,
3127 artifacts,
3128 )
3129 .unwrap(),
3130 )
3131 }
3132
3133 async fn invoke_with_approval(
3134 runtime: &GuardedToolRuntime<InMemoryApprovalCapabilityStore>,
3135 invocation: ToolInvocation,
3136 grant: RunToolGrant,
3137 ) -> GuardedToolResult {
3138 let GuardedToolResult::ApprovalRequired { binding, .. } = runtime
3139 .invoke(
3140 invocation.clone(),
3141 grant.clone(),
3142 None,
3143 CancellationToken::new(),
3144 )
3145 .await
3146 else {
3147 panic!("MCP invocation bypassed exact Host approval");
3148 };
3149 let capability = HostApprovalIssuer::new(GATE_SIGNING_KEY)
3150 .unwrap()
3151 .issue(binding, i64::MAX)
3152 .unwrap();
3153 runtime
3154 .invoke(
3155 invocation,
3156 grant,
3157 Some(capability),
3158 CancellationToken::new(),
3159 )
3160 .await
3161 }
3162
3163 #[tokio::test]
3164 async fn one_thousand_mcp_calls_are_journaled_filtered_and_share_one_connection() {
3165 let factory = FaultFactory::new(FaultStage::Healthy);
3166 let mut config = fault_config(factory.clone(), Duration::from_secs(2));
3167 config.server_id = McpServerId::new("gate");
3168 config.enabled_tools =
3169 BTreeSet::from(["echo".to_owned(), "beta".to_owned(), "hidden".to_owned()]);
3170 config.disabled_tools = BTreeSet::from(["hidden".to_owned()]);
3171 let bounds = authority_bounds(factory.authority());
3172 let journal = Arc::new(InMemoryToolEffectJournalStore::default());
3173 let runtime = gate_runtime(bounds.clone(), journal.clone());
3174 let registry = McpToolsAdapterRegistry::register(
3175 runtime.as_ref(),
3176 vec![config],
3177 ToolRestriction {
3178 bounds: bounds.clone(),
3179 },
3180 CancellationToken::new(),
3181 )
3182 .await
3183 .unwrap();
3184 let manager = registry.manager(&McpServerId::new("gate")).unwrap();
3185
3186 assert_eq!(registry.tool_count(), 2);
3187 assert_eq!(
3188 manager
3189 .snapshot()
3190 .tools
3191 .iter()
3192 .map(|tool| tool.name.as_str())
3193 .collect::<BTreeSet<_>>(),
3194 BTreeSet::from(["beta", "echo"])
3195 );
3196 assert!(runtime
3197 .resolve_tool_id("mcp__gate__hidden")
3198 .unwrap()
3199 .is_none());
3200
3201 let mut committed = 0;
3202 for index in 0..GATE_CASES {
3203 let dispatches_before = factory.state.tool_call_dispatches.load(Ordering::SeqCst);
3204 for rejected_name in ["hidden", "stale-cached-tool"] {
3205 assert!(matches!(
3206 manager
3207 .invoke(rejected_name, json!({}), CancellationToken::new())
3208 .await,
3209 Err(GuardedMcpCallError::Rejected(_))
3210 ));
3211 }
3212 assert_eq!(
3213 factory.state.tool_call_dispatches.load(Ordering::SeqCst),
3214 dispatches_before
3215 );
3216
3217 let tool = if index % 2 == 0 { "echo" } else { "beta" };
3218 let run_id = RunId::new("m4-mcp-gate");
3219 let call_id = ToolCallId::new(format!("call-{index}"));
3220 let invocation = ToolInvocation {
3221 run_id: run_id.clone(),
3222 call_id: call_id.clone(),
3223 tool_id: ToolId::new(format!("mcp/gate/{tool}/v1")),
3224 arguments: json!({}),
3225 };
3226 let result = invoke_with_approval(
3227 runtime.as_ref(),
3228 invocation,
3229 RunToolGrant {
3230 bounds: bounds.clone(),
3231 },
3232 )
3233 .await;
3234 assert!(matches!(
3235 result,
3236 GuardedToolResult::Outcome {
3237 outcome: ToolOutcome::Completed {
3238 output: ToolOutput::Inline(_)
3239 },
3240 cached: false,
3241 }
3242 ));
3243
3244 let key = ToolEffectKey::new(run_id, call_id);
3245 let records = journal.load_effect(&key).await.unwrap();
3246 let projection = replay_tool_effect(&key, &records).unwrap().unwrap();
3247 assert!(matches!(
3248 projection.phase,
3249 ToolEffectPhase::Committed { .. }
3250 ));
3251 committed += 1;
3252 }
3253
3254 assert_eq!(
3255 factory.state.tool_call_dispatches.load(Ordering::SeqCst),
3256 committed
3257 );
3258 assert_eq!(factory.state.connects.load(Ordering::SeqCst), 1);
3259 assert_eq!(manager.connection_generation(), 1);
3260 assert_eq!(factory.state.active_connections.load(Ordering::SeqCst), 1);
3261 registry.shutdown().await;
3262 factory.state.assert_released();
3263 assert_eq!(factory.state.explicit_closes.load(Ordering::SeqCst), 1);
3264 }
3265
3266 #[tokio::test]
3267 async fn one_thousand_oversized_mcp_results_keep_only_verified_artifact_summaries() {
3268 let factory = FaultFactory::new(FaultStage::LargeResult);
3269 let mut config = fault_config(factory.clone(), Duration::from_secs(2));
3270 config.server_id = McpServerId::new("spill-gate");
3271 config.enabled_tools = BTreeSet::from(["echo".to_owned()]);
3272 let mut bounds = authority_bounds(factory.authority());
3273 bounds.max_output_bytes = Some(128);
3274 let journal = Arc::new(InMemoryToolEffectJournalStore::default());
3275 let artifacts =
3276 ToolArtifactStore::new(Arc::new(InMemoryBlobStore::default()), 128 * 1024, 96).unwrap();
3277 let runtime =
3278 gate_runtime_with_artifacts(bounds.clone(), journal.clone(), artifacts.clone());
3279 let registry = McpToolsAdapterRegistry::register(
3280 runtime.as_ref(),
3281 vec![config],
3282 ToolRestriction {
3283 bounds: bounds.clone(),
3284 },
3285 CancellationToken::new(),
3286 )
3287 .await
3288 .unwrap();
3289
3290 for index in 0..GATE_CASES {
3291 let run_id = RunId::new("m4-mcp-spill-gate");
3292 let call_id = ToolCallId::new(format!("spill-{index}"));
3293 let result = invoke_with_approval(
3294 runtime.as_ref(),
3295 ToolInvocation {
3296 run_id: run_id.clone(),
3297 call_id: call_id.clone(),
3298 tool_id: ToolId::new("mcp/spill-gate/echo/v1"),
3299 arguments: json!({}),
3300 },
3301 RunToolGrant {
3302 bounds: bounds.clone(),
3303 },
3304 )
3305 .await;
3306 let GuardedToolResult::Outcome {
3307 outcome:
3308 ToolOutcome::Completed {
3309 output: ToolOutput::Artifact(artifact),
3310 },
3311 cached: false,
3312 } = result
3313 else {
3314 panic!("oversized MCP output escaped artifact spill")
3315 };
3316 artifact.validate().unwrap();
3317 assert!(!artifact.summary.trim().is_empty());
3318 assert!(artifact.byte_size > 128);
3319 let bytes = artifacts.resolve(&artifact).await.unwrap();
3320 let resolved: Value = serde_json::from_slice(&bytes).unwrap();
3321 assert!(resolved["result"]["content"][0]["text"]
3322 .as_str()
3323 .is_some_and(|text| text.starts_with("large-mcp-result/")));
3324
3325 let key = ToolEffectKey::new(run_id, call_id);
3326 let projection = replay_tool_effect(&key, &journal.load_effect(&key).await.unwrap())
3327 .unwrap()
3328 .unwrap();
3329 assert!(matches!(
3330 projection.phase,
3331 ToolEffectPhase::Committed {
3332 outcome: ToolOutcome::Completed {
3333 output: ToolOutput::Artifact(_)
3334 },
3335 ..
3336 }
3337 ));
3338 }
3339 assert_eq!(
3340 factory.state.tool_call_dispatches.load(Ordering::SeqCst),
3341 GATE_CASES
3342 );
3343 registry.shutdown().await;
3344 factory.state.assert_released();
3345 }
3346
3347 #[tokio::test]
3348 async fn one_thousand_unauthorized_environment_and_credential_bindings_connect_zero_times() {
3349 let mut baseline_error = None;
3350 for index in 0..GATE_CASES {
3351 let authority = if index % 2 == 0 {
3352 McpTransportAuthority {
3353 kind: McpTransportKind::Stdio,
3354 binding_digest: Digest::sha256(format!("unauthorized-env-{index}")),
3355 effect_scopes: BTreeSet::from([
3356 EffectScope::Process,
3357 EffectScope::FilesystemRead,
3358 EffectScope::FilesystemWrite,
3359 EffectScope::ExternalSideEffect,
3360 EffectScope::SecretRead,
3361 ]),
3362 process_programs: BTreeSet::from(["/fault/mcp".to_owned()]),
3363 allow_child_processes: false,
3364 allow_host_ui: false,
3365 filesystem_read_roots: BTreeSet::from(["/fault/read".to_owned()]),
3366 filesystem_write_roots: BTreeSet::from(["/fault/write".to_owned()]),
3367 sandbox_profiles: BTreeSet::from([MCP_STDIO_SANDBOX_PROFILE.to_owned()]),
3368 network_targets: BTreeSet::new(),
3369 allow_unrestricted_network: false,
3370 environment_variables: BTreeSet::from([format!("SENTINEL_ENV_{index}")]),
3371 credential_references: BTreeSet::new(),
3372 }
3373 } else {
3374 McpTransportAuthority {
3375 kind: McpTransportKind::StreamableHttp,
3376 binding_digest: Digest::sha256(format!("unauthorized-credential-{index}")),
3377 effect_scopes: BTreeSet::from([
3378 EffectScope::Network,
3379 EffectScope::ExternalSideEffect,
3380 EffectScope::SecretRead,
3381 ]),
3382 process_programs: BTreeSet::new(),
3383 allow_child_processes: false,
3384 allow_host_ui: false,
3385 filesystem_read_roots: BTreeSet::new(),
3386 filesystem_write_roots: BTreeSet::new(),
3387 sandbox_profiles: BTreeSet::new(),
3388 network_targets: BTreeSet::from(["http://127.0.0.1/fault-mcp".to_owned()]),
3389 allow_unrestricted_network: false,
3390 environment_variables: BTreeSet::new(),
3391 credential_references: BTreeSet::from([format!("env:SENTINEL_TOKEN_{index}")]),
3392 }
3393 };
3394 authority.validate().unwrap();
3395 let factory = FaultFactory::with_authority(authority);
3396 let mut config = fault_config(factory.clone(), Duration::from_secs(1));
3397 config.server_id = McpServerId::new("authority-gate");
3398 let mut restriction_bounds = authority_bounds(factory.authority());
3399 restriction_bounds.environment.allowed_variables.clear();
3400 restriction_bounds.allowed_credentials.clear();
3401 let runtime = gate_runtime(
3402 restriction_bounds.clone(),
3403 Arc::new(InMemoryToolEffectJournalStore::default()),
3404 );
3405 let result = McpToolsAdapterRegistry::register(
3406 runtime.as_ref(),
3407 vec![config],
3408 ToolRestriction {
3409 bounds: restriction_bounds,
3410 },
3411 CancellationToken::new(),
3412 )
3413 .await;
3414 let error = match result {
3415 Err(McpToolsAdapterError::InvalidConfig(message)) => message,
3416 Err(other) => panic!("unexpected authority rejection: {other:?}"),
3417 Ok(registry) => {
3418 registry.shutdown().await;
3419 panic!("unauthorized MCP authority reached connect")
3420 }
3421 };
3422 if let Some(baseline) = &baseline_error {
3423 assert_eq!(&error, baseline);
3424 } else {
3425 baseline_error = Some(error);
3426 }
3427 assert_eq!(factory.state.connects.load(Ordering::SeqCst), 0);
3428 factory.state.assert_released();
3429 }
3430 }
3431
3432 #[tokio::test]
3433 async fn one_thousand_disconnects_reconnect_once_and_schema_drift_never_calls_stale_tools() {
3434 for stage in [FaultStage::ReconnectStable, FaultStage::SchemaChanged] {
3435 let mut baseline_schema_error = None;
3436 for _ in 0..GATE_CASES {
3437 let factory = FaultFactory::new(stage);
3438 let manager = connect_result(
3439 fault_config(factory.clone(), Duration::from_secs(1)),
3440 CancellationToken::new(),
3441 )
3442 .await
3443 .unwrap();
3444 assert_unknown_effect(
3445 manager
3446 .invoke("echo", json!({}), CancellationToken::new())
3447 .await,
3448 );
3449 assert_eq!(manager.health(), McpServerHealth::Degraded);
3450 factory.state.assert_released();
3451 assert_eq!(factory.state.explicit_closes.load(Ordering::SeqCst), 1);
3452
3453 let second = manager
3454 .invoke("echo", json!({}), CancellationToken::new())
3455 .await;
3456 match stage {
3457 FaultStage::ReconnectStable => {
3458 assert!(second.is_ok());
3459 assert_eq!(manager.connection_generation(), 2);
3460 assert_eq!(factory.state.active_connections.load(Ordering::SeqCst), 1);
3461 assert_eq!(factory.state.tool_call_dispatches.load(Ordering::SeqCst), 2);
3462 manager.shutdown().await;
3463 assert_eq!(factory.state.explicit_closes.load(Ordering::SeqCst), 2);
3464 }
3465 FaultStage::SchemaChanged => {
3466 let message = match second {
3467 Err(GuardedMcpCallError::Failed(message)) => message,
3468 _ => panic!("schema drift did not fail closed"),
3469 };
3470 if let Some(baseline) = &baseline_schema_error {
3471 assert_eq!(&message, baseline);
3472 } else {
3473 baseline_schema_error = Some(message);
3474 }
3475 assert_eq!(manager.connection_generation(), 1);
3476 assert_eq!(factory.state.tool_call_dispatches.load(Ordering::SeqCst), 1);
3477 assert_eq!(factory.state.explicit_closes.load(Ordering::SeqCst), 2);
3478 }
3479 _ => unreachable!(),
3480 }
3481 assert_eq!(factory.state.connects.load(Ordering::SeqCst), 2);
3482 factory.state.assert_released();
3483 }
3484 }
3485 }
3486
3487 #[tokio::test]
3488 async fn one_thousand_required_start_and_name_conflicts_are_deterministic_and_clean() {
3489 let start_factory = FaultFactory::new(FaultStage::Connect);
3490 let start_config = fault_config(start_factory.clone(), Duration::from_millis(1));
3491 let start_bounds = authority_bounds(start_factory.authority());
3492 let start_runtime = gate_runtime(
3493 start_bounds.clone(),
3494 Arc::new(InMemoryToolEffectJournalStore::default()),
3495 );
3496 let mut baseline_start_error = None;
3497 for _ in 0..GATE_CASES {
3498 let result = McpToolsAdapterRegistry::register(
3499 start_runtime.as_ref(),
3500 vec![start_config.clone()],
3501 ToolRestriction {
3502 bounds: start_bounds.clone(),
3503 },
3504 CancellationToken::new(),
3505 )
3506 .await;
3507 let error = match result {
3508 Err(error) => error.to_string(),
3509 Ok(registry) => {
3510 registry.shutdown().await;
3511 panic!("required MCP startup failure was silently skipped")
3512 }
3513 };
3514 if let Some(baseline) = &baseline_start_error {
3515 assert_eq!(&error, baseline);
3516 } else {
3517 baseline_start_error = Some(error);
3518 }
3519 start_factory.state.assert_released();
3520 }
3521 assert_eq!(
3522 start_factory.state.connects.load(Ordering::SeqCst),
3523 GATE_CASES
3524 );
3525
3526 let conflict_factory = FaultFactory::new(FaultStage::NameConflict);
3527 let conflict_config = fault_config(conflict_factory.clone(), Duration::from_secs(1));
3528 let conflict_bounds = authority_bounds(conflict_factory.authority());
3529 let conflict_runtime = gate_runtime(
3530 conflict_bounds.clone(),
3531 Arc::new(InMemoryToolEffectJournalStore::default()),
3532 );
3533 let mut baseline_conflict_error = None;
3534 for _ in 0..GATE_CASES {
3535 let result = McpToolsAdapterRegistry::register(
3536 conflict_runtime.as_ref(),
3537 vec![conflict_config.clone()],
3538 ToolRestriction {
3539 bounds: conflict_bounds.clone(),
3540 },
3541 CancellationToken::new(),
3542 )
3543 .await;
3544 let error = match result {
3545 Err(McpToolsAdapterError::Conflict(message)) => message,
3546 Err(other) => panic!("unexpected name conflict result: {other:?}"),
3547 Ok(registry) => {
3548 registry.shutdown().await;
3549 panic!("MCP name collision was silently registered")
3550 }
3551 };
3552 if let Some(baseline) = &baseline_conflict_error {
3553 assert_eq!(&error, baseline);
3554 } else {
3555 baseline_conflict_error = Some(error);
3556 }
3557 conflict_factory.state.assert_released();
3558 }
3559 assert_eq!(
3560 conflict_factory.state.connects.load(Ordering::SeqCst),
3561 GATE_CASES
3562 );
3563 assert_eq!(
3564 conflict_factory
3565 .state
3566 .explicit_closes
3567 .load(Ordering::SeqCst),
3568 GATE_CASES
3569 );
3570 }
3571
3572 #[tokio::test]
3573 async fn one_thousand_deadlines_per_mcp_stage_never_escape_or_leak_a_session() {
3574 let startup_stages = [
3575 FaultStage::Connect,
3576 FaultStage::Discover,
3577 FaultStage::Initialize,
3578 FaultStage::List,
3579 ];
3580 for stage in startup_stages {
3581 let factory = FaultFactory::new(stage);
3582 let config = fault_config(factory.clone(), Duration::from_millis(1));
3583 let mut baseline_error = None;
3584 for _ in 0..GATE_CASES {
3585 let error = match connect_result(config.clone(), CancellationToken::new()).await {
3586 Ok(manager) => {
3587 manager.shutdown().await;
3588 panic!("fault stage {stage:?} escaped its deadline")
3589 }
3590 Err(error) => error.to_string(),
3591 };
3592 if let Some(baseline) = &baseline_error {
3593 assert_eq!(&error, baseline);
3594 } else {
3595 baseline_error = Some(error);
3596 }
3597 factory.state.assert_released();
3598 }
3599 assert_eq!(factory.state.connects.load(Ordering::SeqCst), GATE_CASES);
3600 }
3601
3602 for stage in [FaultStage::Call, FaultStage::BodyDecode] {
3603 let factory = FaultFactory::new(stage);
3604 let config = fault_config(factory.clone(), Duration::from_millis(1));
3605 let mut baseline_error = None;
3606 for _ in 0..GATE_CASES {
3607 let manager = connect_result(config.clone(), CancellationToken::new())
3608 .await
3609 .unwrap();
3610 let result = manager
3611 .invoke("echo", json!({}), CancellationToken::new())
3612 .await;
3613 let message = match &result {
3614 Err(GuardedMcpCallError::UnknownEffect(message)) => message.clone(),
3615 _ => panic!("fault stage {stage:?} did not become UnknownEffect"),
3616 };
3617 if let Some(baseline) = &baseline_error {
3618 assert_eq!(&message, baseline);
3619 } else {
3620 baseline_error = Some(message);
3621 }
3622 assert_unknown_effect(result);
3623 assert_eq!(manager.health(), McpServerHealth::Degraded);
3624 assert_eq!(
3625 factory
3626 .state
3627 .completed_call_responses
3628 .load(Ordering::SeqCst),
3629 0
3630 );
3631 factory.state.assert_released();
3632 }
3633 assert_eq!(factory.state.connects.load(Ordering::SeqCst), GATE_CASES);
3634 }
3635 }
3636
3637 #[tokio::test]
3638 async fn one_thousand_connect_read_and_call_cancellations_release_within_one_second() {
3639 let mut connect_latencies = Vec::with_capacity(GATE_CASES);
3640 let mut read_latencies = Vec::with_capacity(GATE_CASES);
3641 let mut call_latencies = Vec::with_capacity(GATE_CASES);
3642
3643 for _ in 0..GATE_CASES {
3644 let factory = FaultFactory::new(FaultStage::Connect);
3645 let config = fault_config(factory.clone(), Duration::from_secs(30));
3646 let cancellation = CancellationToken::new();
3647 let task_cancellation = cancellation.clone();
3648 let task = tokio::spawn(async move { connect_result(config, task_cancellation).await });
3649 factory.state.entered.notified().await;
3650 let started = Instant::now();
3651 cancellation.cancel();
3652 let result = timeout(Duration::from_secs(1), task)
3653 .await
3654 .expect("MCP connect cancellation exceeded one second")
3655 .unwrap();
3656 assert!(matches!(result, Err(McpToolsAdapterError::Cancelled)));
3657 connect_latencies.push(started.elapsed());
3658 factory.state.assert_released();
3659 }
3660
3661 for _ in 0..GATE_CASES {
3662 let factory = FaultFactory::new(FaultStage::List);
3663 let config = fault_config(factory.clone(), Duration::from_secs(30));
3664 let cancellation = CancellationToken::new();
3665 let task_cancellation = cancellation.clone();
3666 let task = tokio::spawn(async move { connect_result(config, task_cancellation).await });
3667 factory.state.entered.notified().await;
3668 let started = Instant::now();
3669 cancellation.cancel();
3670 let result = timeout(Duration::from_secs(1), task)
3671 .await
3672 .expect("MCP read cancellation exceeded one second")
3673 .unwrap();
3674 assert!(matches!(result, Err(McpToolsAdapterError::Cancelled)));
3675 read_latencies.push(started.elapsed());
3676 factory.state.assert_released();
3677 }
3678
3679 for _ in 0..GATE_CASES {
3680 let factory = FaultFactory::new(FaultStage::Call);
3681 let manager = connect_result(
3682 fault_config(factory.clone(), Duration::from_secs(30)),
3683 CancellationToken::new(),
3684 )
3685 .await
3686 .unwrap();
3687 let cancellation = CancellationToken::new();
3688 let task_cancellation = cancellation.clone();
3689 let task =
3690 tokio::spawn(
3691 async move { manager.invoke("echo", json!({}), task_cancellation).await },
3692 );
3693 factory.state.entered.notified().await;
3694 let started = Instant::now();
3695 cancellation.cancel();
3696 let result = timeout(Duration::from_secs(1), task)
3697 .await
3698 .expect("MCP call cancellation exceeded one second")
3699 .unwrap();
3700 assert_unknown_effect(result);
3701 call_latencies.push(started.elapsed());
3702 assert_eq!(
3703 factory
3704 .state
3705 .completed_call_responses
3706 .load(Ordering::SeqCst),
3707 0
3708 );
3709 factory.state.assert_released();
3710 }
3711
3712 assert!(p99(&mut connect_latencies) <= Duration::from_secs(1));
3713 assert!(p99(&mut read_latencies) <= Duration::from_secs(1));
3714 assert!(p99(&mut call_latencies) <= Duration::from_secs(1));
3715 }
3716
3717 fn p99(latencies: &mut [Duration]) -> Duration {
3718 latencies.sort_unstable();
3719 latencies[(latencies.len() * 99).div_ceil(100) - 1]
3720 }
3721}