1use mcp_utils::client::{
2 InMemoryServerSpec, McpClientEvent, McpConfig, McpConnectionDetails, McpError, McpManager, McpServer, McpTransport,
3 OAuthHandlerFactory, PROGRESSIVE_DISCOVERY_INSTRUCTION_NAME, ParseError, RuntimeMcpServer, RuntimeMcpTransport,
4 ToolFilter,
5};
6use mcp_utils::tool_gateway::{AETHER_MCP_IPC_SOCKET, UnixSocketMcpTransport, UnixSocketPath, UnixSocketServer};
7use utils::{SettingsStore, variables::Vars};
8
9use crate::agent_spec::McpConfigSource;
10use crate::core::AgentDeps;
11use crate::events::{AgentCommand, Command};
12
13use super::{
14 gateway_service::GatewayService,
15 mcp_handle::McpHandle,
16 run_mcp_task::{ManagerCommand, run_mcp_task},
17};
18use futures::future::BoxFuture;
19use rmcp::{RoleServer, service::DynService};
20use std::collections::{BTreeMap, BTreeSet, HashMap};
21use std::path::{Path, PathBuf};
22use std::sync::Arc;
23use tokio::{
24 sync::{
25 mpsc::{self, Receiver},
26 watch,
27 },
28 task::JoinHandle,
29};
30
31pub fn mcp(root_dir: impl AsRef<Path>) -> McpBuilder {
32 McpBuilder::new(root_dir)
33}
34
35#[derive(Clone)]
36pub struct RuntimeServices {
37 pub mcp: McpHandle,
38 pub root_dir: PathBuf,
39 pub agent_deps: AgentDeps,
40 pub shell_environment: BTreeMap<String, String>,
41}
42
43pub type ServerFactory = Box<
44 dyn Fn(InMemoryServerSpec, RuntimeServices) -> BoxFuture<'static, Box<dyn DynService<RoleServer>>> + Send + Sync,
45>;
46
47pub struct McpRuntime {
49 mcp: McpHandle,
50 handle: JoinHandle<()>,
51 agent_sync_handle: Option<JoinHandle<()>>,
52 gateway: Option<UnixSocketServer>,
53}
54
55impl McpRuntime {
56 pub fn handle(&self) -> &McpHandle {
57 &self.mcp
58 }
59
60 pub fn gateway_endpoint(&self) -> Option<&Path> {
61 self.gateway.as_ref().map(UnixSocketServer::path)
62 }
63}
64
65impl Drop for McpRuntime {
66 fn drop(&mut self) {
67 self.handle.abort();
68 if let Some(handle) = &self.agent_sync_handle {
69 handle.abort();
70 }
71 }
72}
73
74pub struct McpSession {
80 runtime: McpRuntime,
81 event_rx: Receiver<McpClientEvent>,
82}
83
84impl McpSession {
85 pub fn handle(&self) -> &McpHandle {
86 self.runtime.handle()
87 }
88
89 pub fn gateway_endpoint(&self) -> Option<&Path> {
90 self.runtime.gateway_endpoint()
91 }
92
93 pub async fn connect_agent(mut self, agent_tx: mpsc::Sender<Command>) -> Self {
96 assert!(self.runtime.agent_sync_handle.is_none(), "an MCP session can only connect one agent");
97 let mut snapshots = self.runtime.handle().subscribe();
98 let initial = snapshots.borrow_and_update().clone();
99 let mut previous_tools = initial.tool_definitions();
100 let mut previous_instructions = initial.model_instructions();
101 if agent_tx.send(Command::agent(AgentCommand::UpdateTools(previous_tools.clone()))).await.is_err() {
102 return self;
103 }
104 for (server, body) in &previous_instructions {
105 if agent_tx
106 .send(Command::agent(AgentCommand::UpdateMcpInstructions {
107 server: server.clone(),
108 body: Some(body.clone()),
109 }))
110 .await
111 .is_err()
112 {
113 return self;
114 }
115 }
116
117 let agent_tx = agent_tx.downgrade();
118 self.runtime.agent_sync_handle = Some(tokio::spawn(async move {
119 while snapshots.changed().await.is_ok() {
120 let Some(agent_tx) = agent_tx.upgrade() else {
121 break;
122 };
123 let snapshot = snapshots.borrow_and_update().clone();
124 let tools = snapshot.tool_definitions();
125 if tools != previous_tools {
126 if agent_tx.send(Command::agent(AgentCommand::UpdateTools(tools.clone()))).await.is_err() {
127 break;
128 }
129 previous_tools = tools;
130 }
131
132 let instructions = snapshot.model_instructions();
133 let servers = previous_instructions.keys().chain(instructions.keys()).cloned().collect::<BTreeSet<_>>();
134 for server in servers {
135 let previous = previous_instructions.get(&server);
136 let next = instructions.get(&server);
137 if previous != next
138 && agent_tx
139 .send(Command::agent(AgentCommand::UpdateMcpInstructions { server, body: next.cloned() }))
140 .await
141 .is_err()
142 {
143 return;
144 }
145 }
146 previous_instructions = instructions;
147 }
148 }));
149 self
150 }
151
152 pub async fn block_until_ready(&mut self) -> Option<McpConnectionDetails> {
156 while let Some(event) = self.event_rx.recv().await {
157 if let McpClientEvent::ConnectionReady(snapshot) = event {
158 return Some(snapshot);
159 }
160 }
161 None
162 }
163
164 pub fn split(self) -> (McpRuntime, Receiver<McpClientEvent>) {
165 (self.runtime, self.event_rx)
166 }
167}
168
169pub struct McpBuilder {
170 servers: Vec<McpServer>,
171 factories: HashMap<String, ServerFactory>,
172 mcp_channel_capacity: usize,
173 root_dir: PathBuf,
174 oauth_handler_factory: Option<OAuthHandlerFactory>,
175 agent_deps: AgentDeps,
176 aether_home: Option<PathBuf>,
177 vars: Vars,
178 tool_filter: ToolFilter,
179 progressive_discovery_instructions: Option<String>,
180}
181
182impl McpBuilder {
183 pub fn new(root_dir: impl AsRef<Path>) -> Self {
184 let mut vars = Vars::new().with("WORKSPACE", root_dir.as_ref().to_string_lossy().into_owned());
185
186 if let Some(store) = SettingsStore::new("AETHER_HOME", ".aether") {
187 vars.insert("AETHER_HOME", store.home().to_string_lossy().into_owned());
188 }
189
190 Self {
191 servers: Vec::new(),
192 factories: HashMap::new(),
193 mcp_channel_capacity: 1000,
194 root_dir: root_dir.as_ref().to_path_buf(),
195 oauth_handler_factory: None,
196 agent_deps: AgentDeps::default(),
197 aether_home: None,
198 vars,
199 tool_filter: ToolFilter::default(),
200 progressive_discovery_instructions: None,
201 }
202 }
203
204 pub fn with_servers(mut self, servers: Vec<McpServer>) -> Self {
205 self.servers.extend(servers);
206 self
207 }
208
209 pub fn with_tool_filter(mut self, filter: ToolFilter) -> Self {
210 self.tool_filter = filter;
211 self
212 }
213
214 pub fn with_progressive_discovery_instructions(mut self, instructions: impl Into<String>) -> Self {
215 self.progressive_discovery_instructions = Some(instructions.into());
216 self
217 }
218
219 pub fn register_in_memory_server(mut self, name: impl Into<String>, factory: ServerFactory) -> Self {
220 self.factories.insert(name.into(), factory);
221 self
222 }
223
224 pub fn root_dir(&self) -> &Path {
225 &self.root_dir
226 }
227
228 pub fn agent_deps(&self) -> AgentDeps {
231 self.agent_deps.clone()
232 }
233
234 pub fn with_agent_deps(mut self, deps: AgentDeps) -> Self {
235 self.agent_deps = deps;
236 self
237 }
238
239 pub fn with_oauth_handler_factory(mut self, factory: OAuthHandlerFactory) -> Self {
240 self.oauth_handler_factory = Some(factory);
241 self
242 }
243
244 pub fn with_aether_home(mut self, aether_home: impl Into<PathBuf>) -> Self {
245 let aether_home = aether_home.into();
246 self.vars.insert("AETHER_HOME", aether_home.to_string_lossy().into_owned());
247 self.aether_home = Some(aether_home);
248 self
249 }
250
251 pub fn from_json_files<T: AsRef<Path>>(mut self, paths: &[T]) -> Result<Self, ParseError> {
252 if paths.is_empty() {
253 return Ok(self);
254 }
255 let raw = McpConfig::from_json_files(paths)?;
256 self.servers.extend(raw.into_servers(&self.vars)?);
257 Ok(self)
258 }
259
260 pub fn from_mcp_config_sources(mut self, sources: &[McpConfigSource]) -> Result<Self, ParseError> {
261 if sources.is_empty() {
262 return Ok(self);
263 }
264
265 let mut merged = McpConfig::default();
266 for source in sources {
267 let config = match source {
268 McpConfigSource::File { path, defer_tools } => {
269 let mut config = McpConfig::from_json_file(path)?;
270 if *defer_tools {
271 config.defer_all_tools();
272 }
273 config
274 }
275 McpConfigSource::Json(json) => McpConfig::from_json(json)?,
276 McpConfigSource::Inline(config) => config.clone(),
277 };
278 merged.servers.extend(config.servers);
279 }
280
281 self.servers.extend(merged.into_servers(&self.vars)?);
282 Ok(self)
283 }
284
285 pub async fn spawn(self) -> Result<McpSession, McpError> {
286 let McpBuilder {
287 servers,
288 factories,
289 mcp_channel_capacity,
290 root_dir,
291 oauth_handler_factory,
292 agent_deps,
293 aether_home: _,
294 vars: _,
295 tool_filter,
296 progressive_discovery_instructions,
297 } = self;
298 if servers.iter().any(|server| server.tool_exposure.has_deferred_tools())
299 && servers.iter().any(|server| server.name == PROGRESSIVE_DISCOVERY_INSTRUCTION_NAME)
300 {
301 return Err(McpError::ReservedServerName(PROGRESSIVE_DISCOVERY_INSTRUCTION_NAME.to_string()));
302 }
303 let (manager_tx, manager_rx) = mpsc::channel::<ManagerCommand>(mcp_channel_capacity);
304 let (snapshot_tx, snapshot_rx) = watch::channel(Arc::new(mcp_utils::client::McpSnapshot::default()));
305 let (event_tx, event_rx) = mpsc::channel::<McpClientEvent>(mcp_channel_capacity);
306 let mcp = McpHandle::new(manager_tx, snapshot_rx);
307 let gateway_transport = if servers.iter().any(|server| server.tool_exposure.has_deferred_tools()) {
308 let path = UnixSocketPath::new().map_err(|error| McpError::TransportError(error.to_string()))?;
309 Some(UnixSocketMcpTransport::bind(path).map_err(|error| McpError::TransportError(error.to_string()))?)
310 } else {
311 None
312 };
313 let shell_environment = gateway_transport
314 .as_ref()
315 .map(|transport| {
316 BTreeMap::from([(AETHER_MCP_IPC_SOCKET.to_string(), transport.path().to_string_lossy().into_owned())])
317 })
318 .unwrap_or_default();
319 let services = RuntimeServices { mcp: mcp.clone(), root_dir: root_dir.clone(), agent_deps, shell_environment };
320 let servers = resolve_servers(servers, &factories, &services).await?;
321
322 let mut mcp_manager = McpManager::new(event_tx, oauth_handler_factory)
323 .with_tool_filter(tool_filter)
324 .with_snapshot_sender(snapshot_tx);
325 if let Some(capabilities) = services.agent_deps.mcp_client_capabilities.clone() {
326 mcp_manager = mcp_manager.with_client_capabilities(capabilities);
327 }
328 if let Some(instructions) = progressive_discovery_instructions {
329 mcp_manager = mcp_manager.with_progressive_discovery_instructions(instructions);
330 }
331 if let Some(store) = services.agent_deps.oauth_credential_store.clone() {
332 mcp_manager = mcp_manager.with_oauth_credential_store(store);
333 }
334 mcp_manager = mcp_manager.with_root_dir(root_dir);
335 let pending = mcp_manager.register_pending(servers).await?;
336 let task = tokio::spawn(run_mcp_task(mcp_manager, manager_rx, pending));
337 let gateway = gateway_transport.map(|transport| transport.spawn(GatewayService::new(mcp.clone())));
338
339 Ok(McpSession { runtime: McpRuntime { mcp, handle: task, agent_sync_handle: None, gateway }, event_rx })
340 }
341}
342
343async fn resolve_servers(
344 servers: Vec<McpServer>,
345 factories: &HashMap<String, ServerFactory>,
346 services: &RuntimeServices,
347) -> Result<Vec<RuntimeMcpServer>, McpError> {
348 let mut resolved = Vec::with_capacity(servers.len());
349 for McpServer { name, transport, tool_exposure } in servers {
350 let transport = match transport {
351 McpTransport::Stdio { command, args, env } => RuntimeMcpTransport::Stdio { command, args, env },
352 McpTransport::Http(config) => RuntimeMcpTransport::Http(config),
353 McpTransport::InMemory { spec } => {
354 let factory = factories.get(&spec.factory).ok_or_else(|| McpError::InMemoryFactoryNotFound {
355 server: name.clone(),
356 factory: spec.factory.clone(),
357 })?;
358 RuntimeMcpTransport::InMemory { server: factory(spec, services.clone()).await }
359 }
360 };
361 resolved.push(RuntimeMcpServer::new(name, transport, tool_exposure));
362 }
363 Ok(resolved)
364}
365
366#[cfg(test)]
367mod tests {
368 use super::*;
369 use aether_auth::{FakeOAuthCredentialStore, OAuthCredentialStorage};
370 use futures::FutureExt;
371 use mcp_utils::testing::FakeMcpServer;
372 use mcp_utils::{
373 client::{McpServerConfig, McpTransport, StdioServerConfig, StdioType, ToolExposure},
374 status::McpServerStatus,
375 };
376 use std::collections::{BTreeMap, HashMap};
377 use std::sync::atomic::{AtomicUsize, Ordering};
378 use std::sync::{Arc, Mutex};
379
380 fn write_config_file(name: &str, json: &str) -> (tempfile::TempDir, PathBuf) {
381 let dir = tempfile::tempdir().unwrap();
382 let path = dir.path().join(name);
383 std::fs::write(&path, json).unwrap();
384 (dir, path)
385 }
386
387 fn json_source(json: &str) -> McpConfigSource {
388 McpConfigSource::Json(json.to_string())
389 }
390
391 fn builder_from_sources(sources: &[McpConfigSource]) -> McpBuilder {
392 McpBuilder::new("/workspace").from_mcp_config_sources(sources).unwrap()
393 }
394
395 #[tokio::test]
396 async fn in_memory_factory_runs_once_at_spawn_with_runtime_services() {
397 let calls = Arc::new(AtomicUsize::new(0));
398 let received = Arc::new(Mutex::new(None::<RuntimeServices>));
399 let factory_calls = Arc::clone(&calls);
400 let factory_received = Arc::clone(&received);
401 let oauth_store: Arc<dyn OAuthCredentialStorage> = Arc::new(FakeOAuthCredentialStore::new());
402 let deps = AgentDeps::new(Arc::clone(&oauth_store), None);
403 let factory: ServerFactory = Box::new(move |spec, services| {
404 factory_calls.fetch_add(1, Ordering::SeqCst);
405 assert_eq!(spec.args, ["--root", "/workspace/tools"]);
406 assert_eq!(spec.input, Some(serde_json::json!({"enabled": true})));
407 *factory_received.lock().unwrap() = Some(services);
408 async move { FakeMcpServer::new().into_dyn() }.boxed()
409 });
410
411 let builder = McpBuilder::new("/workspace")
412 .with_agent_deps(deps)
413 .register_in_memory_server("test", factory)
414 .from_mcp_config_sources(&[json_source(
415 r#"{"servers":{"test":{"type":"in-memory","args":["--root","${WORKSPACE}/tools"],"input":{"enabled":true}}}}"#,
416 )])
417 .unwrap();
418 assert_eq!(calls.load(Ordering::SeqCst), 0);
419
420 let spawn = builder.spawn().await.unwrap();
421 assert_eq!(calls.load(Ordering::SeqCst), 1);
422 let services = received.lock().unwrap().clone().expect("factory received runtime services");
423 assert_eq!(services.root_dir, PathBuf::from("/workspace"));
424 assert!(Arc::ptr_eq(&services.mcp.snapshot(), &spawn.handle().snapshot()));
425 assert!(Arc::ptr_eq(
426 services.agent_deps.oauth_credential_store.as_ref().expect("factory received agent dependencies"),
427 &oauth_store,
428 ));
429 assert!(services.shell_environment.is_empty());
430 }
431
432 #[tokio::test]
433 async fn deferred_gateway_is_bound_before_in_memory_factories_run() {
434 let received = Arc::new(Mutex::new(None::<RuntimeServices>));
435 let factory_received = Arc::clone(&received);
436 let factory: ServerFactory = Box::new(move |_, services| {
437 *factory_received.lock().unwrap() = Some(services);
438 async move { FakeMcpServer::new().into_dyn() }.boxed()
439 });
440 let spawn = McpBuilder::new("/workspace")
441 .register_in_memory_server("test", factory)
442 .from_mcp_config_sources(&[json_source(r#"{"servers":{"test":{"type":"in-memory","deferTools":true}}}"#)])
443 .unwrap()
444 .spawn()
445 .await
446 .unwrap();
447
448 let services = received.lock().unwrap().clone().expect("factory received runtime services");
449 let inherited =
450 services.shell_environment.get(AETHER_MCP_IPC_SOCKET).expect("factory receives gateway endpoint");
451 assert_eq!(Path::new(inherited), spawn.gateway_endpoint().expect("gateway endpoint exists"));
452 assert!(Path::new(inherited).exists());
453 }
454
455 #[tokio::test]
456 async fn snapshots_are_immutable_and_watch_observes_connection_changes() {
457 let factory: ServerFactory = Box::new(|_, _| async move { FakeMcpServer::new().into_dyn() }.boxed());
458 let mut spawn = McpBuilder::new("/workspace")
459 .register_in_memory_server("test", factory)
460 .from_mcp_config_sources(&[json_source(r#"{"servers":{"test":{"type":"in-memory"}}}"#)])
461 .unwrap()
462 .spawn()
463 .await
464 .unwrap();
465 let old = spawn.handle().snapshot();
466 let mut updates = spawn.handle().subscribe();
467
468 let ready = spawn.block_until_ready().await.expect("bootstrap completes");
469 updates.changed().await.expect("connection publishes a snapshot");
470 let observed = updates.borrow().clone();
471
472 assert!(old.tool_definitions().is_empty());
473 assert_eq!(ready.tool_definitions()[0].name, "test__add_numbers");
474 assert_eq!(observed.tool_definitions(), ready.tool_definitions());
475 assert!(!Arc::ptr_eq(&old, &ready));
476 }
477
478 #[tokio::test]
479 async fn missing_in_memory_factory_fails_at_spawn_with_server_and_factory() {
480 let builder = McpBuilder::new("/workspace")
481 .from_mcp_config_sources(&[json_source(r#"{"servers":{"custom":{"type":"in-memory"}}}"#)])
482 .unwrap();
483
484 let Err(error) = builder.spawn().await else {
485 panic!("spawn should reject an unregistered factory");
486 };
487 assert!(matches!(
488 error,
489 McpError::InMemoryFactoryNotFound { ref server, ref factory }
490 if server == "custom" && factory == "custom"
491 ));
492 }
493
494 #[tokio::test]
495 async fn mixed_direct_sources_preserve_last_wins_order() {
496 let (_dir, file_path) =
497 write_config_file("mcp.json", r#"{"servers":{"coding":{"type":"stdio","command":"from_file"}}}"#);
498 let inline = McpConfig::new(BTreeMap::from([(
499 "coding".to_string(),
500 McpServerConfig::Stdio(StdioServerConfig {
501 type_: StdioType::Stdio,
502 command: "from_inline".to_string(),
503 args: Vec::new(),
504 env: HashMap::new(),
505 defer_tools: ToolExposure::ModelVisible,
506 }),
507 )]));
508 let sources = vec![
509 McpConfigSource::model_visible(file_path),
510 json_source(r#"{"servers":{"coding":{"type":"stdio","command":"from_json"}}}"#),
511 McpConfigSource::Inline(inline),
512 ];
513
514 let builder = builder_from_sources(&sources);
515
516 assert_eq!(command_for(&builder, "coding"), Some("from_inline"));
517 assert_eq!(deferred_tools_for(&builder, "coding"), Some(false));
518 }
519
520 #[tokio::test]
521 async fn file_sources_keep_their_position_relative_to_json_sources() {
522 let (_dir, file_path) =
523 write_config_file("mcp.json", r#"{"servers":{"coding":{"type":"stdio","command":"from_file"}}}"#);
524 let sources = vec![
525 json_source(r#"{"servers":{"coding":{"type":"stdio","command":"from_json"}}}"#),
526 McpConfigSource::model_visible(file_path),
527 ];
528
529 let builder = builder_from_sources(&sources);
530
531 assert_eq!(command_for(&builder, "coding"), Some("from_file"));
532 }
533
534 #[tokio::test]
535 async fn file_source_defer_tools_marks_all_file_servers_deferred() {
536 let (_dir, file_path) = write_config_file(
537 "deferred.json",
538 r#"{"servers":{"github":{"type":"stdio","command":"g","deferTools":{"exclude":["status"]}},"browser":{"type":"stdio","command":"b"}}}"#,
539 );
540
541 let builder = McpBuilder::new("/workspace")
542 .from_mcp_config_sources(&[McpConfigSource::File { path: file_path, defer_tools: true }])
543 .unwrap();
544
545 assert_eq!(deferred_tools_for(&builder, "github"), Some(true));
546 assert_eq!(deferred_tools_for(&builder, "browser"), Some(true));
547 assert!(is_direct_tool(&builder, "github", "status"));
548 }
549
550 #[tokio::test]
551 async fn later_sources_override_defer_tools_flag() {
552 let (_dir, file_path) =
553 write_config_file("deferred.json", r#"{"servers":{"coding":{"type":"stdio","command":"from_file"}}}"#);
554 let sources = vec![
555 McpConfigSource::File { path: file_path, defer_tools: true },
556 json_source(r#"{"servers":{"coding":{"type":"stdio","command":"from_json","deferTools":false}}}"#),
557 ];
558
559 let builder = builder_from_sources(&sources);
560
561 assert_eq!(command_for(&builder, "coding"), Some("from_json"));
562 assert_eq!(deferred_tools_for(&builder, "coding"), Some(false));
563 }
564
565 #[tokio::test]
566 async fn spawn_returns_immediately_and_emits_initial_connecting_status() {
567 let spawn = McpBuilder::new("/workspace")
568 .from_mcp_config_sources(&[json_source(
569 r#"{"servers":{"slow":{"type":"stdio","command":"sleep","args":["30"]}}}"#,
570 )])
571 .unwrap()
572 .spawn()
573 .await
574 .expect("spawn should succeed");
575
576 let (_runtime, mut event_rx) = spawn.split();
577 let event = event_rx.try_recv().expect("spawn() should buffer an initial ServerStatusesChanged");
578 let McpClientEvent::ServerStatusesChanged(statuses) = event else {
579 panic!("expected ServerStatusesChanged, got {event:?}");
580 };
581 assert!(matches!(statuses[0].status, McpServerStatus::Connecting));
582 }
583
584 #[tokio::test]
585 async fn from_mcp_config_sources_expands_workspace_var_in_stdio_args() {
586 let builder = McpBuilder::new("/work")
587 .from_mcp_config_sources(&[json_source(
588 r#"{"servers":{"notes":{"type":"stdio","command":"server","args":["--dir","${WORKSPACE}/notes"]}}}"#,
589 )])
590 .unwrap();
591
592 assert_eq!(args_for(&builder, "notes"), Some(vec!["--dir".to_string(), "/work/notes".to_string()]));
593 }
594
595 #[tokio::test]
596 async fn from_mcp_config_sources_expands_aether_home_var_in_stdio_args() {
597 let home = tempfile::tempdir().unwrap();
598
599 let builder = McpBuilder::new("/work")
600 .with_aether_home(home.path())
601 .from_mcp_config_sources(&[json_source(
602 r#"{"servers":{"skills":{"type":"stdio","command":"server","args":["--dir","${AETHER_HOME}/skills"]}}}"#,
603 )])
604 .unwrap();
605
606 assert_eq!(
607 args_for(&builder, "skills"),
608 Some(vec!["--dir".to_string(), home.path().join("skills").to_string_lossy().into_owned()])
609 );
610 }
611
612 #[tokio::test]
613 async fn reserved_progressive_discovery_server_is_rejected_when_gateway_is_enabled() {
614 let result = McpBuilder::new("/workspace")
615 .from_mcp_config_sources(&[json_source(
616 r#"{"servers":{"progressive-discovery":{"type":"stdio","command":"server"},"deferred":{"type":"stdio","command":"server","deferTools":true}}}"#,
617 )])
618 .unwrap()
619 .spawn()
620 .await;
621
622 assert!(matches!(result, Err(McpError::ReservedServerName(name)) if name == "progressive-discovery"));
623 }
624
625 #[test]
626 fn new_sets_root_directory_from_workspace_root() {
627 let builder = McpBuilder::new("/workspace");
628 assert_eq!(builder.root_dir, PathBuf::from("/workspace"));
629 }
630
631 fn command_for<'a>(builder: &'a McpBuilder, name: &str) -> Option<&'a str> {
632 builder.servers.iter().find_map(|server| match &server.transport {
633 McpTransport::Stdio { command, .. } if server.name == name => Some(command.as_str()),
634 _ => None,
635 })
636 }
637
638 fn args_for(builder: &McpBuilder, name: &str) -> Option<Vec<String>> {
639 builder.servers.iter().find_map(|server| match &server.transport {
640 McpTransport::Stdio { args, .. } if server.name == name => Some(args.clone()),
641 _ => None,
642 })
643 }
644
645 fn is_direct_tool(builder: &McpBuilder, server_name: &str, tool_name: &str) -> bool {
646 builder
647 .servers
648 .iter()
649 .find(|server| server.name == server_name)
650 .is_some_and(|server| server.tool_exposure.is_model_visible_tool(tool_name))
651 }
652
653 fn deferred_tools_for(builder: &McpBuilder, name: &str) -> Option<bool> {
654 builder.servers.iter().find(|server| server.name == name).map(mcp_utils::client::McpServer::has_deferred_tools)
655 }
656}