feat(mcp): extend the server catalog to resources, templates, and prompts

Implements the unified catalog from plans/mcp-resources-prompts-design.md §4.2 (T2): CatalogItem gains kind/uri/mime_type/size keyed as {kind}:{id}; catalog_items() lists per kind gated by advertised capabilities with warn-and-degrade; mcp_search results carry kind; mcp_describe gains an optional kind param (default tool); write-only registry ServerCatalog removed.
This commit is contained in:
2026-08-25 11:37:41 -06:00
parent 01ada1da18
commit d68f4ecaeb
5 changed files with 701 additions and 76 deletions
+1
View File
@@ -141,6 +141,7 @@ arboard = { version = "3.3.0", default-features = false }
[dev-dependencies]
pretty_assertions = "1.4.0"
rmcp = { version = "3.1.2", features = ["server"] }
serial_test = "3"
[[bin]]
+2
View File
@@ -51,6 +51,8 @@ pub use self::skill::Skill;
pub use self::skill_policy::SkillPolicy;
#[allow(unused_imports)]
pub use self::skill_registry::SkillRegistry;
#[cfg(test)]
pub(crate) use self::tool_scope::test_fixtures;
pub use self::update::run_self_update;
use crate::client::{
self, ClientConfig, MessageContentToolCalls, Model, ModelType, OPENAI_COMPATIBLE_PROVIDERS,
+583 -15
View File
@@ -1,9 +1,11 @@
use crate::function::{Functions, ToolCallTracker};
use crate::mcp::{CatalogItem, ConnectedServer, McpRegistry};
use crate::mcp::{CatalogItem, CatalogItemKind, ConnectedServer, McpRegistry};
use anyhow::{Context, Result, anyhow};
use bm25::{Document, Language, SearchEngineBuilder};
use rmcp::model::{CallToolRequestParams, CallToolResult};
use rmcp::model::{
CallToolRequestParams, CallToolResult, Prompt, Resource, ResourceTemplate, Tool,
};
use serde_json::{Value, json};
use std::collections::HashMap;
use std::sync::Arc;
@@ -63,16 +65,56 @@ impl McpRuntime {
.get(server)
.cloned()
.with_context(|| format!("{server} MCP server not found in runtime"))?;
let tools = server_handle.list_all_tools().await?;
let capabilities = server_handle
.peer_info()
.map(|info| info.capabilities.clone());
let mut items = HashMap::new();
for tool in tools {
let item = CatalogItem {
name: tool.name.to_string(),
server: server.to_string(),
description: tool.description.unwrap_or_default().to_string(),
};
items.insert(item.name.clone(), item);
if capabilities.as_ref().is_none_or(|c| c.tools.is_some()) {
match server_handle.list_all_tools().await {
Ok(tools) => merge_catalog_items(
&mut items,
tools
.into_iter()
.map(|tool| tool_catalog_item(server, tool)),
),
Err(e) => warn!("Failed to list tools on MCP server {server}: {e}"),
}
}
if capabilities.as_ref().is_some_and(|c| c.resources.is_some()) {
match server_handle.list_all_resources().await {
Ok(resources) => merge_catalog_items(
&mut items,
resources
.into_iter()
.map(|resource| resource_catalog_item(server, resource)),
),
Err(e) => warn!("Failed to list resources on MCP server {server}: {e}"),
}
match server_handle.list_all_resource_templates().await {
Ok(templates) => merge_catalog_items(
&mut items,
templates
.into_iter()
.map(|template| resource_template_catalog_item(server, template)),
),
Err(e) => {
warn!("Failed to list resource templates on MCP server {server}: {e}")
}
}
}
if capabilities.as_ref().is_some_and(|c| c.prompts.is_some()) {
match server_handle.list_all_prompts().await {
Ok(prompts) => merge_catalog_items(
&mut items,
prompts
.into_iter()
.map(|prompt| prompt_catalog_item(server, prompt)),
),
Err(e) => warn!("Failed to list prompts on MCP server {server}: {e}"),
}
}
Ok(items)
@@ -85,11 +127,15 @@ impl McpRuntime {
top_k: usize,
) -> Result<Vec<CatalogItem>> {
let items = self.catalog_items(server).await?;
let docs = items.values().map(|item| Document {
id: item.name.clone(),
let docs = items.iter().map(|(key, item)| Document {
id: key.clone(),
contents: format!(
"{}\n{}\nserver:{}",
item.name, item.description, item.server
"{}\n{}\n{}\n{}\nserver:{}",
item.name,
item.description,
item.kind,
item.uri.as_deref().unwrap_or_default(),
item.server
),
});
let engine = SearchEngineBuilder::<String>::with_documents(Language::English, docs).build();
@@ -103,12 +149,14 @@ impl McpRuntime {
.collect())
}
pub async fn describe(&self, server: &str, tool: &str) -> Result<Value> {
pub async fn describe(&self, server: &str, kind: &str, tool: &str) -> Result<Value> {
let server_handle = self
.get(server)
.cloned()
.with_context(|| format!("{server} MCP server not found in runtime"))?;
match kind {
"tool" => {
let tool_schema = server_handle
.list_all_tools()
.await?
@@ -127,6 +175,78 @@ impl McpRuntime {
}
}))
}
"resource" => {
let resource = server_handle
.list_all_resources()
.await?
.into_iter()
.find(|item| item.uri == tool)
.ok_or_else(|| {
anyhow!("{tool} not found in {server} MCP server resource catalog")
})?;
Ok(json!({
"uri": resource.uri,
"name": resource.name,
"title": resource.title,
"description": resource.description,
"mime_type": resource.mime_type,
"size": resource.size,
}))
}
"resource_template" => {
let template = server_handle
.list_all_resource_templates()
.await?
.into_iter()
.find(|item| item.uri_template == tool)
.ok_or_else(|| {
anyhow!("{tool} not found in {server} MCP server resource template catalog")
})?;
Ok(json!({
"uri_template": template.uri_template,
"name": template.name,
"title": template.title,
"description": template.description,
"mime_type": template.mime_type,
"variables": uri_template_variables(&template.uri_template),
}))
}
"prompt" => {
let prompt = server_handle
.list_all_prompts()
.await?
.into_iter()
.find(|item| item.name == tool)
.ok_or_else(|| {
anyhow!("{tool} not found in {server} MCP server prompt catalog")
})?;
let arguments: Vec<Value> = prompt
.arguments
.unwrap_or_default()
.into_iter()
.map(|arg| {
json!({
"name": arg.name,
"description": arg.description,
"required": arg.required,
})
})
.collect();
Ok(json!({
"name": prompt.name,
"description": prompt.description,
"arguments": arguments,
}))
}
other => Err(anyhow!(
"Unknown kind '{other}'. Valid kinds: tool, resource, resource_template, prompt"
)),
}
}
pub async fn invoke(
&self,
@@ -146,10 +266,228 @@ impl McpRuntime {
}
}
fn catalog_key(item: &CatalogItem) -> String {
let id = item.uri.as_deref().unwrap_or(&item.name);
format!("{}:{id}", item.kind)
}
fn merge_catalog_items(
items: &mut HashMap<String, CatalogItem>,
new_items: impl IntoIterator<Item = CatalogItem>,
) {
for item in new_items {
items.insert(catalog_key(&item), item);
}
}
fn tool_catalog_item(server: &str, tool: Tool) -> CatalogItem {
CatalogItem {
kind: CatalogItemKind::Tool,
name: tool.name.to_string(),
server: server.to_string(),
description: tool.description.unwrap_or_default().to_string(),
..Default::default()
}
}
fn resource_catalog_item(server: &str, resource: Resource) -> CatalogItem {
CatalogItem {
kind: CatalogItemKind::Resource,
name: resource.name,
server: server.to_string(),
description: resource.description.unwrap_or_default(),
uri: Some(resource.uri),
mime_type: resource.mime_type,
size: resource.size,
}
}
fn resource_template_catalog_item(server: &str, template: ResourceTemplate) -> CatalogItem {
CatalogItem {
kind: CatalogItemKind::ResourceTemplate,
name: template.name,
server: server.to_string(),
description: template.description.unwrap_or_default(),
uri: Some(template.uri_template),
mime_type: template.mime_type,
size: None,
}
}
fn prompt_catalog_item(server: &str, prompt: Prompt) -> CatalogItem {
CatalogItem {
kind: CatalogItemKind::Prompt,
name: prompt.name,
server: server.to_string(),
description: prompt.description.unwrap_or_default(),
..Default::default()
}
}
fn uri_template_variables(template: &str) -> Vec<String> {
let mut variables = Vec::new();
let mut rest = template;
while let Some(start) = rest.find('{') {
let Some(len) = rest[start + 1..].find('}') else {
break;
};
variables.push(rest[start + 1..start + 1 + len].to_string());
rest = &rest[start + len + 2..];
}
variables
}
#[cfg(test)]
pub(crate) mod test_fixtures {
use super::*;
use rmcp::model::{
ErrorData, ListPromptsResult, ListResourceTemplatesResult, ListResourcesResult,
ListToolsResult, PaginatedRequestParams, PromptArgument, PromptsCapability,
ResourcesCapability, ServerCapabilities, ServerInfo,
};
use rmcp::service::{RequestContext, RunningService};
use rmcp::{RoleServer, ServerHandler, ServiceExt};
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Clone, Default)]
pub(crate) struct FixtureServer {
pub(crate) resources_capability: bool,
pub(crate) prompts_capability: bool,
pub(crate) fail_resource_listings: bool,
pub(crate) list_resources_calls: Arc<AtomicUsize>,
pub(crate) list_prompts_calls: Arc<AtomicUsize>,
}
impl ServerHandler for FixtureServer {
fn get_info(&self) -> ServerInfo {
let mut capabilities = ServerCapabilities::builder().enable_tools().build();
capabilities.resources = self.resources_capability.then(ResourcesCapability::default);
capabilities.prompts = self.prompts_capability.then(PromptsCapability::default);
ServerInfo::new(capabilities)
}
async fn list_tools(
&self,
_request: Option<PaginatedRequestParams>,
_context: RequestContext<RoleServer>,
) -> Result<ListToolsResult, ErrorData> {
let schema = json!({
"type": "object",
"properties": { "q": { "type": "string" } }
})
.as_object()
.cloned()
.unwrap();
Ok(ListToolsResult::with_all_items(vec![Tool::new(
"dup",
"Duplicate-named tool",
schema,
)]))
}
async fn list_resources(
&self,
_request: Option<PaginatedRequestParams>,
_context: RequestContext<RoleServer>,
) -> Result<ListResourcesResult, ErrorData> {
self.list_resources_calls.fetch_add(1, Ordering::SeqCst);
if self.fail_resource_listings {
return Err(ErrorData::internal_error("resource listing exploded", None));
}
Ok(ListResourcesResult::with_all_items(vec![
Resource::new("dup", "dup-resource")
.with_description("Duplicate-named resource")
.with_mime_type("text/plain")
.with_size(42),
]))
}
async fn list_resource_templates(
&self,
_request: Option<PaginatedRequestParams>,
_context: RequestContext<RoleServer>,
) -> Result<ListResourceTemplatesResult, ErrorData> {
if self.fail_resource_listings {
return Err(ErrorData::internal_error("template listing exploded", None));
}
Ok(ListResourceTemplatesResult::with_all_items(vec![
ResourceTemplate::new("file:///{path}/{name}", "file-template")
.with_description("Read a file")
.with_mime_type("text/plain"),
]))
}
async fn list_prompts(
&self,
_request: Option<PaginatedRequestParams>,
_context: RequestContext<RoleServer>,
) -> Result<ListPromptsResult, ErrorData> {
self.list_prompts_calls.fetch_add(1, Ordering::SeqCst);
Ok(ListPromptsResult::with_all_items(vec![Prompt::new(
"summarize",
Some("Summarize a document"),
Some(vec![
PromptArgument::new("path")
.with_description("Document path")
.with_required(true),
PromptArgument::new("style"),
]),
)]))
}
}
pub(crate) async fn fixture_runtime(
fixture: FixtureServer,
) -> (McpRuntime, RunningService<RoleServer, FixtureServer>) {
let (client_io, server_io) = tokio::io::duplex(4096);
let (server, client) = tokio::join!(fixture.serve(server_io), ().serve(client_io));
let mut runtime = McpRuntime::new();
runtime.insert("fixture".to_string(), Arc::new(client.unwrap()));
(runtime, server.unwrap())
}
}
#[cfg(test)]
mod tests {
use super::test_fixtures::{FixtureServer, fixture_runtime};
use super::*;
use crate::function::ToolCall;
use log::{Level, LevelFilter, Log, Metadata, Record};
use std::sync::atomic::Ordering;
use std::sync::{Mutex, Once, OnceLock};
struct WarnCollector;
static WARN_MESSAGES: OnceLock<Mutex<Vec<String>>> = OnceLock::new();
fn warn_messages() -> &'static Mutex<Vec<String>> {
WARN_MESSAGES.get_or_init(Mutex::default)
}
impl Log for WarnCollector {
fn enabled(&self, metadata: &Metadata) -> bool {
metadata.level() <= Level::Warn
}
fn log(&self, record: &Record) {
if self.enabled(record.metadata()) {
warn_messages()
.lock()
.unwrap()
.push(record.args().to_string());
}
}
fn flush(&self) {}
}
fn install_warn_collector() {
static INSTALL: Once = Once::new();
INSTALL.call_once(|| {
log::set_logger(&WarnCollector).expect("no other logger should be installed");
log::set_max_level(LevelFilter::Warn);
});
}
#[test]
fn mcp_runtime_new_is_empty() {
@@ -188,4 +526,234 @@ mod tests {
let dummy_call = ToolCall::default();
assert!(scope.tool_tracker.check_loop(&dummy_call).is_none());
}
#[test]
fn uri_template_variables_extracts_placeholders() {
assert_eq!(
uri_template_variables("file:///{path}/{name}"),
vec!["path", "name"]
);
assert!(uri_template_variables("file:///static").is_empty());
}
#[tokio::test]
async fn catalog_items_keeps_tool_and_resource_with_same_id() {
let fixture = FixtureServer {
resources_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let items = runtime.catalog_items("fixture").await.unwrap();
assert!(items.contains_key("tool:dup"));
assert!(items.contains_key("resource:dup"));
assert!(items.contains_key("resource_template:file:///{path}/{name}"));
}
#[tokio::test]
async fn catalog_items_degrades_when_resource_listing_fails() {
let fixture = FixtureServer {
resources_capability: true,
fail_resource_listings: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let items = runtime.catalog_items("fixture").await.unwrap();
assert!(items.contains_key("tool:dup"));
assert!(!items.keys().any(|key| key.starts_with("resource")));
}
#[tokio::test]
async fn catalog_items_warns_when_resource_listing_fails() {
install_warn_collector();
let fixture = FixtureServer {
resources_capability: true,
fail_resource_listings: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
runtime.catalog_items("fixture").await.unwrap();
let messages = warn_messages().lock().unwrap();
assert!(
messages
.iter()
.any(|msg| msg.contains("Failed to list resources on MCP server fixture")),
"missing resource-listing warning in: {messages:?}"
);
}
#[tokio::test]
async fn catalog_items_skips_unadvertised_capabilities() {
let fixture = FixtureServer::default();
let resources_calls = Arc::clone(&fixture.list_resources_calls);
let prompts_calls = Arc::clone(&fixture.list_prompts_calls);
let (runtime, _server) = fixture_runtime(fixture).await;
let items = runtime.catalog_items("fixture").await.unwrap();
assert_eq!(items.keys().collect::<Vec<_>>(), vec!["tool:dup"]);
assert_eq!(resources_calls.load(Ordering::SeqCst), 0);
assert_eq!(prompts_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn catalog_items_includes_prompts_when_advertised() {
let fixture = FixtureServer {
prompts_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let items = runtime.catalog_items("fixture").await.unwrap();
let prompt = items.get("prompt:summarize").unwrap();
assert_eq!(prompt.kind, CatalogItemKind::Prompt);
assert_eq!(prompt.description, "Summarize a document");
}
#[tokio::test]
async fn search_results_carry_kind_and_uri() {
let fixture = FixtureServer {
resources_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let results = runtime.search("fixture", "dup", 10).await.unwrap();
let values: Vec<Value> = results
.iter()
.map(|item| serde_json::to_value(item).unwrap())
.collect();
let tool = values.iter().find(|v| v["kind"] == "tool").unwrap();
assert!(tool.get("uri").is_none());
let resource = values.iter().find(|v| v["kind"] == "resource").unwrap();
assert_eq!(resource["uri"], "dup");
assert_eq!(resource["mime_type"], "text/plain");
assert_eq!(resource["size"], 42);
}
#[tokio::test]
async fn describe_tool_keeps_existing_schema_shape() {
let fixture = FixtureServer::default();
let (runtime, _server) = fixture_runtime(fixture).await;
let result = runtime.describe("fixture", "tool", "dup").await.unwrap();
assert_eq!(
result,
json!({
"type": "object",
"properties": {
"tool": { "type": "string" },
"arguments": {
"type": "object",
"properties": { "q": { "type": "string" } }
}
}
})
);
}
#[tokio::test]
async fn describe_resource_returns_metadata() {
let fixture = FixtureServer {
resources_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let result = runtime
.describe("fixture", "resource", "dup")
.await
.unwrap();
assert_eq!(result["uri"], "dup");
assert_eq!(result["name"], "dup-resource");
assert_eq!(result["description"], "Duplicate-named resource");
assert_eq!(result["mime_type"], "text/plain");
assert_eq!(result["size"], 42);
}
#[tokio::test]
async fn describe_resource_template_returns_variables() {
let fixture = FixtureServer {
resources_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let result = runtime
.describe("fixture", "resource_template", "file:///{path}/{name}")
.await
.unwrap();
assert_eq!(result["uri_template"], "file:///{path}/{name}");
assert_eq!(result["name"], "file-template");
assert_eq!(result["variables"], json!(["path", "name"]));
}
#[tokio::test]
async fn describe_prompt_returns_arguments() {
let fixture = FixtureServer {
prompts_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let result = runtime
.describe("fixture", "prompt", "summarize")
.await
.unwrap();
assert_eq!(result["name"], "summarize");
assert_eq!(result["description"], "Summarize a document");
assert_eq!(result["arguments"][0]["name"], "path");
assert_eq!(result["arguments"][0]["description"], "Document path");
assert_eq!(result["arguments"][0]["required"], true);
assert_eq!(result["arguments"][1]["name"], "style");
assert_eq!(result["arguments"][1]["required"], Value::Null);
}
#[tokio::test]
async fn describe_unknown_kind_lists_valid_kinds() {
let fixture = FixtureServer::default();
let (runtime, _server) = fixture_runtime(fixture).await;
let err = runtime
.describe("fixture", "widget", "dup")
.await
.unwrap_err()
.to_string();
assert!(err.contains("widget"));
for kind in ["tool", "resource", "resource_template", "prompt"] {
assert!(err.contains(kind), "missing {kind} in: {err}");
}
}
#[tokio::test]
async fn describe_missing_resource_names_kind_in_error() {
let fixture = FixtureServer {
resources_capability: true,
..Default::default()
};
let (runtime, _server) = fixture_runtime(fixture).await;
let err = runtime
.describe("fixture", "resource", "file:///missing")
.await
.unwrap_err()
.to_string();
assert_eq!(
err,
"file:///missing not found in fixture MCP server resource catalog"
);
}
}
+61 -2
View File
@@ -669,6 +669,18 @@ impl Functions {
..Default::default()
},
);
describe_function_properties.insert(
"kind".to_string(),
JsonSchema {
type_value: Some("string".to_string()),
description: Some(
"Catalog item kind: tool (default), resource, resource_template, or prompt"
.into(),
),
default: Some(Value::from("tool")),
..Default::default()
},
);
for server in mcp_servers {
let search_function_name = format!("{}_{server}", MCP_SEARCH_META_FUNCTION_NAME_PREFIX);
@@ -709,7 +721,9 @@ impl Functions {
};
let describe_functions_declaration = FunctionDeclaration {
name: describe_function_name.clone(),
description: "Get the full JSON schema for exactly one MCP tool.".to_string(),
description: "Get the full schema or metadata for exactly one MCP catalog item: \
a tool, resource, resource template, or prompt."
.to_string(),
parameters: JsonSchema {
type_value: Some("object".to_string()),
properties: Some(describe_function_properties.clone()),
@@ -1378,10 +1392,16 @@ impl ToolCall {
.ok_or_else(|| anyhow!("Missing 'tool' in arguments"))?
.as_str()
.ok_or_else(|| anyhow!("Invalid 'tool' in arguments"))?;
let kind = match json_data.get("kind") {
Some(value) => value
.as_str()
.ok_or_else(|| anyhow!("Invalid 'kind' in arguments"))?,
None => "tool",
};
let result = ctx
.tool_scope
.mcp_runtime
.describe(&server_id, tool)
.describe(&server_id, kind, tool)
.await?;
Ok(serde_json::to_value(result)?)
}
@@ -1825,6 +1845,7 @@ fn format_json_colored_keys(value: &serde_json::Value) -> String {
#[cfg(test)]
mod tests {
use super::*;
use crate::config::test_fixtures::{FixtureServer, fixture_runtime};
use crate::config::{AppState, WorkingMode};
use crate::supervisor::escalation::{EscalationQueue, EscalationRequest};
use serde_json::json;
@@ -2292,6 +2313,44 @@ mod tests {
assert!(props.contains_key("tool"));
}
#[test]
fn functions_mcp_describe_declaration_has_optional_kind_param() {
let mut f = Functions::default();
f.append_mcp_meta_functions(vec!["srv".to_string()]);
let decl = f.find("mcp_describe_srv").unwrap();
let props = decl.parameters.properties.as_ref().unwrap();
let kind = props.get("kind").unwrap();
assert_eq!(kind.default, Some(Value::from("tool")));
let required = decl.parameters.required.as_ref().unwrap();
assert_eq!(required, &vec!["tool".to_string()]);
}
#[test]
fn eval_mcp_describe_without_kind_defaults_to_tool() {
let output = run_async(async {
let (runtime, _server) = fixture_runtime(FixtureServer::default()).await;
let mut ctx = RequestContext::new(Arc::new(AppState::test_default()), WorkingMode::Cmd);
ctx.tool_scope.mcp_runtime = runtime;
let call = call_with_args("mcp_describe_fixture", json!({"tool": "dup"}));
call.eval_mcp(&ctx).await
})
.unwrap();
assert_eq!(
output,
json!({
"type": "object",
"properties": {
"tool": { "type": "string" },
"arguments": {
"type": "object",
"properties": { "q": { "type": "string" } }
}
}
})
);
}
#[test]
fn functions_supervisor_includes_task_queue_tools() {
let mut f = Functions::default();
+39 -44
View File
@@ -60,24 +60,45 @@ pub fn mcp_meta_function_names(server: &str) -> Vec<String> {
pub type ConnectedServer = RunningService<RoleClient, ()>;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum CatalogItemKind {
#[default]
Tool,
Resource,
ResourceTemplate,
Prompt,
}
impl CatalogItemKind {
pub fn as_str(self) -> &'static str {
match self {
Self::Tool => "tool",
Self::Resource => "resource",
Self::ResourceTemplate => "resource_template",
Self::Prompt => "prompt",
}
}
}
impl Display for CatalogItemKind {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Default, Serialize)]
pub struct CatalogItem {
pub kind: CatalogItemKind,
pub name: String,
pub server: String,
pub description: String,
}
#[derive(Debug)]
struct ServerCatalog {
items: HashMap<String, CatalogItem>,
}
impl Clone for ServerCatalog {
fn clone(&self) -> Self {
Self {
items: self.items.clone(),
}
}
#[serde(skip_serializing_if = "Option::is_none")]
pub uri: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub mime_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub size: Option<u64>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
@@ -183,7 +204,6 @@ pub struct McpRegistry {
log_path: Option<PathBuf>,
config: Option<McpServersConfig>,
servers: HashMap<String, Arc<ConnectedServer>>,
catalogs: HashMap<String, ServerCatalog>,
}
impl McpRegistry {
@@ -326,7 +346,7 @@ impl McpRegistry {
debug!("Starting selected MCP servers: {:?}", ids_to_start);
let results: Vec<Option<(String, Arc<ConnectedServer>, ServerCatalog)>> = stream::iter(
let results: Vec<Option<(String, Arc<ConnectedServer>)>> = stream::iter(
ids_to_start
.into_iter()
.map(|id| async { self.start_server(id).await }),
@@ -335,18 +355,14 @@ impl McpRegistry {
.try_collect()
.await?;
for (id, server, catalog) in results.into_iter().flatten() {
self.servers.insert(id.clone(), server);
self.catalogs.insert(id, catalog);
for (id, server) in results.into_iter().flatten() {
self.servers.insert(id, server);
}
Ok(())
}
async fn start_server(
&self,
id: String,
) -> Result<Option<(String, Arc<ConnectedServer>, ServerCatalog)>> {
async fn start_server(&self, id: String) -> Result<Option<(String, Arc<ConnectedServer>)>> {
let spec = self
.config
.as_ref()
@@ -370,30 +386,9 @@ impl McpRegistry {
Err(e) => return Err(e),
};
let tools = service.list_all_tools().await?;
debug!("Available tools for MCP server {id}: {tools:?}");
let mut items_vec = Vec::new();
for t in tools {
let name = t.name.to_string();
let description = t.description.unwrap_or_default().to_string();
items_vec.push(CatalogItem {
name,
server: id.clone(),
description,
});
}
let mut items_map = HashMap::new();
items_vec.into_iter().for_each(|it| {
items_map.insert(it.name.clone(), it);
});
let catalog = ServerCatalog { items: items_map };
info!("Started MCP server: {id}");
Ok(Some((id.to_string(), service, catalog)))
Ok(Some((id, service)))
}
fn resolve_server_ids(&self, enabled_mcp_servers: Option<Vec<String>>) -> Vec<String> {