Skip to main content

orchestral_runtime/tools/
mcp.rs

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/// Immutable Host filesystem boundary for one stdio MCP server process.
236#[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    /// Builds one explicit local-MCP process boundary. The executable itself
264    /// is bound separately by [`StdioMcpTransportFactory`]; these roots cover
265    /// only the server's data access and never widen generic command Tools.
266    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    /// Allows launchers such as npx/uvx/sh to form a process tree. Every
285    /// descendant remains confined by this server's filesystem/network policy.
286    pub fn with_child_processes(mut self, allow: bool) -> Self {
287        self.allow_child_processes = allow;
288        self
289    }
290
291    /// Allows this registered MCP transport to invoke the operating system's
292    /// URL/application opener. The MCP still owns the protocol (for example,
293    /// OAuth); the Host only materializes the declared OS capability.
294    pub fn with_host_ui(mut self, allow: bool) -> Self {
295        self.allow_host_ui = allow;
296        self
297    }
298
299    /// Uses the Host network for an explicitly registered local MCP process.
300    /// This is transport trust established by configuration, not authority
301    /// supplied by a model Tool call.
302    pub fn with_unrestricted_network(mut self, allow: bool) -> Self {
303        self.allow_unrestricted_network = allow;
304        self
305    }
306
307    /// Supplies HOME/TMP-style process state without inheriting the user's
308    /// ambient home directory. The directory must remain inside one declared
309    /// writable root after canonicalization.
310    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/// Built-in stdio transport factory. External transports implement the core
394/// [`McpTransportFactory`] SPI from a plugin and are injected by the Host.
395#[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/// One explicitly configured MCP Tool provider. Discovery and invocation use
551/// the same immutable transport authority for the lifetime of this registry.
552#[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
683/// Owns exactly one negotiated transport handle for all Tools on one MCP server.
684pub 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            // A legacy stdio server may wait forever for initialize instead of
907            // rejecting server/discover. Restart it before legacy fallback.
908            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
1190/// Keeps server managers alive after their thin Tool executors are registered.
1191pub 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
1385/// Narrows the aggregate Host MCP ceiling to one immutable server binding.
1386/// A remote transport therefore cannot inherit a local server's filesystem or
1387/// process authority, and local servers cannot inherit each other's roots.
1388fn 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
1447/// Describes the business operation initiated through an already-authorized
1448/// MCP transport. Process launch, transport cache roots, environment, and
1449/// credentials belong to the immutable server binding and must not be copied
1450/// into every user-facing Tool approval request.
1451fn 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            // Like Codex's server + Tool session key, but sealed to the full
1519            // discovered schema and annotations so a changed Tool contract
1520            // cannot inherit an earlier decision.
1521            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        // This watcher outlives a dropped executor future. GuardedToolRuntime
1575        // may finish the Run cancellation branch first, but the server process
1576        // must still be terminated and reaped.
1577        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
1675/// MCP Tool failures commonly carry actionable JSON inside text content. Keep
1676/// that payload readable for both the model and TUI instead of stringifying
1677/// the whole `CallToolResult`, which double-escapes the actual error.
1678fn 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        // MCP stdio servers receive only the Host-configured environment.
2049        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        // Use NDJSON (newline-delimited JSON) — compatible with all MCP servers.
2145        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        // Auto-detect: NDJSON (line = JSON) or LSP (Content-Length header).
2169        loop {
2170            let line = self.read_bounded_line().await?;
2171            let trimmed = line.trim();
2172            if trimmed.is_empty() {
2173                continue;
2174            }
2175            // If line starts with '{', it's NDJSON
2176            if trimmed.starts_with('{') {
2177                return serde_json::from_str::<Value>(trimmed)
2178                    .map_err(|err| format!("parse mcp NDJSON failed: {}", err));
2179            }
2180            // Otherwise treat as Content-Length header (LSP style)
2181            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                        // Read blank line after headers
2191                        let _ = self.read_bounded_line().await?;
2192                        // Read exact body
2193                        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        // SAFETY: this child was spawned as leader of a fresh process group.
2401        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}