Skip to content

Commit 081605d

Browse files
Fix stdio readiness barrier to wait for both Mono<Void> signals
Mono.zip cancels the slower Mono<Void> source when the faster one completes empty, so sendMessage could proceed before both inbound and outbound readiness signals fired. Mono.when waits for both. Fixes #303.
1 parent c7fef64 commit 081605d

2 files changed

Lines changed: 31 additions & 1 deletion

File tree

‎mcp-core/src/main/java/io/modelcontextprotocol/server/transport/StdioServerTransportProvider.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,7 @@ public StdioMcpSessionTransport() {
174174
@Override
175175
public Mono<Void> sendMessage(McpSchema.JSONRPCMessage message) {
176176

177-
return Mono.zip(inboundReady.asMono(), outboundReady.asMono()).then(Mono.defer(() -> {
177+
return Mono.when(inboundReady.asMono(), outboundReady.asMono()).then(Mono.defer(() -> {
178178
try {
179179
outboundSink.emitNext(message, Sinks.EmitFailureHandler.busyLooping(Duration.ofMillis(100)));
180180
return Mono.empty();

‎mcp-test/src/test/java/io/modelcontextprotocol/server/transport/StdioServerTransportProviderTests.java‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313
import java.io.InputStreamReader;
1414
import java.io.OutputStream;
1515
import java.io.PrintStream;
16+
import java.lang.reflect.Field;
1617
import java.nio.charset.StandardCharsets;
1718
import java.time.Duration;
1819
import java.util.Map;
@@ -30,6 +31,7 @@
3031
import org.junit.jupiter.api.Test;
3132
import reactor.core.publisher.Flux;
3233
import reactor.core.publisher.Mono;
34+
import reactor.core.publisher.Sinks;
3335
import reactor.core.scheduler.Schedulers;
3436
import reactor.test.StepVerifier;
3537

@@ -307,6 +309,34 @@ void shouldHandleSessionClose() {
307309
verify(mockSession).closeGracefully();
308310
}
309311

312+
@Test
313+
@SuppressWarnings("unchecked")
314+
void sendMessageWaitsForBothInboundAndOutboundReadinessSignals() throws Exception {
315+
transportProvider = new StdioServerTransportProvider(McpJsonDefaults.getMapper(), System.in,
316+
testOutPrintStream);
317+
Class<?> transportClass = Class
318+
.forName(StdioServerTransportProvider.class.getName() + "$StdioMcpSessionTransport");
319+
var constructor = transportClass.getDeclaredConstructor(StdioServerTransportProvider.class);
320+
constructor.setAccessible(true);
321+
McpServerTransport transport = (McpServerTransport) constructor.newInstance(transportProvider);
322+
323+
Field inboundReadyField = StdioServerTransportProvider.class.getDeclaredField("inboundReady");
324+
inboundReadyField.setAccessible(true);
325+
Sinks.One<Void> inboundReady = (Sinks.One<Void>) inboundReadyField.get(transportProvider);
326+
327+
Field outboundReadyField = transportClass.getDeclaredField("outboundReady");
328+
outboundReadyField.setAccessible(true);
329+
Sinks.One<Void> outboundReady = (Sinks.One<Void>) outboundReadyField.get(transport);
330+
331+
StepVerifier
332+
.create(transport.sendMessage(
333+
new McpSchema.JSONRPCNotification(McpSchema.JSONRPC_VERSION, "test/notification", Map.of())))
334+
.then(() -> inboundReady.tryEmitValue(null))
335+
.expectNoEvent(Duration.ofMillis(100))
336+
.then(() -> outboundReady.tryEmitValue(null))
337+
.verifyComplete();
338+
}
339+
310340
@Test
311341
void shouldHandleConcurrentSendMessage() throws Exception {
312342
int messageCount = 500;

0 commit comments

Comments
 (0)