#![allow(clippy::unwrap_used, clippy::expect_used)]
use std::sync::Arc;
use serde_json::json;
use wiremock::matchers::{header, method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
use crate::mcp_source::source::{McpToolSource, dispatch_route};
use crate::mcp_source::types::{MCP_ROLE_AGENTIC_WORKER, MCP_ROLE_FLOW_EDITOR};
use crate::tenant::TenantContext;
pub(super) fn initialize_ok() -> ResponseTemplate {
ResponseTemplate::new(200)
.insert_header("Mcp-Session-Id", "sess-1")
.set_body_json(json!({
"jsonrpc": "2.0", "id": 1,
"result": {
"protocolVersion": "2025-06-18",
"serverInfo": { "name": "fake", "version": "1.0.0" }
}
}))
}
pub(super) async fn fake_mcp_server(
tools_json: serde_json::Value,
call_result_json: serde_json::Value,
) -> MockServer {
use wiremock::matchers::body_partial_json;
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(body_partial_json(json!({ "method": "initialize" })))
.respond_with(initialize_ok())
.mount(&server)
.await;
Mock::given(method("POST"))
.and(body_partial_json(
json!({ "method": "notifications/initialized" }),
))
.respond_with(ResponseTemplate::new(202))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(body_partial_json(json!({ "method": "tools/list" })))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"jsonrpc": "2.0", "id": 2,
"result": { "tools": tools_json }
})))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(body_partial_json(json!({ "method": "tools/call" })))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"jsonrpc": "2.0", "id": 3,
"result": call_result_json
})))
.mount(&server)
.await;
server
}
pub(super) fn two_tools() -> serde_json::Value {
json!([
{ "name": "get_issue", "description": "Get an issue",
"inputSchema": { "type": "object", "properties": { "id": { "type": "string" } } } },
{ "name": "search_code", "description": "Search code",
"inputSchema": { "type": "object" } }
])
}
pub(super) fn admin_body_agentic(
id: &str,
url: &str,
allowed: Option<Vec<&str>>,
roles: Vec<&str>,
) -> serde_json::Value {
json!({
"servers": [
{
"id": id,
"name": "Server",
"transport_url": url,
"auth_header_name": null,
"auth_token": null,
"allowed_tools": allowed.map(|v| v.into_iter().collect::<Vec<_>>()),
"roles": roles,
}
]
})
}
pub(super) async fn mount_admin(server: &MockServer, body: serde_json::Value) {
Mock::given(method("GET"))
.and(path("/api/v1/designer/tenant/me/mcp-servers"))
.and(header("authorization", "Bearer gtc_live_x"))
.respond_with(ResponseTemplate::new(200).set_body_json(body))
.mount(server)
.await;
}
pub(super) fn tenant() -> TenantContext {
TenantContext::new("acme", "prod")
}
#[tokio::test]
async fn source_fetches_and_filters_agentic_worker() {
let mcp = fake_mcp_server(two_tools(), json!({})).await;
let admin = MockServer::start().await;
let body = json!({
"servers": [
{
"id": "worker", "name": "Worker", "transport_url": mcp.uri(),
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["agentic_worker"]
},
{
"id": "editor", "name": "Editor", "transport_url": "http://127.0.0.1:1/",
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["flow_editor"]
}
]
});
mount_admin(&admin, body).await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source.catalog(&tenant()).await;
assert_eq!(catalog.len(), 2);
assert!(catalog.tool_entry("worker", "get_issue").is_some());
assert!(catalog.tool_entry("worker", "search_code").is_some());
assert!(catalog.tool_entry("editor", "get_issue").is_none());
}
#[tokio::test]
async fn catalog_for_role_filters_flow_editor() {
let mcp = fake_mcp_server(two_tools(), json!({})).await;
let admin = MockServer::start().await;
let body = json!({
"servers": [
{
"id": "editor", "name": "Editor", "transport_url": mcp.uri(),
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["flow_editor"]
},
{
"id": "worker", "name": "Worker", "transport_url": "http://127.0.0.1:1/",
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["agentic_worker"]
}
]
});
mount_admin(&admin, body).await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source
.catalog_for_role(&tenant(), MCP_ROLE_FLOW_EDITOR)
.await;
assert_eq!(catalog.len(), 2);
assert!(catalog.tool_entry("editor", "get_issue").is_some());
assert!(catalog.tool_entry("editor", "search_code").is_some());
assert!(catalog.tool_entry("worker", "get_issue").is_none());
}
#[tokio::test]
async fn role_catalogs_are_cached_independently() {
let editor_mcp = fake_mcp_server(two_tools(), json!({})).await;
let admin = MockServer::start().await;
let body = json!({
"servers": [
{
"id": "editor", "name": "Editor", "transport_url": editor_mcp.uri(),
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["flow_editor"]
},
{
"id": "worker", "name": "Worker", "transport_url": "http://127.0.0.1:1/",
"auth_header_name": null, "auth_token": null,
"allowed_tools": null, "roles": ["agentic_worker"]
}
]
});
mount_admin(&admin, body).await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let t = tenant();
let editor = source.catalog_for_role(&t, MCP_ROLE_FLOW_EDITOR).await;
let worker = source.catalog_for_role(&t, MCP_ROLE_AGENTIC_WORKER).await;
assert_eq!(editor.len(), 2);
assert_eq!(worker.len(), 0);
assert!(!Arc::ptr_eq(&editor, &worker));
}
#[tokio::test]
async fn catalog_applies_allowed_tools() {
let mcp = fake_mcp_server(two_tools(), json!({})).await;
let admin = MockServer::start().await;
mount_admin(
&admin,
admin_body_agentic(
"s1",
&mcp.uri(),
Some(vec!["get_issue"]),
vec!["agentic_worker"],
),
)
.await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source.catalog(&tenant()).await;
assert_eq!(catalog.len(), 1);
assert!(catalog.tool_entry("s1", "get_issue").is_some());
assert!(catalog.tool_entry("s1", "search_code").is_none());
}
#[tokio::test]
async fn unreachable_server_degrades() {
let admin = MockServer::start().await;
mount_admin(
&admin,
admin_body_agentic("s1", "http://127.0.0.1:1/", None, vec!["agentic_worker"]),
)
.await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source.catalog(&tenant()).await;
assert_eq!(catalog.len(), 0);
assert!(catalog.is_empty());
}
#[tokio::test]
async fn admin_unreachable_degrades_to_empty() {
let admin = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/api/v1/designer/tenant/me/mcp-servers"))
.respond_with(ResponseTemplate::new(500))
.mount(&admin)
.await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source.catalog(&tenant()).await;
assert!(catalog.is_empty());
}
#[tokio::test]
async fn ttl_cache_reuses_within_window() {
let mcp = fake_mcp_server(two_tools(), json!({})).await;
let admin = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/api/v1/designer/tenant/me/mcp-servers"))
.and(header("authorization", "Bearer gtc_live_x"))
.respond_with(ResponseTemplate::new(200).set_body_json(admin_body_agentic(
"s1",
&mcp.uri(),
None,
vec!["agentic_worker"],
)))
.expect(1) .mount(&admin)
.await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let t = tenant();
let first = source.catalog(&t).await;
let second = source.catalog(&t).await;
assert!(Arc::ptr_eq(&first, &second));
}
#[tokio::test]
async fn dispatch_route_calls_server_and_wraps() {
let mcp = fake_mcp_server(two_tools(), json!({ "structuredContent": { "ok": 1 } })).await;
let admin = MockServer::start().await;
mount_admin(
&admin,
admin_body_agentic(
"s1",
&mcp.uri(),
Some(vec!["get_issue"]),
vec!["agentic_worker"],
),
)
.await;
let source = McpToolSource::new(admin.uri(), "gtc_live_x");
let catalog = source.catalog(&tenant()).await;
let route = catalog.route("s1", "get_issue").expect("route present");
let out = dispatch_route(route, "{}").await;
assert_eq!(out, json!({ "ok": 1 }), "got: {out}");
assert!(!out.to_string().contains("\"error\""), "got: {out}");
let mcp_err = fake_mcp_server(
two_tools(),
json!({ "isError": true, "content": [{ "type": "text", "text": "boom" }] }),
)
.await;
let admin_err = MockServer::start().await;
mount_admin(
&admin_err,
admin_body_agentic(
"s1",
&mcp_err.uri(),
Some(vec!["get_issue"]),
vec!["agentic_worker"],
),
)
.await;
let source_err = McpToolSource::new(admin_err.uri(), "gtc_live_x");
let catalog_err = source_err.catalog(&TenantContext::new("acme", "stg")).await;
let route_err = catalog_err.route("s1", "get_issue").expect("route present");
let out_err = dispatch_route(route_err, "{}").await;
assert!(out_err.to_string().contains("error"), "got: {out_err}");
}