1use crate::events::{TaskOutcomeState, TraceContext, task_created_result};
2use crate::mcp::tool_bridge::{convert_tool_result, map_task_result_to_outcome};
3use crate::mcp::{McpHandle, McpRuntime, ServerFactory, ToolCallStream, mcp};
4use futures::{FutureExt, StreamExt};
5use mcp_utils::client::{
6 CallToolOptions, CancellationToken, InMemoryServerSpec, McpConnectionDetails, McpServer, McpTransport,
7 ToolCallEvent, ToolExposure,
8};
9use mcp_utils::testing::ElicitationScript;
10use rmcp::model::{CreateTaskResult, ElicitResult, ProgressNotificationParam};
11use rmcp::{RoleServer, ServerHandler, service::DynService};
12use serde_json::Value;
13use std::collections::{HashMap, VecDeque};
14use std::sync::Mutex;
15use std::sync::atomic::{AtomicU64, Ordering};
16use std::time::Duration;
17use tokio::sync::watch;
18
19pub use mcp_utils::testing::CapturedElicitation;
20
21const DEFAULT_TOOL_TIMEOUT: Duration = Duration::from_secs(10);
22
23#[derive(Default)]
24pub struct McpTestBuilder {
25 servers: Vec<McpServer>,
26 factories: Vec<(String, ServerFactory)>,
27 elicitation_responses: Vec<ElicitResult>,
28 trace_context: Option<TraceContext>,
29 tool_timeout: Duration,
30}
31
32fn task_outcome(outcome: crate::events::TaskOutcome) -> TaskOutcome {
33 let (status, body) = match outcome.state {
34 TaskOutcomeState::Completed { result, .. } => ("completed", result.result),
35 TaskOutcomeState::Failed { error } => ("failed", error.error),
36 TaskOutcomeState::Cancelled => {
37 ("cancelled", "The background task was cancelled and will not produce a result.".into())
38 }
39 };
40 TaskOutcome { task_id: outcome.task_id, status: status.into(), body }
41}
42
43pub struct McpTest {
44 mcp: McpHandle,
45 _runtime: McpRuntime,
46 snapshot: McpConnectionDetails,
47 elicitations: ElicitationScript,
48 deferred_tools: tokio::sync::Mutex<VecDeque<DeferredTool>>,
49 cancel_tokens: Mutex<HashMap<String, CancellationToken>>,
50 trace_context: Option<TraceContext>,
51 tool_timeout: Duration,
52 next_call_id: AtomicU64,
53}
54
55pub struct TaskOutcome {
56 pub task_id: String,
57 pub status: String,
58 pub body: String,
59}
60
61pub struct ToolCallOutcome {
62 pub result: Result<llm::ToolCallResult, llm::ToolCallError>,
63 pub progress: Vec<ProgressNotificationParam>,
64 pub deferred_task: Option<CreateTaskResult>,
65}
66
67struct DeferredTool {
68 request: llm::ToolCallRequest,
69 events: ToolCallStream,
70}
71
72impl McpTestBuilder {
73 pub fn new() -> Self {
74 Self::default()
75 }
76
77 pub fn server<S>(self, name: impl Into<String>, server: S) -> Self
78 where
79 S: ServerHandler + Clone + Send + Sync + 'static,
80 {
81 self.server_with_exposure(name, server, ToolExposure::ModelVisible)
82 }
83
84 pub fn deferred_server<T>(self, name: impl Into<String>, server: T) -> Self
85 where
86 T: ServerHandler + Clone + Send + Sync + 'static,
87 {
88 self.server_with_exposure(name, server, ToolExposure::deferred_all())
89 }
90
91 pub fn server_with_exposure<T>(mut self, name: impl Into<String>, server: T, exposure: ToolExposure) -> Self
92 where
93 T: ServerHandler + Clone + Send + Sync + 'static,
94 {
95 let name = name.into();
96 let factory_name = format!("test-{}", self.factories.len());
97 let factory_server = server;
98 let factory: ServerFactory = Box::new(move |_spec, _services| {
99 let server = factory_server.clone();
100 async move { Box::new(server) as Box<dyn DynService<RoleServer>> }.boxed()
101 });
102 self.factories.push((factory_name.clone(), factory));
103 self.servers.push(McpServer::new(
104 name,
105 McpTransport::InMemory {
106 spec: InMemoryServerSpec { factory: factory_name, args: Vec::new(), input: None },
107 },
108 exposure,
109 ));
110 self
111 }
112
113 pub fn elicitation_response(mut self, response: ElicitResult) -> Self {
114 self.elicitation_responses.push(response);
115 self
116 }
117
118 pub fn trace_context(mut self, trace_context: TraceContext) -> Self {
119 self.trace_context = Some(trace_context);
120 self
121 }
122
123 pub fn tool_timeout(mut self, timeout: Duration) -> Self {
124 self.tool_timeout = timeout;
125 self
126 }
127
128 pub async fn build(self) -> McpTest {
129 let builder = self
130 .factories
131 .into_iter()
132 .fold(mcp("/workspace").with_servers(self.servers), |builder, (name, factory)| {
133 builder.register_in_memory_server(name, factory)
134 });
135 let mut spawn = builder.spawn().await.expect("MCP test manager spawns");
136 let snapshot = spawn.block_until_ready().await.expect("MCP test manager becomes ready");
137 let (runtime, event_rx) = spawn.split();
138
139 McpTest {
140 mcp: runtime.handle().clone(),
141 _runtime: runtime,
142 snapshot,
143 elicitations: ElicitationScript::spawn(event_rx, self.elicitation_responses),
144 deferred_tools: tokio::sync::Mutex::new(VecDeque::new()),
145 cancel_tokens: Mutex::new(HashMap::new()),
146 trace_context: self.trace_context,
147 tool_timeout: if self.tool_timeout.is_zero() { DEFAULT_TOOL_TIMEOUT } else { self.tool_timeout },
148 next_call_id: AtomicU64::new(1),
149 }
150 }
151}
152
153impl McpTest {
154 pub async fn call(&self, server: &str, tool: &str, arguments: Value) -> ToolCallOutcome {
155 let id = self.next_call_id.fetch_add(1, Ordering::Relaxed);
156 let request = llm::ToolCallRequest {
157 id: format!("mcp-test-{id}"),
158 name: format!("{server}__{tool}"),
159 arguments: arguments.to_string(),
160 };
161 let request_for_outcome = request.clone();
162 let cancel = CancellationToken::new();
163 self.cancel_tokens.lock().expect("cancel token lock").insert(request.id.clone(), cancel.clone());
164 let options = CallToolOptions {
165 timeout: self.tool_timeout,
166 meta: self.trace_context.as_ref().map(TraceContext::to_meta),
167 cancel,
168 };
169 let mut events = self.mcp.call_model_visible(request.name, &request.arguments, options);
170
171 let mut progress = Vec::new();
172 while let Some(event) = events.next().await {
173 match event {
174 ToolCallEvent::Progress(event) => progress.push(event),
175 ToolCallEvent::TaskCreated(task) => {
176 self.deferred_tools
177 .lock()
178 .await
179 .push_back(DeferredTool { request: request_for_outcome.clone(), events });
180 return ToolCallOutcome {
181 result: Ok(task_created_result(&request_for_outcome, &task.task.task_id)),
182 progress,
183 deferred_task: Some(task),
184 };
185 }
186 ToolCallEvent::Complete(outcome) => {
187 let result = convert_tool_result(&request_for_outcome, outcome).map(|(result, _)| result);
188 return ToolCallOutcome { result, progress, deferred_task: None };
189 }
190 ToolCallEvent::TaskStatus(_) | ToolCallEvent::TaskComplete { .. } | ToolCallEvent::Cancelled { .. } => {
191 panic!("MCP task lifecycle event arrived before deferral")
192 }
193 }
194 }
195 panic!("MCP test tool event stream ended before completion");
196 }
197
198 pub fn cancel_tool(&self, tool_id: &str) {
199 let tokens = self.cancel_tokens.lock().expect("cancel token lock");
200 tokens.get(tool_id).expect("cancel_tool targets a tool started with call()").cancel();
201 }
202
203 pub async fn next_tool_event(&self) -> Option<ToolCallEvent> {
204 self.next_deferred_event().await.map(|(_, event)| event)
205 }
206
207 pub async fn next_task_outcome(&self) -> Option<TaskOutcome> {
208 while let Some((request, event)) = self.next_deferred_event().await {
209 match event {
210 ToolCallEvent::TaskComplete { task, result } => {
211 return Some(task_outcome(map_task_result_to_outcome(request, task, result)));
212 }
213 ToolCallEvent::Cancelled { task_id } => {
214 return Some(task_outcome(crate::events::TaskOutcome {
215 request,
216 task_id: task_id.unwrap_or_else(|| "pending".to_string()),
217 state: TaskOutcomeState::Cancelled,
218 }));
219 }
220 ToolCallEvent::Progress(_)
221 | ToolCallEvent::TaskCreated(_)
222 | ToolCallEvent::TaskStatus(_)
223 | ToolCallEvent::Complete(_) => {}
224 }
225 }
226 None
227 }
228
229 async fn next_deferred_event(&self) -> Option<(llm::ToolCallRequest, ToolCallEvent)> {
230 let mut deferred = self.deferred_tools.lock().await;
231 loop {
232 let front = deferred.front_mut()?;
233 match front.events.next().await {
234 Some(event) => return Some((front.request.clone(), event)),
235 None => {
236 deferred.pop_front();
237 }
238 }
239 }
240 }
241
242 pub fn snapshot(&self) -> &McpConnectionDetails {
243 &self.snapshot
244 }
245
246 pub fn subscribe(&self) -> watch::Receiver<McpConnectionDetails> {
247 self.mcp.subscribe()
248 }
249
250 pub fn elicitations(&self) -> Vec<CapturedElicitation> {
251 self.elicitations.captured()
252 }
253}