Skip to main content

aether_core/mcp/
mcp_builder.rs

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
47/// Owns a spawned MCP manager. Dropping this value aborts the manager task.
48pub 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
74/// A freshly spawned MCP manager paired with its event stream. Consumers
75/// receive incremental updates over the event stream (starting with an initial
76/// `ServerStatusesChanged` reflecting every configured server in `Connecting`)
77/// and can call [`split`](Self::split) to separate the stream from the
78/// [`McpRuntime`] that keeps the manager task alive.
79pub 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    /// Synchronize this session's current and future tools and instructions with
94    /// one agent. Initial state is sent before this method returns.
95    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    /// Block until the manager finishes bootstrapping every initially-configured
153    /// server, then return the consolidated snapshot. Returns `None` if the
154    /// event channel closes before `ConnectionReady` is received.
155    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    /// Cross-cutting dependencies handed to every agent spawned behind this
229    /// builder's in-memory servers.
230    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}