Skip to main content

aether_cli/
runtime.rs

1use crate::error::CliError;
2use aether_core::agent_spec::{AgentSpec, McpConfigSource};
3use aether_core::core::{AgentBuilder, AgentDeps, AgentHandle, Prompt};
4use aether_core::events::{AgentEvent, Command};
5use aether_core::mcp::McpBuilder;
6use aether_core::mcp::mcp;
7use aether_core::mcp::{McpHandle, McpRuntime, McpSession};
8use llm::{ChatMessage, SessionUsageEvent, ToolDefinition};
9use mcp_servers::McpBuilderExt;
10use mcp_utils::client::{McpClientEvent, McpConnectionDetails, McpServer, OAuthHandlerFactory};
11use std::path::PathBuf;
12use tokio::sync::mpsc::{Receiver, Sender};
13use tracing::debug;
14
15pub struct RuntimeBuilder {
16    cwd: PathBuf,
17    spec: AgentSpec,
18    mcp_config_sources: Vec<McpConfigSource>,
19    extra_mcp_servers: Vec<McpServer>,
20    oauth_applicator: Option<Box<dyn FnOnce(McpBuilder) -> McpBuilder + Send>>,
21    agent_deps: AgentDeps,
22    usage_seed: Option<SessionUsageEvent>,
23}
24
25pub struct Runtime {
26    pub agent_tx: Sender<Command>,
27    pub agent_rx: Receiver<AgentEvent>,
28    pub agent_handle: AgentHandle,
29    pub event_rx: Receiver<McpClientEvent>,
30    pub mcp_runtime: McpRuntime,
31}
32
33pub struct PromptInfo {
34    pub spec: AgentSpec,
35    pub tool_definitions: Vec<ToolDefinition>,
36}
37
38impl RuntimeBuilder {
39    pub fn from_spec(cwd: PathBuf, spec: AgentSpec) -> Self {
40        Self {
41            cwd,
42            spec,
43            mcp_config_sources: Vec::new(),
44            extra_mcp_servers: Vec::new(),
45            oauth_applicator: None,
46            agent_deps: AgentDeps::default(),
47            usage_seed: None,
48        }
49    }
50
51    pub fn agent_deps(mut self, deps: AgentDeps) -> Self {
52        self.agent_deps = deps;
53        self
54    }
55
56    /// Continue session usage totals from the last persisted usage event.
57    pub fn resume_usage(mut self, last: SessionUsageEvent) -> Self {
58        self.usage_seed = Some(last);
59        self
60    }
61
62    /// Set MCP config source overrides. When non-empty, these completely
63    /// replace any sources resolved from the agent's `AgentSpec`.
64    pub fn mcp_sources(mut self, sources: Vec<McpConfigSource>) -> Self {
65        self.mcp_config_sources = sources;
66        self
67    }
68
69    pub fn extra_servers(mut self, servers: Vec<McpServer>) -> Self {
70        self.extra_mcp_servers = servers;
71        self
72    }
73
74    pub fn oauth_handler_factory(mut self, factory: OAuthHandlerFactory) -> Self {
75        self.oauth_applicator = Some(Box::new(|builder| builder.with_oauth_handler_factory(factory)));
76        self
77    }
78
79    pub async fn build(
80        self,
81        custom_prompt: Option<Prompt>,
82        messages: Option<Vec<ChatMessage>>,
83    ) -> Result<Runtime, CliError> {
84        let deps = self.agent_deps.clone();
85        let usage_seed = self.usage_seed.clone();
86        let (spec, session) = self.spawn_mcp().await?;
87        let mcp = session.handle().clone();
88
89        let (agent_tx, agent_rx, agent_handle) = spawn_agent(&spec, &deps, mcp, Vec::new(), |mut agent_builder| {
90            if let Some(last) = &usage_seed {
91                agent_builder = agent_builder.resume_usage(last);
92            }
93            if let Some(prompt) = custom_prompt {
94                agent_builder = agent_builder.system_prompt(prompt);
95            }
96            if let Some(msgs) = messages {
97                agent_builder = agent_builder.messages(msgs);
98            }
99            agent_builder
100        })
101        .await?;
102        let (mcp_runtime, event_rx) = session.connect_agent(agent_tx.clone()).await.split();
103
104        Ok(Runtime { agent_tx, agent_rx, agent_handle, event_rx, mcp_runtime })
105    }
106
107    /// Spawn MCP, block until every server finishes its initial connection,
108    /// then connect the agent to the session's filtered tools and MCP
109    /// instructions before returning. Returns the live [`Runtime`] plus the
110    /// bootstrap snapshot for callers that need the agent ready to use tools on
111    /// its first turn.
112    pub async fn build_ready(self, messages: Vec<ChatMessage>) -> Result<(Runtime, McpConnectionDetails), CliError> {
113        let deps = self.agent_deps.clone();
114        let usage_seed = self.usage_seed.clone();
115        let (spec, mut session) = self.spawn_mcp().await?;
116        let snapshot = session
117            .block_until_ready()
118            .await
119            .ok_or_else(|| CliError::McpError("MCP bootstrap aborted before completion".to_string()))?;
120        let mcp = session.handle().clone();
121
122        let (agent_tx, agent_rx, agent_handle) = spawn_agent(&spec, &deps, mcp, Vec::new(), |mut agent_builder| {
123            if let Some(last) = &usage_seed {
124                agent_builder = agent_builder.resume_usage(last);
125            }
126            agent_builder.messages(messages)
127        })
128        .await?;
129        let (mcp_runtime, event_rx) = session.connect_agent(agent_tx.clone()).await.split();
130
131        Ok((Runtime { agent_tx, agent_rx, agent_handle, event_rx, mcp_runtime }, snapshot))
132    }
133
134    pub async fn build_prompt_info(self) -> Result<PromptInfo, CliError> {
135        let (spec, mut session) = self.spawn_mcp().await?;
136        let details = session
137            .block_until_ready()
138            .await
139            .ok_or_else(|| CliError::McpError("MCP bootstrap aborted before completion".to_string()))?;
140        let filtered_tools = details.tool_definitions();
141        Ok(PromptInfo { spec, tool_definitions: filtered_tools })
142    }
143
144    async fn spawn_mcp(self) -> Result<(AgentSpec, McpSession), CliError> {
145        let deps = self.agent_deps.clone();
146        let mut builder = mcp(&self.cwd).with_tool_filter(self.spec.tools.clone());
147
148        if let Some(apply_oauth) = self.oauth_applicator {
149            builder = apply_oauth(builder);
150        }
151
152        builder = builder.with_agent_deps(deps).with_builtin_servers();
153
154        if !self.extra_mcp_servers.is_empty() {
155            builder = builder.with_servers(self.extra_mcp_servers);
156        }
157
158        let mcp_config_sources: Vec<McpConfigSource> = if self.mcp_config_sources.is_empty() {
159            self.spec.mcp_config_sources.clone()
160        } else {
161            self.mcp_config_sources
162        };
163
164        if !mcp_config_sources.is_empty() {
165            debug!("Loading MCP configs from: {:?}", mcp_config_sources);
166            builder =
167                builder.from_mcp_config_sources(&mcp_config_sources).map_err(|e| CliError::McpError(e.to_string()))?;
168        }
169
170        let spawn = builder.spawn().await.map_err(|e| CliError::McpError(e.to_string()))?;
171        Ok((self.spec, spawn))
172    }
173}
174
175async fn spawn_agent(
176    spec: &AgentSpec,
177    deps: &AgentDeps,
178    mcp: McpHandle,
179    tool_definitions: Vec<ToolDefinition>,
180    configure: impl FnOnce(AgentBuilder) -> AgentBuilder,
181) -> Result<(Sender<Command>, Receiver<AgentEvent>, AgentHandle), CliError> {
182    let builder = AgentBuilder::from_spec(spec, vec![], deps)
183        .await
184        .map_err(|error| CliError::AgentError(error.to_string()))?
185        .tools(mcp, tool_definitions);
186
187    configure(builder).spawn().await.map_err(|error| CliError::AgentError(error.to_string()))
188}