Concurrent stdio responses can go missing under load #941

Description

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:

  1. Start the helper server over stdin/stdout.
  2. Send initialize and wait for its response.
  3. Send 200 tools/call requests back-to-back without waiting between them.
  4. Each request returns a 64 KiB text response.
  5. Read stdout lines and collect JSON-RPC response IDs.
  6. 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... 😄

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething is not working

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions

      , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
       blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
      }
      } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
      })();
      (function(){
      try {
      var __m = "github.com";
      var __re = new RegExp('^' + "github\\.com" + '
      
      Skip to content

      Concurrent stdio responses can go missing under load #941

      Description

      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:

      1. Start the helper server over stdin/stdout.
      2. Send initialize and wait for its response.
      3. Send 200 tools/call requests back-to-back without waiting between them.
      4. Each request returns a 64 KiB text response.
      5. Read stdout lines and collect JSON-RPC response IDs.
      6. 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... 😄

      Metadata

      Metadata

      Assignees

      No one assigned

        Labels

        bugSomething is not working

        Type

        No type

        Projects

        No projects

          Milestone

          No milestone

          Relationships

          None yet

          Development

          No branches or pull requests

          Issue actions

          , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
          Skip to content

          Concurrent stdio responses can go missing under load #941

          Description

          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:

          1. Start the helper server over stdin/stdout.
          2. Send initialize and wait for its response.
          3. Send 200 tools/call requests back-to-back without waiting between them.
          4. Each request returns a 64 KiB text response.
          5. Read stdout lines and collect JSON-RPC response IDs.
          6. 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... 😄

          Metadata

          Metadata

          Assignees

          No one assigned

            Labels

            bugSomething is not working

            Type

            No type

            Projects

            No projects

              Milestone

              No milestone

              Relationships

              None yet

              Development

              No branches or pull requests

              Issue actions

              , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
              Skip to content

              Concurrent stdio responses can go missing under load #941

              Description

              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:

              1. Start the helper server over stdin/stdout.
              2. Send initialize and wait for its response.
              3. Send 200 tools/call requests back-to-back without waiting between them.
              4. Each request returns a 64 KiB text response.
              5. Read stdout lines and collect JSON-RPC response IDs.
              6. 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... 😄

              Metadata

              Metadata

              Assignees

              No one assigned

                Labels

                bugSomething is not working

                Type

                No type

                Projects

                No projects

                  Milestone

                  No milestone

                  Relationships

                  None yet

                  Development

                  No branches or pull requests

                  Issue actions

                  , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
                  Skip to content

                  Concurrent stdio responses can go missing under load #941

                  Description

                  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:

                  1. Start the helper server over stdin/stdout.
                  2. Send initialize and wait for its response.
                  3. Send 200 tools/call requests back-to-back without waiting between them.
                  4. Each request returns a 64 KiB text response.
                  5. Read stdout lines and collect JSON-RPC response IDs.
                  6. 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... 😄

                  Metadata

                  Metadata

                  Assignees

                  No one assigned

                    Labels

                    bugSomething is not working

                    Type

                    No type

                    Projects

                    No projects

                      Milestone

                      No milestone

                      Relationships

                      None yet

                      Development

                      No branches or pull requests

                      Issue actions

                      , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
                      Skip to content

                      Concurrent stdio responses can go missing under load #941

                      Description

                      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:

                      1. Start the helper server over stdin/stdout.
                      2. Send initialize and wait for its response.
                      3. Send 200 tools/call requests back-to-back without waiting between them.
                      4. Each request returns a 64 KiB text response.
                      5. Read stdout lines and collect JSON-RPC response IDs.
                      6. 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... 😄

                      Metadata

                      Metadata

                      Assignees

                      No one assigned

                        Labels

                        bugSomething is not working

                        Type

                        No type

                        Projects

                        No projects

                          Milestone

                          No milestone

                          Relationships

                          None yet

                          Development

                          No branches or pull requests

                          Issue actions

                          , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
                          Skip to content

                          Concurrent stdio responses can go missing under load #941

                          Description

                          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:

                          1. Start the helper server over stdin/stdout.
                          2. Send initialize and wait for its response.
                          3. Send 200 tools/call requests back-to-back without waiting between them.
                          4. Each request returns a 64 KiB text response.
                          5. Read stdout lines and collect JSON-RPC response IDs.
                          6. 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... 😄

                          Metadata

                          Metadata

                          Assignees

                          No one assigned

                            Labels

                            bugSomething is not working

                            Type

                            No type

                            Projects

                            No projects

                              Milestone

                              No milestone

                              Relationships

                              None yet

                              Development

                              No branches or pull requests

                              Issue actions

                              , 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
                              Skip to content

                              Concurrent stdio responses can go missing under load #941

                              Description

                              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:

                              1. Start the helper server over stdin/stdout.
                              2. Send initialize and wait for its response.
                              3. Send 200 tools/call requests back-to-back without waiting between them.
                              4. Each request returns a 64 KiB text response.
                              5. Read stdout lines and collect JSON-RPC response IDs.
                              6. 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... 😄

                              Metadata

                              Metadata

                              Assignees

                              No one assigned

                                Labels

                                bugSomething is not working

                                Type

                                No type

                                Projects

                                No projects

                                  Milestone

                                  No milestone

                                  Relationships

                                  None yet

                                  Development

                                  No branches or pull requests

                                  Issue actions