Describe the bug
When an rmcp stdio server handles many concurrent requests that produce large responses, one or more JSON-RPC responses can be missing from stdout. The server keeps running and most responses are received, but at least one expected response ID is never observed by the client before the read deadline.
The reproducer below uses tools/call because it is an easy way to generate controlled large responses. The interesting part is not tool execution itself; it is the generic request/response path under concurrent response pressure.
To Reproduce
The test below starts a helper process running a minimal rmcp stdio server. The parent test acts as a raw JSON-RPC client:
- Start the helper server over stdin/stdout.
- Send
initialize and wait for its response. - Send 200
tools/call requests back-to-back without waiting between them. - Each request returns a 64 KiB text response.
- Read stdout lines and collect JSON-RPC response IDs.
- Fail if any request ID is still missing after the deadline.
A failing run looks like this conceptually:
missing response ids: {1173}
Usually most responses are observed. The failure is that one or more expected IDs never arrive.
Expected behavior
Every accepted JSON-RPC request with an id should eventually produce exactly one response for that id, unless the transport closes or an explicit error is returned.
Actual behavior
Under concurrent large responses, at least one response can be absent from stdout. The client does not need to observe malformed JSON or an unexpected ID; the response for a valid pending request is simply missing.
Note that this was still reproducible with a 60s read deadline, so it does not appear to be only a short-timeout issue.
Minimal reproducer
Add the following test as crates/rmcp/tests/test_stdio_response_concurrency.rs:
#![cfg(not(feature = "local"))]use std::{collections::BTreeSet, process::Stdio, time::Duration};use rmcp::{ErrorDataasMcpError,ServerHandler,ServiceExt,
model::{CallToolRequestParams,CallToolResult,ContentBlock,ServerCapabilities,ServerInfo},};use serde_json::{Value, json};use tokio::{
io::{AsyncBufReadExt,AsyncWrite,AsyncWriteExt,BufReader},
process::{Child,Command},};constHELPER_ENV:&str = "RMCP_STDIO_RESPONSE_CONCURRENCY_HELPER";constREQUESTS:usize = 200;constRESPONSE_BYTES:usize = 64*1024;constREAD_TIMEOUT:Duration = Duration::from_secs(10);#[tokio::test(flavor = "multi_thread", worker_threads = 8)]asyncfnraw_client_concurrent_large_stdio_tool_responses_are_not_lost() -> anyhow::Result<()>{// Spawn the same test binary as a child process so the server uses real// stdio pipes, not an in-process transport.letmut child = spawn_helper();letmut writer = child.stdin.take().expect("helper stdin");let stdout = child.stdout.take().expect("helper stdout");letmut reader = BufReader::new(stdout);// Complete the normal MCP initialization flow before stressing tools/call.// This keeps the repro focused on response delivery after initialization.send_json(&mut writer,&json!({"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2024-11-05","capabilities":{},"clientInfo":{"name":"raw-test-client","version":"0.0.0"}}})).await?;read_response_for_id(&mut reader,1).await?;send_json(&mut writer,&json!({"jsonrpc":"2.0","method":"notifications/initialized"})).await?;// Send the whole batch before reading responses. This creates concurrent// request handling and concurrent response production inside rmcp.for id inrequest_ids(){send_json(&mut writer,&json!({"jsonrpc":"2.0","id": id,"method":"tools/call","params":{"name":"large-response","arguments":{}}})).await?;}// Accept responses in any order. The assertion only cares that every// request ID eventually appears on stdout.let missing_ids = read_responses_for_ids(&mut reader,request_ids(),READ_TIMEOUT).await?;assert!(
missing_ids.is_empty(),"missing response ids: {missing_ids:?}",);drop(writer);wait_for_child(&mut child).await;Ok(())}structLargeResponseServer;implServerHandlerforLargeResponseServer{fnget_info(&self) -> ServerInfo{ServerInfo::new(ServerCapabilities::builder().enable_tools().build())}asyncfncall_tool(&self,request:CallToolRequestParams,_context: rmcp::service::RequestContext<rmcp::RoleServer>,) -> Result<CallToolResult,McpError>{assert_eq!("large-response", request.name.as_ref());// Large responses make stdout backpressure and response scheduling// visible with a small number of concurrent requests.Ok(CallToolResult::success(vec![ContentBlock::text("x".repeat(RESPONSE_BYTES),)]))}}#[tokio::test]asyncfnstdio_response_concurrency_helper() -> anyhow::Result<()>{// This is not an independent assertion. The parent test above starts this// same test binary with HELPER_ENV=1 so it can act as a small MCP server// connected over real stdin/stdout pipes.if std::env::var(HELPER_ENV).as_deref() != Ok("1"){returnOk(());}let server = LargeResponseServer.serve(rmcp::transport::stdio()).await?;
server.waiting().await?;Ok(())}fnspawn_helper() -> Child{let exe = std::env::current_exe().expect("current test exe");Command::new(exe).arg("--exact").arg("stdio_response_concurrency_helper").arg("--quiet").arg("--no-capture").arg("--test-threads").arg("1").env(HELPER_ENV,"1").stdin(Stdio::piped()).stdout(Stdio::piped()).stderr(Stdio::null()).kill_on_drop(true).spawn().expect("spawn helper")}asyncfnwait_for_child(child:&mutChild){let _ = tokio::time::timeout(Duration::from_secs(2), child.wait()).await;if child.id().is_some(){let _ = child.kill().await;}}fnrequest_ids() -> BTreeSet<u64>{(1000..1000 + REQUESTSasu64).collect()}asyncfnsend_json<W>(writer:&mutW,message:&Value) -> anyhow::Result<()>whereW:AsyncWrite + Unpin,{let serialized = serde_json::to_string(message)?;
writer.write_all(serialized.as_bytes()).await?;
writer.write_all(b"\n").await?;
writer.flush().await?;Ok(())}asyncfnread_response_for_id<R>(reader:&mutBufReader<R>,expected_id:u64) -> anyhow::Result<()>whereR: tokio::io::AsyncRead + Unpin,{let missing =
read_responses_for_ids(reader,BTreeSet::from([expected_id]),READ_TIMEOUT).await?;if missing.is_empty(){Ok(())}else{
anyhow::bail!("missing response id {expected_id}")}}asyncfnread_responses_for_ids<R>(reader:&mutBufReader<R>,mutpending_ids:BTreeSet<u64>,timeout:Duration,) -> anyhow::Result<BTreeSet<u64>>whereR: tokio::io::AsyncRead + Unpin,{let deadline = tokio::time::Instant::now() + timeout;while !pending_ids.is_empty(){let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());if remaining.is_zero(){break;}letmut line = String::new();letOk(read_result) = tokio::time::timeout(remaining, reader.read_line(&mut line)).awaitelse{break;};let read = read_result?;if read == 0{break;}let trimmed = line.trim();if trimmed.is_empty(){continue;}letOk(message) = serde_json::from_str::<Value>(trimmed)else{// Skip non-JSON lines coming from the test harness (e.g. "running 1 test", ...)continue;};ifletSome(id) = message.get("id").and_then(Value::as_u64){
pending_ids.remove(&id);}}Ok(pending_ids)}I could observe test failures on both debug and release builds:
cargo test -p rmcp --features transport-io --test test_stdio_response_concurrency raw_client_concurrent_large_stdio_tool_responses_are_not_lost -- --exact --nocapture
cargo test --release -p rmcp --features transport-io --test test_stdio_response_concurrency raw_client_concurrent_large_stdio_tool_responses_are_not_lost -- --exact --nocapture
If it does not reproduce immediately, increase REQUESTS or run either command repeatedly.
Thanks in advance. Happy to provide any additional information if needed... 😄
Describe the bug
When an rmcp stdio server handles many concurrent requests that produce large responses, one or more JSON-RPC responses can be missing from stdout. The server keeps running and most responses are received, but at least one expected response ID is never observed by the client before the read deadline.
The reproducer below uses
tools/callbecause it is an easy way to generate controlled large responses. The interesting part is not tool execution itself; it is the generic request/response path under concurrent response pressure.To Reproduce
The test below starts a helper process running a minimal rmcp stdio server. The parent test acts as a raw JSON-RPC client:
initializeand wait for its response.tools/callrequests back-to-back without waiting between them.A failing run looks like this conceptually:
Usually most responses are observed. The failure is that one or more expected IDs never arrive.
Expected behavior
Every accepted JSON-RPC request with an
idshould eventually produce exactly one response for thatid, unless the transport closes or an explicit error is returned.Actual behavior
Under concurrent large responses, at least one response can be absent from stdout. The client does not need to observe malformed JSON or an unexpected ID; the response for a valid pending request is simply missing.
Note that this was still reproducible with a 60s read deadline, so it does not appear to be only a short-timeout issue.
Minimal reproducer
Add the following test as
crates/rmcp/tests/test_stdio_response_concurrency.rs:I could observe test failures on both debug and release builds:
cargo test -p rmcp --features transport-io --test test_stdio_response_concurrency raw_client_concurrent_large_stdio_tool_responses_are_not_lost -- --exact --nocapturecargo test --release -p rmcp --features transport-io --test test_stdio_response_concurrency raw_client_concurrent_large_stdio_tool_responses_are_not_lost -- --exact --nocaptureIf it does not reproduce immediately, increase
REQUESTSor run either command repeatedly.Thanks in advance. Happy to provide any additional information if needed... 😄