(&self, name: &str)
| 82 | } |
| 83 | |
| 84 | async fn do_connect(&self, name: &str) -> Result<()> { |
| 85 | // Get config |
| 86 | let config = { |
| 87 | let configs = self.configs.read().await; |
| 88 | configs |
| 89 | .get(name) |
| 90 | .cloned() |
| 91 | .ok_or_else(|| anyhow!("MCP server not found: {}", name))? |
| 92 | }; |
| 93 | |
| 94 | if !config.enabled { |
| 95 | return Err(anyhow!("MCP server is disabled: {}", name)); |
| 96 | } |
| 97 | |
| 98 | // Resolve OAuth token into an Authorization header (if configured) |
| 99 | let auth_header = Self::resolve_auth_header(config.oauth.as_ref()).await?; |
| 100 | |
| 101 | // Create transport based on config |
| 102 | let transport: Arc<dyn McpTransport> = match &config.transport { |
| 103 | McpTransportConfig::Stdio { command, args } => Arc::new( |
| 104 | StdioTransport::spawn_with_timeout( |
| 105 | command, |
| 106 | args, |
| 107 | &config.env, |
| 108 | config.tool_timeout_secs, |
| 109 | ) |
| 110 | .await?, |
| 111 | ), |
| 112 | McpTransportConfig::Http { url, headers } => { |
| 113 | let mut merged = headers.clone(); |
| 114 | if let Some((k, v)) = &auth_header { |
| 115 | merged.insert(k.clone(), v.clone()); |
| 116 | } |
| 117 | Arc::new( |
| 118 | HttpSseTransport::connect_with_timeout(url, merged, config.tool_timeout_secs) |
| 119 | .await?, |
| 120 | ) |
| 121 | } |
| 122 | McpTransportConfig::StreamableHttp { url, headers } => { |
| 123 | let mut merged = headers.clone(); |
| 124 | if let Some((k, v)) = &auth_header { |
| 125 | merged.insert(k.clone(), v.clone()); |
| 126 | } |
| 127 | Arc::new( |
| 128 | StreamableHttpTransport::connect_with_timeout( |
| 129 | url, |
| 130 | merged, |
| 131 | config.tool_timeout_secs, |
| 132 | ) |
| 133 | .await?, |
| 134 | ) |
| 135 | } |
| 136 | }; |
| 137 | |
| 138 | // Create client |
| 139 | let client = Arc::new(McpClient::new(name.to_string(), transport)); |
| 140 | |
| 141 | // Initialize |
no test coverage detected