vtcode_core/tools/registry/
mcp_facade.rs1use rustc_hash::{FxHashMap, FxHashSet};
4use std::sync::Arc;
5use std::time::Duration;
6
7use anyhow::{Context, Result, anyhow};
8use serde_json::Value;
9use tracing::{debug, warn};
10
11use crate::mcp::{McpClient, McpToolExecutor, McpToolInfo};
12use crate::tools::mcp::build_mcp_registration;
13use vtcode_commons::classify_anyhow_error;
14
15use super::ToolRegistry;
16use super::mcp_helpers::normalize_mcp_tool_identifier;
17use super::registration::ToolCatalogSource;
18
19fn sorted_providers(index: &FxHashMap<String, Vec<String>>) -> Vec<(&String, &Vec<String>)> {
21 let mut providers: Vec<_> = index.iter().collect();
22 providers.sort_unstable_by(|a, b| a.0.cmp(b.0));
23 providers
24}
25
26fn mcp_refresh_retry_allowed(error: &anyhow::Error) -> bool {
27 vtcode_commons::detect_misconfiguration_in_anyhow(error).is_none()
28}
29
30impl ToolRegistry {
31 fn remove_all_mcp_proxy_tools(&self) {
39 let stale: Vec<String> = self
40 .inventory
41 .registrations_snapshot()
42 .into_iter()
43 .filter(|registration| registration.catalog_source() == ToolCatalogSource::Mcp)
44 .map(|registration| registration.name().to_string())
45 .collect();
46 for name in stale {
47 if let Err(err) = self.inventory.remove_tool(&name) {
48 warn!(tool = %name, error = %err, "failed to remove stale MCP proxy tool");
49 }
50 }
51 }
52
53 pub async fn with_mcp_client(self, mcp_client: Arc<McpClient>) -> Self {
55 self.remove_all_mcp_proxy_tools();
56 *self.mcp_client.write() = Some(mcp_client);
57 self.mcp_tool_index.write().await.clear();
58 self.mcp_reverse_index.write().await.clear();
59 *self.cached_available_tools.write() = None;
60 self.initialized.store(false, std::sync::atomic::Ordering::Relaxed);
61 self
62 }
63
64 pub async fn set_mcp_client(&self, mcp_client: Arc<McpClient>) {
66 self.remove_all_mcp_proxy_tools();
67 *self.mcp_client.write() = Some(mcp_client);
68 self.mcp_tool_index.write().await.clear();
69 self.mcp_reverse_index.write().await.clear();
70 *self.cached_available_tools.write() = None;
71 self.initialized.store(false, std::sync::atomic::Ordering::Relaxed);
72 }
73
74 pub async fn clear_mcp_client(&self) {
82 self.remove_all_mcp_proxy_tools();
83 *self.mcp_client.write() = None;
84 self.mcp_tool_index.write().await.clear();
85 self.mcp_reverse_index.write().await.clear();
86 *self.cached_available_tools.write() = None;
87 self.initialized.store(false, std::sync::atomic::Ordering::Relaxed);
88 self.rebuild_tool_assembly().await;
89 self.tool_catalog_state.note_explicit_refresh("mcp_client_cleared");
90 self.sync_policy_catalog().await;
91 }
92
93 pub fn mcp_client(&self) -> Option<Arc<McpClient>> {
95 self.mcp_client.read().clone()
96 }
97
98 pub async fn list_mcp_tools(&self) -> Result<Vec<McpToolInfo>> {
100 let index = self.mcp_tool_index.read().await;
101 if index.is_empty() {
102 return Ok(Vec::new());
103 }
104
105 let providers = sorted_providers(&index);
106
107 let mut mcp_tools = Vec::new();
108 for (provider, tools) in providers {
109 for tool_name in tools {
110 let canonical_name = format!("mcp::{provider}::{tool_name}");
111 if let Some(registration) = self.inventory.get_registration(&canonical_name) {
112 mcp_tools.push(McpToolInfo {
113 name: tool_name.clone(),
114 description: registration.metadata().description().unwrap_or("").to_string(),
115 provider: provider.clone(),
116 input_schema: registration.parameter_schema().cloned().unwrap_or(Value::Null),
117 output_schema: None,
121 });
122 }
123 }
124 }
125
126 Ok(mcp_tools)
127 }
128
129 pub async fn has_mcp_tool(&self, tool_name: &str) -> bool {
131 self.mcp_reverse_index.read().await.contains_key(tool_name)
132 }
133
134 pub async fn execute_mcp_tool(&self, tool_name: &str, args: Value) -> Result<Value> {
136 let client_opt = self.mcp_client.read().clone();
137 if let Some(mcp_client) = client_opt {
138 mcp_client.execute_mcp_tool(tool_name, &args).await
139 } else {
140 Err(anyhow!(
141 "MCP client not available (no active MCP connections). The requested MCP tool '{tool_name}' cannot run while disconnected. Use `/mcp repair` or `mcp connect <server>` to reconnect, then retry."
142 ))
143 }
144 }
145
146 pub(super) async fn resolve_mcp_tool_alias(&self, tool_name: &str) -> Option<String> {
147 let normalized = normalize_mcp_tool_identifier(tool_name);
148 if normalized.is_empty() {
149 return None;
150 }
151
152 let index = self.mcp_tool_index.read().await;
153 for tools in index.values() {
154 for tool in tools {
155 if normalize_mcp_tool_identifier(tool) == normalized {
156 return Some(tool.clone());
157 }
158 }
159 }
160
161 None
162 }
163
164 pub async fn refresh_mcp_tools(&self) -> Result<()> {
166 let mcp_client_opt = self.mcp_client.read().clone();
167 if let Some(mcp_client) = mcp_client_opt {
168 debug!("Refreshing MCP tools for {} providers", mcp_client.get_status().provider_count);
169
170 let mut tools: Option<Vec<McpToolInfo>> = None;
171 let mut last_err: Option<anyhow::Error> = None;
172 for attempt in 0..3 {
173 match mcp_client.list_mcp_tools().await {
174 Ok(list) => {
175 tools = Some(list);
176 break;
177 }
178 Err(err) => {
179 let retry_allowed = mcp_refresh_retry_allowed(&err);
180 last_err = Some(err);
181 if !retry_allowed {
182 warn!(
183 attempt = attempt + 1,
184 "MCP tool refresh failed due to configuration; skipping retries"
185 );
186 break;
187 }
188 let jitter = (attempt * 37) % 80;
189 let pow = 2_u64.saturating_pow(attempt.min(4) as u32); let backoff = Duration::from_millis(200 * pow + jitter).min(Duration::from_secs(3));
191 warn!(
192 attempt = attempt + 1,
193 delay_ms = %backoff.as_millis(),
194 "Failed to list MCP tools, retrying with backoff"
195 );
196 tokio::time::sleep(backoff).await;
197 }
198 }
199 }
200
201 let tools = match tools {
202 Some(list) => list,
203 None => {
204 let Some(error) = last_err else {
205 warn!("Failed to refresh MCP tools without an error payload; keeping existing cache");
206 return Ok(());
207 };
208 if let Some(guidance) = vtcode_commons::detect_misconfiguration_in_anyhow(&error) {
209 return Err(anyhow!("{error:#}: {}", guidance.user_message()))
210 .context("MCP tool refresh is blocked by configuration");
211 }
212 warn!(
213 error = %error,
214 "Failed to refresh MCP tools after retries; keeping existing cache"
215 );
216 let category = classify_anyhow_error(&error);
217 self.mcp_circuit_breaker.record_failure_category(category);
218 return Ok(());
219 }
220 };
221 let existing_tools: Vec<String> = {
222 let index = self.mcp_tool_index.read().await;
223 index
224 .iter()
225 .flat_map(|(provider, names)| names.iter().map(move |name| format!("mcp::{provider}::{name}")))
226 .collect()
227 };
228 for canonical_name in existing_tools {
229 if let Err(err) = self.inventory.remove_tool(&canonical_name) {
230 warn!(
231 tool = %canonical_name,
232 error = %err,
233 "failed to remove stale MCP proxy tool"
234 );
235 }
236 }
237
238 let mut provider_map: FxHashMap<String, Vec<String>> = FxHashMap::default();
239 let mut seen_tools = FxHashSet::default();
240
241 for tool in &tools {
242 let canonical_name = format!("mcp::{}::{}", tool.provider, tool.name);
243 if !seen_tools.insert(canonical_name) {
244 continue;
245 }
246 let registration = match build_mcp_registration(Arc::clone(&mcp_client), &tool.provider, tool, None) {
247 Ok(registration) => registration,
248 Err(error) => {
249 warn!(%error, "Rejected MCP proxy metadata");
250 continue;
251 }
252 };
253 if let Err(error) = self.inventory.register_tool(registration) {
254 warn!(%error, "Failed to register MCP proxy tool");
255 continue;
256 }
257 provider_map.entry(tool.provider.clone()).or_default().push(tool.name.clone());
258 }
259
260 for tools in provider_map.values_mut() {
261 tools.sort();
262 tools.dedup();
263 }
264
265 *self.mcp_tool_index.write().await = provider_map;
266 {
267 let mut reverse_index = self.mcp_reverse_index.write().await;
268 reverse_index.clear();
269 let index = self.mcp_tool_index.read().await;
270 for (provider, tools) in sorted_providers(&index) {
272 for tool in tools {
273 let _ = reverse_index.entry(tool.clone()).or_insert_with(|| provider.clone());
274 }
275 }
276 }
277
278 let mcp_index = self.mcp_tool_index.read().await;
279 let std_index: hashbrown::HashMap<String, Vec<String>> =
280 mcp_index.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
281 let allowlist = {
282 let gateway = self.policy_gateway.clone();
283 gateway.update_mcp_tools(&std_index).await?
284 };
285 if let Some(allowlist) = allowlist {
286 mcp_client.update_allowlist(allowlist);
287 }
288
289 *self.cached_available_tools.write() = None;
290 self.rebuild_tool_assembly().await;
291 self.tool_catalog_state.note_explicit_refresh("mcp_tool_refresh");
292 self.sync_policy_catalog().await;
293 self.mcp_circuit_breaker.record_success();
295 Ok(())
296 } else {
297 debug!("No MCP client configured, pruning stale MCP proxy tools");
298 self.remove_all_mcp_proxy_tools();
299 self.mcp_tool_index.write().await.clear();
300 self.mcp_reverse_index.write().await.clear();
301 *self.cached_available_tools.write() = None;
302 self.rebuild_tool_assembly().await;
303 self.tool_catalog_state.note_explicit_refresh("mcp_tool_refresh_no_client");
304 self.sync_policy_catalog().await;
305 Ok(())
306 }
307 }
308}
309
310#[cfg(test)]
311mod tests {
312 use super::{mcp_refresh_retry_allowed, sorted_providers};
313 use anyhow::anyhow;
314 use rustc_hash::FxHashMap;
315
316 #[test]
317 fn mcp_refresh_skips_configuration_failures_but_retries_transient_errors() {
318 assert!(!mcp_refresh_retry_allowed(&anyhow!("MCP server URL invalid: endpoint must use https")));
319 assert!(mcp_refresh_retry_allowed(&anyhow!("MCP server connection reset by peer")));
320 }
321
322 #[test]
323 fn sorted_providers_orders_by_name_independent_of_insertion() {
324 let mut index: FxHashMap<String, Vec<String>> = FxHashMap::default();
325 for name in ["zeta", "Alpha", "alpha", "beta"] {
326 drop(index.insert(name.to_owned(), vec![format!("{name}_tool")]));
327 }
328
329 let names: Vec<&str> = sorted_providers(&index).iter().map(|(name, _)| name.as_str()).collect();
330 assert_eq!(names, ["Alpha", "alpha", "beta", "zeta"]);
331 }
332}