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 pub fn resume_usage(mut self, last: SessionUsageEvent) -> Self {
58 self.usage_seed = Some(last);
59 self
60 }
61
62 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 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}