use std::sync::Arc;
use std::time::Duration;
use mcpmesh::backends::spawn::SpawnBackend;
use mcpmesh_net::PeerIdentity;
use mcpmesh_net::transport::NdjsonTransport;
use serde_json::json;
use tokio::io::{duplex, split};
use tokio::sync::Semaphore;
use tokio::time::timeout;
const MAX_FRAME: usize = 16 * 1024 * 1024;
const STUB: &str = env!("CARGO_BIN_EXE_echo_mcp_stub");
#[tokio::test]
async fn run_backend_pumps_frames_and_injects_identity() {
timeout(Duration::from_secs(30), async {
let (server_io, client_io) = duplex(64 * 1024);
let (sr, sw) = split(server_io);
let backend_transport = NdjsonTransport::new(sr, sw, MAX_FRAME);
let (cr, cw) = split(client_io);
let mut client = NdjsonTransport::new(cr, cw, MAX_FRAME);
let backend = SpawnBackend {
cmd: vec![STUB.to_string()],
concurrency: Arc::new(Semaphore::new(4)),
service: "test".into(),
audit: mcpmesh::audit::AuditSink::disabled(),
limiter: mcpmesh::limits::RateLimiter::unlimited_shared(),
env: Default::default(),
cwd: None,
};
let identity = Some(PeerIdentity {
endpoint: [0u8; 32].into(),
name: "bob".into(),
user_id: None,
groups: vec![],
});
let initialize = json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": "2025-11-25", "capabilities": {}}
});
let session = tokio::spawn(async move {
backend
.run_over(identity, initialize, backend_transport)
.await
});
let init_resp = client.recv_value().await.unwrap().unwrap();
assert_eq!(init_resp["id"], 1);
assert_eq!(init_resp["result"]["serverInfo"]["name"], "echo-stub");
client
.send_value(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"arguments": {"text": "hello mesh"}}
}))
.await
.unwrap();
let call_resp = client.recv_value().await.unwrap().unwrap();
assert_eq!(call_resp["id"], 2);
assert_eq!(call_resp["result"]["content"][0]["text"], "hello mesh");
assert_eq!(call_resp["result"]["peer_name"], "bob");
drop(client);
session
.await
.unwrap()
.expect("run_over returns Ok on transport EOF");
})
.await
.expect("spawn backend test timed out");
}
#[tokio::test]
async fn concurrency_cap_is_enforced() {
timeout(Duration::from_secs(30), async {
let sem = Arc::new(Semaphore::new(1));
let (server_io, client_io) = duplex(64 * 1024);
let (sr, sw) = split(server_io);
let backend_transport = NdjsonTransport::new(sr, sw, MAX_FRAME);
let (cr, cw) = split(client_io);
let mut client = NdjsonTransport::new(cr, cw, MAX_FRAME);
let backend1 = SpawnBackend {
cmd: vec![STUB.to_string()],
concurrency: sem.clone(),
service: "test".into(),
audit: mcpmesh::audit::AuditSink::disabled(),
limiter: mcpmesh::limits::RateLimiter::unlimited_shared(),
env: Default::default(),
cwd: None,
};
let init = json!({"jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {}});
let session1 =
tokio::spawn(async move { backend1.run_over(None, init, backend_transport).await });
let resp = client.recv_value().await.unwrap().unwrap();
assert_eq!(resp["result"]["serverInfo"]["name"], "echo-stub");
let (server_io2, client_io2) = duplex(64 * 1024);
let (sr2, sw2) = split(server_io2);
let backend_transport2 = NdjsonTransport::new(sr2, sw2, MAX_FRAME);
let (cr2, cw2) = split(client_io2);
let mut client2 = NdjsonTransport::new(cr2, cw2, MAX_FRAME);
let backend2 = SpawnBackend {
cmd: vec![STUB.to_string()],
concurrency: sem.clone(),
service: "test".into(),
audit: mcpmesh::audit::AuditSink::disabled(),
limiter: mcpmesh::limits::RateLimiter::unlimited_shared(),
env: Default::default(),
cwd: None,
};
backend2
.run_over(
None,
json!({"jsonrpc": "2.0", "id": 7, "method": "initialize", "params": {}}),
backend_transport2,
)
.await
.expect("cap refusal is handled cleanly (Ok), not a session error");
let refusal = client2
.recv_value()
.await
.unwrap()
.expect("a -32053 refusal frame, not EOF");
assert_eq!(refusal["error"]["code"], -32053);
assert_eq!(refusal["error"]["data"]["source"], "mcpmesh");
assert_eq!(refusal["id"], 7, "the refusal echoes the initialize id");
drop(client);
session1.await.unwrap().expect("first session returns Ok");
})
.await
.expect("concurrency test timed out");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_service_env_cannot_set_mcpmesh_home() {
let _ = timeout(Duration::from_secs(30), async {}).await;
for key in ["MCPMESH_HOME", "Mcpmesh_Home"] {
let (server_io, client_io) = duplex(64 * 1024);
let (sr, sw) = split(server_io);
let backend_transport = NdjsonTransport::new(sr, sw, MAX_FRAME);
let (cr, cw) = split(client_io);
let mut client = NdjsonTransport::new(cr, cw, MAX_FRAME);
let mut env = std::collections::BTreeMap::new();
env.insert(key.to_string(), "/tmp/attacker-controlled".to_string());
let backend = SpawnBackend {
cmd: vec![STUB.to_string()],
concurrency: Arc::new(Semaphore::new(4)),
service: "test".into(),
audit: mcpmesh::audit::AuditSink::disabled(),
limiter: mcpmesh::limits::RateLimiter::unlimited_shared(),
env,
cwd: None,
};
let identity = Some(PeerIdentity {
endpoint: [0u8; 32].into(),
name: "bob".into(),
user_id: None,
groups: vec![],
});
let initialize = json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": "2025-11-25", "capabilities": {}}
});
let session = tokio::spawn(async move {
backend
.run_over(identity, initialize, backend_transport)
.await
});
let _init = client.recv_value().await.unwrap().unwrap();
client
.send_value(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"arguments": {"text": "hi"}}
}))
.await
.unwrap();
let call = client.recv_value().await.unwrap().unwrap();
assert_eq!(
call["result"]["mesh_home"], "",
"a service env must not be able to set MCPMESH_HOME via `{key}` — it sandboxes all \
daemon state, and the identity injector never overwrites it"
);
drop(client);
let _ = session.await.unwrap();
}
}
#[tokio::test]
async fn run_backend_env_reaches_child_but_identity_is_not_spoofable() {
timeout(Duration::from_secs(30), async {
let (server_io, client_io) = duplex(64 * 1024);
let (sr, sw) = split(server_io);
let backend_transport = NdjsonTransport::new(sr, sw, MAX_FRAME);
let (cr, cw) = split(client_io);
let mut client = NdjsonTransport::new(cr, cw, MAX_FRAME);
let mut env = std::collections::BTreeMap::new();
env.insert(
"SPAWN_TEST_ENV".to_string(),
"reached-the-child".to_string(),
);
env.insert("MCPMESH_PEER_USER".to_string(), "b64u:FORGED".to_string());
env.insert("MCPMESH_PEER_EID".to_string(), "eid:FORGED".to_string());
let backend = SpawnBackend {
cmd: vec![STUB.to_string()],
concurrency: Arc::new(Semaphore::new(4)),
service: "test".into(),
audit: mcpmesh::audit::AuditSink::disabled(),
limiter: mcpmesh::limits::RateLimiter::unlimited_shared(),
env,
cwd: None,
};
let identity = Some(PeerIdentity {
endpoint: [0u8; 32].into(),
name: "bob".into(),
user_id: None, groups: vec![],
});
let initialize = json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {"protocolVersion": "2025-11-25", "capabilities": {}}
});
let session = tokio::spawn(async move {
backend
.run_over(identity, initialize, backend_transport)
.await
});
let _init = client.recv_value().await.unwrap().unwrap();
client
.send_value(json!({
"jsonrpc": "2.0", "id": 2, "method": "tools/call",
"params": {"arguments": {"text": "hi"}}
}))
.await
.unwrap();
let call = client.recv_value().await.unwrap().unwrap();
assert_eq!(
call["result"]["test_env"], "reached-the-child",
"per-service env must reach the child (#51)"
);
assert_eq!(
call["result"]["peer_user"], "",
"a service env must never spoof MCPMESH_PEER_USER (#51 security)"
);
assert_eq!(call["result"]["peer_name"], "bob", "real identity injected");
assert_eq!(
call["result"]["peer_eid"],
mcpmesh_net::EndpointId::from_bytes([0u8; 32]).principal(),
"a service env must never spoof MCPMESH_PEER_EID — it is the authz key"
);
drop(client);
let _ = session.await.unwrap();
})
.await
.expect("env test timed out");
}