use std::sync::Arc;
use async_trait::async_trait;
use rmcp::model::{CallToolRequestParams, ClientCapabilities, ClientInfo, Implementation};
use rmcp::transport::StreamableHttpClientTransport;
use rmcp::ServiceExt;
use serde_json::Value;
use crate::config::McpServerConfig;
use crate::tools::Tool;
pub enum Transport {
Http {
url: String,
token: Option<String>,
},
HttpOAuth {
url: String,
},
Stdio {
command: String,
args: Vec<String>,
},
}
impl std::fmt::Debug for Transport {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Transport::Http { url, token } => f
.debug_struct("Http")
.field("url", url)
.field("token", &token.as_ref().map(|_| "<redacted>"))
.finish(),
Transport::HttpOAuth { url } => f.debug_struct("HttpOAuth").field("url", url).finish(),
Transport::Stdio { command, args } => f
.debug_struct("Stdio")
.field("command", command)
.field("args", args)
.finish(),
}
}
}
impl McpServerConfig {
pub fn transport(&self) -> Result<Transport, String> {
if self.name.is_empty() {
return Err("An mcp_servers entry is missing 'name'".to_string());
}
match (&self.url, &self.command) {
(Some(_), Some(_)) => Err(format!(
"'{}' sets both url and command; use one transport per entry",
self.name
)),
(None, None) => Err(format!(
"'{}' sets neither url nor command, so there is nothing to connect to",
self.name
)),
(Some(url), None) if self.auth.as_deref() == Some("oauth") => {
if self.token_env.is_some() {
return Err(format!(
"'{}' asks for oauth and also sets token_env; pick one",
self.name
));
}
Ok(Transport::HttpOAuth { url: url.clone() })
}
(Some(url), None) => {
if let Some(other) = self.auth.as_deref() {
return Err(format!(
"'{}' has auth = \"{}\", which is not a supported value; use \"oauth\" or \
omit it",
self.name, other
));
}
let token = match &self.token_env {
Some(var) => match std::env::var(var) {
Ok(token) if !token.is_empty() => Some(token),
_ => {
return Err(format!(
"{} is not set, so no credential is available for '{}'",
var, self.name
))
}
},
None => None,
};
Ok(Transport::Http {
url: url.clone(),
token,
})
}
(None, Some(command)) => {
if self.auth.is_some() {
return Err(format!(
"'{}' is a stdio server, so auth does not apply",
self.name
));
}
if self.token_env.is_some() {
return Err(format!(
"'{}' is a stdio server, so token_env does not apply; pass secrets through \
the environment it inherits",
self.name
));
}
Ok(Transport::Stdio {
command: command.clone(),
args: self.args.clone(),
})
}
}
}
pub fn endpoint_label(&self) -> String {
match (&self.url, &self.command) {
(Some(url), _) => url.clone(),
(None, Some(command)) if self.args.is_empty() => command.clone(),
(None, Some(command)) => format!("{} {}", command, self.args.join(" ")),
(None, None) => "unconfigured".to_string(),
}
}
}
type Client = rmcp::service::RunningService<rmcp::service::RoleClient, ClientInfo>;
fn qualified_name(server: &str, tool: &str) -> String {
format!("{}__{}", server, tool)
}
pub struct McpConnection {
config: McpServerConfig,
client: tokio::sync::RwLock<Option<Arc<Client>>>,
last_failure: tokio::sync::Mutex<Option<(std::time::Instant, String)>>,
}
const RECONNECT_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(5);
impl McpConnection {
fn new(config: McpServerConfig, client: Arc<Client>) -> Self {
Self {
config,
client: tokio::sync::RwLock::new(Some(client)),
last_failure: tokio::sync::Mutex::new(None),
}
}
pub async fn client(&self) -> Result<Arc<Client>, String> {
if let Some(client) = self.client.read().await.clone() {
return Ok(client);
}
let mut slot = self.client.write().await;
if let Some(client) = slot.clone() {
return Ok(client);
}
let mut failure = self.last_failure.lock().await;
if let Some((at, reason)) = failure.as_ref() {
if at.elapsed() < RECONNECT_COOLDOWN {
return Err(format!(
"'{}' is disconnected: {}. Not retrying for another {}s.",
self.config.name,
reason,
(RECONNECT_COOLDOWN - at.elapsed()).as_secs() + 1
));
}
}
match connect(&self.config).await {
Ok(client) => {
*slot = Some(client.clone());
*failure = None;
Ok(client)
}
Err(e) => {
*failure = Some((std::time::Instant::now(), e.clone()));
Err(e)
}
}
}
pub async fn invalidate(&self) {
*self.client.write().await = None;
}
}
pub struct McpTool {
name: String,
remote_name: String,
description: String,
input_schema: Value,
connection: Arc<McpConnection>,
capability: crate::risk::Capability,
}
fn declared_capability(
annotations: Option<&rmcp::model::ToolAnnotations>,
) -> crate::risk::Capability {
match annotations.and_then(|a| a.read_only_hint) {
Some(true) => crate::risk::Capability::ReadOnly,
_ => crate::risk::Capability::default(),
}
}
#[async_trait]
impl Tool for McpTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
&self.description
}
fn input_schema(&self) -> Value {
self.input_schema.clone()
}
fn capability(&self) -> crate::risk::Capability {
self.capability
}
async fn execute(&self, input: Value) -> Result<String, String> {
let arguments = match input {
Value::Object(map) => Some(map),
Value::Null => None,
other => return Err(format!("Tool input must be a JSON object, got {}", other)),
};
let mut params = CallToolRequestParams::new(self.remote_name.clone());
if let Some(arguments) = arguments {
params = params.with_arguments(arguments);
}
let client = self.connection.client().await?;
let result = match client.call_tool(params).await {
Ok(result) => result,
Err(e) => {
self.connection.invalidate().await;
return Err(format!("{} failed: {}", self.name, e));
}
};
let rendered = render_content(&result);
if result.is_error.unwrap_or(false) {
return Err(if rendered.is_empty() {
format!("{} reported an error with no detail", self.name)
} else {
rendered
});
}
Ok(if rendered.is_empty() {
format!("{} returned no content", self.name)
} else {
rendered
})
}
}
fn render_content(result: &rmcp::model::CallToolResult) -> String {
let mut parts: Vec<String> = Vec::new();
let mut skipped = 0usize;
for block in &result.content {
match block.as_text() {
Some(text) => parts.push(text.text.clone()),
None => skipped += 1,
}
}
if let Some(structured) = &result.structured_content {
parts.push(
serde_json::to_string_pretty(structured).unwrap_or_else(|_| structured.to_string()),
);
}
if skipped > 0 {
parts.push(format!("[{} non-text block(s) omitted]", skipped));
}
parts.join("\n")
}
fn client_info() -> ClientInfo {
ClientInfo::new(
ClientCapabilities::default(),
Implementation::new("procyon", env!("CARGO_PKG_VERSION")),
)
}
const STDERR_TAIL_LINES: usize = 10;
const STDERR_FLUSH_GRACE: std::time::Duration = std::time::Duration::from_millis(300);
type StderrTail = (
Arc<std::sync::Mutex<Vec<String>>>,
tokio::task::JoinHandle<()>,
);
fn drain_stderr(stderr: tokio::process::ChildStderr) -> StderrTail {
let tail = Arc::new(std::sync::Mutex::new(Vec::new()));
let sink = tail.clone();
let handle = tokio::spawn(async move {
use tokio::io::{AsyncBufReadExt, BufReader};
let mut lines = BufReader::new(stderr).lines();
while let Ok(Some(line)) = lines.next_line().await {
if let Ok(mut tail) = sink.lock() {
if tail.len() == STDERR_TAIL_LINES {
tail.remove(0);
}
tail.push(line);
}
}
});
(tail, handle)
}
fn stderr_context(tail: &Arc<std::sync::Mutex<Vec<String>>>) -> String {
match tail.lock() {
Ok(tail) if !tail.is_empty() => format!("\nIts stderr said:\n {}", tail.join("\n ")),
_ => String::new(),
}
}
async fn connect_stdio(name: &str, command: &str, args: &[String]) -> Result<Arc<Client>, String> {
let mut cmd = tokio::process::Command::new(command);
cmd.args(args);
let (process, stderr) = rmcp::transport::TokioChildProcess::builder(cmd)
.stderr(std::process::Stdio::piped())
.spawn()
.map_err(|e| format!("Failed to launch '{}' ({}): {}", name, command, e))?;
let drain = stderr.map(drain_stderr);
match client_info().serve(process).await {
Ok(client) => Ok(Arc::new(client)),
Err(e) => {
let context = match drain {
Some((tail, handle)) => {
let _ = tokio::time::timeout(STDERR_FLUSH_GRACE, handle).await;
stderr_context(&tail)
}
None => String::new(),
};
Err(format!(
"Failed to speak MCP with '{}' ({}): {}{}",
name, command, e, context
))
}
}
}
async fn connect_http(name: &str, url: &str, token: Option<String>) -> Result<Arc<Client>, String> {
let mut transport_config =
rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig::with_uri(
url.to_string(),
)
.reinit_on_expired_session(true);
if let Some(token) = token {
transport_config = transport_config.auth_header(token);
}
let transport = StreamableHttpClientTransport::from_config(transport_config);
client_info()
.serve(transport)
.await
.map(Arc::new)
.map_err(|e| format!("Failed to connect to '{}' at {}: {}", name, url, e))
}
const AUTHORIZE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(180);
fn credential_store(name: &str) -> Result<crate::oauth::FileCredentialStore, String> {
Ok(crate::oauth::FileCredentialStore::new(
crate::oauth::credentials_path(name)?,
))
}
fn build_oauth_transport(
url: &str,
manager: rmcp::transport::auth::AuthorizationManager,
) -> rmcp::transport::StreamableHttpClientTransport<
rmcp::transport::auth::AuthClient<reqwest::Client>,
> {
let auth_client = rmcp::transport::auth::AuthClient::new(reqwest::Client::default(), manager);
let config =
rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig::with_uri(
url.to_string(),
)
.reinit_on_expired_session(true);
rmcp::transport::StreamableHttpClientTransport::with_client(auth_client, config)
}
async fn connect_oauth_stored(name: &str, url: &str) -> Result<Arc<Client>, String> {
let mut manager = rmcp::transport::auth::AuthorizationManager::new(url)
.await
.map_err(|e| {
format!(
"Failed to reach the authorization server of '{}': {}",
name, e
)
})?;
manager.set_credential_store(credential_store(name)?);
let restored = manager
.initialize_from_store()
.await
.map_err(|e| format!("Failed to read stored credentials for '{}': {}", name, e))?;
if !restored {
return Err(format!(
"'{}' is not authorized yet. Run `procyon --authorize {}`.",
name, name
));
}
client_info()
.serve(build_oauth_transport(url, manager))
.await
.map(Arc::new)
.map_err(|e| format!("Failed to connect to '{}' at {}: {}", name, url, e))
}
pub async fn authorize(config: &McpServerConfig) -> Result<(), String> {
let Transport::HttpOAuth { url } = config.transport()? else {
return Err(format!("'{}' is not configured for oauth", config.name));
};
let listener = crate::oauth::bind_redirect_listener().await?;
let mut manager = rmcp::transport::auth::AuthorizationManager::new(&url)
.await
.map_err(|e| format!("Failed to reach the authorization server: {}", e))?;
manager.set_credential_store(credential_store(&config.name)?);
let resolution = manager
.resolve_metadata_from_challenge(None)
.await
.map_err(|e| format!("Failed to discover the authorization server: {}", e))?;
manager.set_metadata(resolution.metadata);
let request = rmcp::transport::auth::AuthorizationRequest::new(crate::oauth::redirect_uri())
.with_client_name("procyon");
let session = rmcp::transport::auth::AuthorizationSession::new(manager, request)
.await
.map_err(|(_manager, e)| format!("Failed to start authorization: {}", e))?;
let auth_url = session.get_authorization_url().to_string();
println!("Opening your browser to authorize '{}'.", config.name);
println!("If it does not open, visit:\n\n {}\n", auth_url);
if !crate::oauth::open_browser(&auth_url).await {
println!("(could not launch a browser automatically)");
}
println!(
"Waiting for the redirect on {} ...",
crate::oauth::redirect_uri()
);
let _ = std::io::Write::flush(&mut std::io::stdout());
let callback = crate::oauth::wait_for_redirect(listener, AUTHORIZE_TIMEOUT).await?;
session
.handle_callback_url(&callback.url)
.await
.map_err(|e| format!("Failed to exchange the authorization code: {}", e))?;
let stored = {
use rmcp::transport::auth::CredentialStore;
credential_store(&config.name)?.load().await.ok().flatten()
};
match stored.as_ref().and_then(crate::oauth::expires_at) {
Some(at) => {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
println!(
"'{}' is authorized. The access token lasts {}s and refreshes automatically.",
config.name,
at.saturating_sub(now)
);
}
None => println!("'{}' is authorized.", config.name),
}
println!(
"Credentials stored in {}",
crate::oauth::credentials_path(&config.name)?.display()
);
Ok(())
}
async fn connect(config: &McpServerConfig) -> Result<Arc<Client>, String> {
match config.transport()? {
Transport::Http { url, token } => connect_http(&config.name, &url, token).await,
Transport::HttpOAuth { url } => connect_oauth_stored(&config.name, &url).await,
Transport::Stdio { command, args } => connect_stdio(&config.name, &command, &args).await,
}
}
pub async fn load_servers(
configs: &[McpServerConfig],
) -> (Vec<Box<dyn Tool>>, Vec<String>, Vec<String>) {
let mut tools: Vec<Box<dyn Tool>> = Vec::new();
let mut connected = Vec::new();
let mut problems = Vec::new();
for config in configs {
let client = match connect(config).await {
Ok(client) => client,
Err(e) => {
problems.push(e);
continue;
}
};
let listed = match client.list_tools(Default::default()).await {
Ok(listed) => listed,
Err(e) => {
problems.push(format!("Could not list tools of '{}': {}", config.name, e));
continue;
}
};
let connection = Arc::new(McpConnection::new(config.clone(), client));
let mut names = Vec::new();
for tool in listed.tools {
let name = qualified_name(&config.name, &tool.name);
names.push(name.clone());
tools.push(Box::new(McpTool {
name,
remote_name: tool.name.to_string(),
description: tool.description.map(|d| d.to_string()).unwrap_or_else(|| {
format!("Tool '{}' on MCP server '{}'", tool.name, config.name)
}),
input_schema: Value::Object((*tool.input_schema).clone()),
connection: connection.clone(),
capability: declared_capability(tool.annotations.as_ref()),
}));
}
let endpoint = config.endpoint_label();
connected.push(if names.is_empty() {
format!("{} ({}) — no tools exposed", config.name, endpoint)
} else {
format!("{} ({}): {}", config.name, endpoint, names.join(", "))
});
}
(tools, connected, problems)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn result_with(text: &str) -> rmcp::model::CallToolResult {
rmcp::model::CallToolResult::success(vec![rmcp::model::ContentBlock::text(text)])
}
#[test]
fn tool_names_are_namespaced_by_server() {
assert_eq!(qualified_name("raven", "search"), "raven__search");
}
#[test]
fn text_content_is_rendered() {
assert_eq!(render_content(&result_with("hello")), "hello");
}
#[test]
fn multiple_text_blocks_are_joined() {
let result = rmcp::model::CallToolResult::success(vec![
rmcp::model::ContentBlock::text("one"),
rmcp::model::ContentBlock::text("two"),
]);
assert_eq!(render_content(&result), "one\ntwo");
}
#[test]
fn an_empty_result_renders_empty() {
let result = rmcp::model::CallToolResult::success(vec![]);
assert!(render_content(&result).is_empty());
}
#[tokio::test]
async fn a_missing_credential_is_reported_not_silently_ignored() {
let config = McpServerConfig {
name: "raven".to_string(),
url: Some("https://example.invalid/mcp".to_string()),
token_env: Some("PROCYON_TEST_TOKEN_THAT_IS_UNSET".to_string()),
..Default::default()
};
let err = connect(&config).await.unwrap_err();
assert!(
err.contains("PROCYON_TEST_TOKEN_THAT_IS_UNSET"),
"got {}",
err
);
assert!(err.contains("raven"), "got {}", err);
}
#[tokio::test]
async fn an_unreachable_server_is_skipped_rather_than_fatal() {
let configs = vec![McpServerConfig {
name: "nowhere".to_string(),
url: Some("http://127.0.0.1:1/mcp".to_string()),
..Default::default()
}];
let (tools, connected, problems) = load_servers(&configs).await;
assert!(tools.is_empty());
assert!(connected.is_empty());
assert_eq!(problems.len(), 1, "got {:?}", problems);
assert!(problems[0].contains("nowhere"));
}
#[tokio::test]
#[ignore]
async fn raven_live_lists_and_calls_search() {
if std::env::var("PROCYON_MCP_TOKEN")
.map(|t| t.is_empty())
.unwrap_or(true)
{
println!(
"skipped: PROCYON_MCP_TOKEN is not set (the OAuth path is covered separately)"
);
return;
}
let url = std::env::var("PROCYON_MCP_URL")
.unwrap_or_else(|_| "https://raven.stellar.org/mcp".to_string());
let configs = vec![McpServerConfig {
name: "raven".to_string(),
url: Some(url),
token_env: Some("PROCYON_MCP_TOKEN".to_string()),
..Default::default()
}];
let (tools, connected, problems) = load_servers(&configs).await;
assert!(problems.is_empty(), "connection problems: {:?}", problems);
assert!(!connected.is_empty());
println!("connected: {:?}", connected);
let names: Vec<&str> = tools.iter().map(|t| t.name()).collect();
println!("tools: {:?}", names);
assert!(
names.contains(&"raven__search"),
"expected a namespaced search tool, got {:?}",
names
);
let search = tools
.iter()
.find(|t| t.name() == "raven__search")
.expect("search tool");
assert_eq!(search.input_schema()["type"], "object");
let out = search
.execute(json!({"query": "soroban storage ttl"}))
.await
.expect("search should succeed");
println!("--- search result ---\n{}", out);
assert!(!out.trim().is_empty());
}
#[test]
fn session_recovery_is_enabled_by_default_in_the_sdk() {
let config =
rmcp::transport::streamable_http_client::StreamableHttpClientTransportConfig::with_uri(
"https://example.invalid/mcp",
);
assert!(
config.reinit_on_expired_session,
"procyon relies on transparent re-initialization after HTTP 404"
);
}
#[test]
fn debug_never_prints_the_credential() {
std::env::set_var("PROCYON_TEST_DEBUG_TOKEN", "super-secret-value");
let config = McpServerConfig {
name: "raven".to_string(),
url: Some("https://x/mcp".to_string()),
token_env: Some("PROCYON_TEST_DEBUG_TOKEN".to_string()),
..Default::default()
};
let rendered = format!("{:?}", config.transport().unwrap());
std::env::remove_var("PROCYON_TEST_DEBUG_TOKEN");
assert!(
!rendered.contains("super-secret-value"),
"leaked: {}",
rendered
);
assert!(rendered.contains("redacted"), "got {}", rendered);
}
#[test]
fn a_stdio_entry_resolves_to_the_stdio_transport() {
let config = McpServerConfig {
name: "fs".to_string(),
command: Some("npx".to_string()),
args: vec!["-y".to_string(), "server-filesystem".to_string()],
..Default::default()
};
match config.transport().unwrap() {
Transport::Stdio { command, args } => {
assert_eq!(command, "npx");
assert_eq!(args, vec!["-y", "server-filesystem"]);
}
_ => panic!("expected stdio"),
}
assert_eq!(config.endpoint_label(), "npx -y server-filesystem");
}
#[test]
fn declaring_both_transports_is_rejected() {
let config = McpServerConfig {
name: "both".to_string(),
url: Some("https://x/mcp".to_string()),
command: Some("npx".to_string()),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("one transport per entry"), "got {}", err);
}
#[test]
fn declaring_neither_transport_is_rejected() {
let config = McpServerConfig {
name: "empty".to_string(),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("neither url nor command"), "got {}", err);
}
#[test]
fn an_entry_without_a_name_is_rejected() {
let config = McpServerConfig {
command: Some("npx".to_string()),
..Default::default()
};
assert!(config.transport().unwrap_err().contains("missing 'name'"));
}
#[test]
fn token_env_on_a_stdio_server_is_rejected_rather_than_ignored() {
let config = McpServerConfig {
name: "fs".to_string(),
command: Some("npx".to_string()),
token_env: Some("SOME_TOKEN".to_string()),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("token_env does not apply"), "got {}", err);
}
#[tokio::test]
async fn a_command_that_does_not_exist_is_reported_with_its_name() {
let configs = vec![McpServerConfig {
name: "ghost".to_string(),
command: Some("procyon-no-such-mcp-server".to_string()),
..Default::default()
}];
let (tools, connected, problems) = load_servers(&configs).await;
assert!(tools.is_empty() && connected.is_empty());
assert_eq!(problems.len(), 1, "got {:?}", problems);
assert!(problems[0].contains("ghost"), "got {}", problems[0]);
}
#[tokio::test]
async fn a_process_that_is_not_an_mcp_server_reports_its_stderr() {
let configs = vec![McpServerConfig {
name: "notmcp".to_string(),
command: Some("ls".to_string()),
args: vec!["/procyon-no-such-path".to_string()],
..Default::default()
}];
let (_, _, problems) = load_servers(&configs).await;
assert_eq!(problems.len(), 1, "got {:?}", problems);
assert!(
problems[0].contains("stderr said"),
"the child's own explanation must survive: {}",
problems[0]
);
}
#[tokio::test]
async fn a_disconnected_server_reports_the_cause_and_backs_off() {
let config = McpServerConfig {
name: "nowhere".to_string(),
url: Some("http://127.0.0.1:1/mcp".to_string()),
..Default::default()
};
let connection = McpConnection {
config,
client: tokio::sync::RwLock::new(None),
last_failure: tokio::sync::Mutex::new(None),
};
let first = connection.client().await.unwrap_err();
assert!(first.contains("nowhere"), "got {}", first);
let second = connection.client().await.unwrap_err();
assert!(
second.contains("Not retrying"),
"expected a backoff message, got {}",
second
);
assert!(second.contains("nowhere"));
}
#[tokio::test]
#[ignore]
async fn stdio_live_reconnects_after_the_server_dies() {
let config = McpServerConfig {
name: "fs".to_string(),
command: Some("npx".to_string()),
args: vec![
"-y".to_string(),
"@modelcontextprotocol/server-filesystem".to_string(),
".".to_string(),
],
..Default::default()
};
let client = connect(&config).await.expect("initial connect");
let connection = Arc::new(McpConnection::new(config, client));
let tool = McpTool {
name: "fs__list_allowed_directories".to_string(),
remote_name: "list_allowed_directories".to_string(),
description: String::new(),
input_schema: json!({"type": "object"}),
connection: connection.clone(),
capability: crate::risk::Capability::ReadOnly,
};
let before = tool.execute(json!({})).await.expect("first call");
println!("before: {}", before.trim());
connection.invalidate().await;
let after = tool.execute(json!({})).await.expect("call after reconnect");
println!("after: {}", after.trim());
assert_eq!(
before.trim(),
after.trim(),
"a reconnected server must answer the same"
);
}
#[tokio::test]
#[ignore]
async fn stdio_live_lists_filesystem_tools() {
let configs = vec![McpServerConfig {
name: "fs".to_string(),
command: Some("npx".to_string()),
args: vec![
"-y".to_string(),
"@modelcontextprotocol/server-filesystem".to_string(),
".".to_string(),
],
..Default::default()
}];
let (tools, connected, problems) = load_servers(&configs).await;
assert!(problems.is_empty(), "problems: {:?}", problems);
println!("connected: {:?}", connected);
let names: Vec<&str> = tools.iter().map(|t| t.name()).collect();
println!("tools: {:?}", names);
assert!(!names.is_empty(), "a filesystem server exposes tools");
assert!(
names.iter().all(|n| n.starts_with("fs__")),
"every tool must be namespaced: {:?}",
names
);
assert!(names.contains(&"fs__read_file"), "got {:?}", names);
let list = tools
.iter()
.find(|t| t.name() == "fs__list_allowed_directories")
.expect("list_allowed_directories");
let out = list.execute(json!({})).await.expect("call should succeed");
println!("--- fs__list_allowed_directories ---\n{}", out);
assert!(!out.trim().is_empty(), "a real call must return content");
}
#[tokio::test]
#[ignore]
async fn raven_oauth_live_connects_with_stored_credentials() {
let configs = vec![McpServerConfig {
name: "raven".to_string(),
url: Some(
std::env::var("PROCYON_MCP_URL")
.unwrap_or_else(|_| "https://raven.stellar.org/mcp".to_string()),
),
auth: Some("oauth".to_string()),
..Default::default()
}];
let (tools, connected, problems) = load_servers(&configs).await;
assert!(problems.is_empty(), "problems: {:?}", problems);
println!("connected: {:?}", connected);
let names: Vec<&str> = tools.iter().map(|t| t.name()).collect();
println!("tools: {:?}", names);
assert!(!names.is_empty(), "an authorized Raven exposes tools");
assert!(
names.iter().all(|n| n.starts_with("raven__")),
"tools must be namespaced: {:?}",
names
);
let search = tools
.iter()
.find(|t| t.name() == "raven__search")
.expect("raven exposes search");
assert_eq!(search.input_schema()["type"], "object");
let out = search
.execute(json!({"query": "soroban storage ttl"}))
.await
.expect("search should succeed over OAuth");
println!("--- search ---\n{}", &out[..out.len().min(600)]);
assert!(!out.trim().is_empty());
}
#[tokio::test]
async fn oauth_and_token_env_together_are_rejected() {
let config = McpServerConfig {
name: "raven".to_string(),
url: Some("https://x/mcp".to_string()),
auth: Some("oauth".to_string()),
token_env: Some("SOME_TOKEN".to_string()),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("pick one"), "got {}", err);
}
#[tokio::test]
async fn an_unknown_auth_value_is_rejected() {
let config = McpServerConfig {
name: "x".to_string(),
url: Some("https://x/mcp".to_string()),
auth: Some("basic".to_string()),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("not a supported value"), "got {}", err);
}
#[tokio::test]
async fn auth_on_a_stdio_server_is_rejected() {
let config = McpServerConfig {
name: "fs".to_string(),
command: Some("npx".to_string()),
auth: Some("oauth".to_string()),
..Default::default()
};
let err = config.transport().unwrap_err();
assert!(err.contains("auth does not apply"), "got {}", err);
}
#[tokio::test]
async fn no_configured_servers_is_not_an_error() {
let (tools, connected, problems) = load_servers(&[]).await;
assert!(tools.is_empty() && connected.is_empty() && problems.is_empty());
}
}