use anyhow::Result;
use regex::Regex;
use rmcp::{
model::{CallToolRequestParam, ReadResourceRequestParam},
service::ServiceExt,
transport::TokioChildProcess,
};
use std::process::Stdio;
use terraphim_config::{ConfigBuilder, Haystack, ServiceType};
use tokio::process::Command;
async fn setup_server_command() -> Result<Command> {
if let Ok(bin) = std::env::var("TERRAPHIM_MCP_SERVER_BIN") {
let path = std::path::PathBuf::from(bin);
if path.exists() {
println!("🚀 Using pre-built server binary at {:?}", path);
let mut command = Command::new(path);
command
.env("RUST_BACKTRACE", "1")
.env("RUST_LOG", "debug")
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
return Ok(command);
}
}
let mut build = Command::new("cargo");
build
.arg("build")
.arg("--package")
.arg("terraphim_mcp_server");
if std::env::var_os("CI").is_some() {
build.arg("--features").arg("zlob");
}
let build_status = build.status().await?;
if !build_status.success() {
return Err(anyhow::anyhow!("Failed to build terraphim_mcp_server"));
}
let crate_dir = std::env::current_dir()?;
let binary_name = if cfg!(target_os = "windows") {
"terraphim_mcp_server.exe"
} else {
"terraphim_mcp_server"
};
let candidate_paths = [
crate_dir
.parent()
.and_then(|p| p.parent())
.map(|workspace| workspace.join("target").join("debug").join(binary_name)),
Some(crate_dir.join("target").join("debug").join(binary_name)),
];
let binary_path = candidate_paths
.into_iter()
.flatten()
.find(|p| p.exists())
.ok_or_else(|| anyhow::anyhow!("Built binary not found in expected locations"))?;
println!("🚀 Using server binary at {:?}", binary_path);
let mut command = Command::new(binary_path);
command
.env("RUST_BACKTRACE", "1")
.env("RUST_LOG", "debug")
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
Ok(command)
}
fn create_test_config() -> String {
let mut config = ConfigBuilder::new()
.build_default_server()
.build()
.expect("Failed to build test configuration");
let docs_src_path = std::env::current_dir()
.expect("Failed to get current directory")
.join("docs/src");
println!("📁 Using docs/src as haystack: {:?}", docs_src_path);
if !docs_src_path.exists() {
println!(
"❌ Warning: docs/src path does not exist: {:?}",
docs_src_path
);
if let Ok(entries) =
std::fs::read_dir(std::env::current_dir().expect("Failed to get current directory"))
{
let dirs: Vec<_> = entries
.filter_map(|entry| entry.ok())
.filter(|entry| entry.file_type().map(|ft| ft.is_dir()).unwrap_or(false))
.collect();
println!(
"📁 Current directory contents: {:?}",
dirs.iter().map(|e| e.path()).collect::<Vec<_>>()
);
}
let workspace_root = std::env::current_dir().expect("Failed to get current directory");
let possible_paths = [
workspace_root.join("docs/src"),
workspace_root.join("..").join("docs/src"),
workspace_root.join("..").join("..").join("docs/src"),
];
for (i, path) in possible_paths.iter().enumerate() {
println!(
"🔍 Trying path {}: {:?} (exists: {})",
i,
path,
path.exists()
);
if path.exists() {
println!("✅ Found docs/src at: {:?}", path);
for role in config.roles.values_mut() {
role.haystacks = vec![Haystack {
location: path.to_string_lossy().to_string(),
service: ServiceType::Ripgrep,
read_only: false,
atomic_server_secret: None,
extra_parameters: std::collections::HashMap::new(),
fetch_content: false,
}];
}
break;
}
}
} else {
println!("✅ docs/src path exists, using it");
for role in config.roles.values_mut() {
role.haystacks = vec![Haystack {
location: docs_src_path.to_string_lossy().to_string(),
service: ServiceType::Ripgrep,
read_only: false,
atomic_server_secret: None,
extra_parameters: std::collections::HashMap::new(),
fetch_content: false,
}];
}
}
serde_json::to_string(&config).expect("Failed to serialize config")
}
fn extract_found_count(message: &str) -> Option<usize> {
let re = Regex::new(r"Found (\d+) documents?").ok()?;
re.captures(message)
.and_then(|cap| cap.get(1))
.and_then(|m| m.as_str().parse::<usize>().ok())
}
#[tokio::test]
async fn test_mcp_server_integration() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
println!("Connected to server: {:?}", service.peer_info());
let tools = service.list_tools(Default::default()).await?;
println!("Available tools: {:#?}", tools);
assert!(!tools.tools.is_empty());
let test_config = create_test_config();
let config_result = service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({
"config_str": test_config
})
.as_object()
.cloned(),
})
.await?;
println!("Update config result: {:#?}", config_result);
assert!(!config_result.is_error.unwrap_or(false));
let search_queries = vec![
("terraphim", "Search for terraphim"),
("machine learning", "Search for machine learning"),
("system operator", "Search for system operator"),
("neural networks", "Search for neural networks"),
];
for (query, description) in search_queries {
println!("Testing: {}", description);
let search_result = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": query,
"role": "Default",
"limit": 5
})
.as_object()
.cloned(),
})
.await?;
println!("Search result for '{}': {:#?}", query, search_result);
assert!(!search_result.is_error.unwrap_or(false));
if let Some(content) = search_result.content.first()
&& let Some(text_content) = content.as_text()
{
println!("Search response: {}", text_content.text);
if let Some(found) = extract_found_count(&text_content.text) {
assert!(
found >= search_result.content.len() - 1,
"Reported document count {} is less than returned resources {}",
found,
search_result.content.len() - 1
);
} else {
panic!("Failed to parse found-count message: {}", text_content.text);
}
}
}
service.cancel().await?;
Ok(())
}
#[tokio::test]
async fn test_search_with_different_roles() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
println!("Connected to server: {:?}", service.peer_info());
let test_config = create_test_config();
let config_result = service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({
"config_str": test_config
})
.as_object()
.cloned(),
})
.await?;
println!("Config update result: {:?}", config_result);
let role_queries = vec![
("Default", "terraphim"),
("Engineer", "system operator"),
("System Operator", "system operator"),
];
for (role, query) in role_queries {
println!("Testing search with role: {}", role);
let search_result = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": query,
"role": role,
"limit": 3
})
.as_object()
.cloned(),
})
.await?;
println!("Search result for role '{}': {:#?}", role, search_result);
if let Some(content) = search_result.content.first()
&& let Some(text_content) = content.as_text()
{
println!("Role '{}' response: {}", role, text_content.text);
if let Some(found) = extract_found_count(&text_content.text) {
if role == "Default" {
assert!(
found > 0,
"Default role should return at least one document"
);
}
}
}
}
Ok(())
}
#[tokio::test]
async fn test_resource_uri_mapping() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
println!("Connected to server: {:?}", service.peer_info());
let test_config = create_test_config();
let config_result = service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({
"config_str": test_config
})
.as_object()
.cloned(),
})
.await?;
println!("Config update result: {:?}", config_result);
let resources = service.list_resources(Default::default()).await?;
println!("Available resources: {:#?}", resources);
if let Some(resource) = resources.resources.first() {
println!("Testing read resource: {}", resource.uri);
let read_result = service
.read_resource(ReadResourceRequestParam {
uri: resource.uri.clone(),
})
.await?;
println!("Read resource result: {:#?}", read_result);
if let Some(content) = read_result.contents.first() {
match content {
rmcp::model::ResourceContents::TextResourceContents { text, .. } => {
println!("Resource text content: {}", text);
assert!(!text.is_empty());
}
rmcp::model::ResourceContents::BlobResourceContents { blob, .. } => {
println!("Resource binary content: {} bytes", blob.len());
assert!(!blob.is_empty());
}
}
}
}
let invalid_result = service
.read_resource(ReadResourceRequestParam {
uri: "invalid://resource/uri".to_string(),
})
.await;
assert!(invalid_result.is_err());
service.cancel().await?;
Ok(())
}
#[tokio::test]
async fn test_simple_search_with_debug() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
println!("Connected to server: {:?}", service.peer_info());
let test_config = create_test_config();
let config_result = service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({
"config_str": test_config
})
.as_object()
.cloned(),
})
.await?;
println!("Config update result: {:?}", config_result);
let search_terms = vec![
"Machine Learning", "Terraphim", "neural", "system", ];
for search_term in search_terms {
println!("Testing search for: '{}'", search_term);
let search_result = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": search_term,
"limit": 10
})
.as_object()
.cloned(),
})
.await?;
println!("Search result for '{}': {:#?}", search_term, search_result);
if let Some(content) = search_result.content.first()
&& let Some(text_content) = content.as_text()
{
println!(
"Search response for '{}': {}",
search_term, text_content.text
);
if text_content.text.contains("Found") && !text_content.text.contains("Found 0") {
println!(
"✅ Found documents for '{}': {}",
search_term, text_content.text
);
} else {
println!(
"❌ No documents found for '{}': {}",
search_term, text_content.text
);
}
}
}
service.cancel().await?;
Ok(())
}
#[tokio::test]
async fn test_search_pagination() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
let test_config = create_test_config();
service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({"config_str": test_config})
.as_object()
.cloned(),
})
.await?;
let first_page = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": "terraphim",
"role": "Default",
"limit": 2
})
.as_object()
.cloned(),
})
.await?;
assert!(!first_page.is_error.unwrap_or(false));
assert!(first_page.content.len() <= 3);
let _first_batch_count = first_page
.content
.iter()
.filter(|c| c.as_resource().is_some())
.count();
let second_page = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": "terraphim",
"role": "Default",
"limit": 2,
"skip": 2
})
.as_object()
.cloned(),
})
.await?;
assert!(!second_page.is_error.unwrap_or(false));
let second_batch_count = second_page
.content
.iter()
.filter(|c| c.as_resource().is_some())
.count();
assert!(second_batch_count <= 2);
Ok(())
}
#[tokio::test]
async fn test_search_invalid_pagination_params() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
let res = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": "test",
"limit": -5
})
.as_object()
.cloned(),
})
.await?;
assert!(res.is_error.unwrap_or(false) || !res.content.is_empty());
let res2 = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": "test",
"limit": 10_000
})
.as_object()
.cloned(),
})
.await?;
assert!(res2.is_error.unwrap_or(false) || !res2.content.is_empty());
Ok(())
}
#[tokio::test]
async fn test_search_read_resource_round_trip() -> Result<()> {
let command = setup_server_command().await?;
let transport = TokioChildProcess::new(command)?;
let service = ().serve(transport).await?;
let test_config = create_test_config();
service
.call_tool(CallToolRequestParam {
name: "update_config_tool".into(),
arguments: serde_json::json!({"config_str": test_config})
.as_object()
.cloned(),
})
.await?;
let search_res = service
.call_tool(CallToolRequestParam {
name: "search".into(),
arguments: serde_json::json!({
"query": "terraphim",
"limit": 1
})
.as_object()
.cloned(),
})
.await?;
assert!(!search_res.is_error.unwrap_or(false));
let resource = search_res
.content
.iter()
.find_map(|c| c.as_resource())
.expect("Expected at least one resource");
let _embedded_text = if let rmcp::model::ResourceContents::TextResourceContents {
text, ..
} = &resource.resource
{
text.clone()
} else {
panic!("Unexpected resource content type");
};
let list_result = service.list_resources(Default::default()).await?;
if list_result.resources.is_empty() {
println!("No resources listed; skipping read_resource round trip test");
return Ok(());
}
let first_uri = list_result.resources[0].uri.clone();
let read_res = service
.read_resource(ReadResourceRequestParam { uri: first_uri })
.await?;
let read_text = match read_res
.contents
.first()
.expect("read_resource returned empty content")
{
rmcp::model::ResourceContents::TextResourceContents { text, .. } => text.clone(),
_ => "".into(),
};
assert!(!read_text.is_empty());
Ok(())
}