Skip to content

Commit b25d2cc

Browse files
committed
fix(core): remove pending response entries on request timeout and cancellation
McpClientSession and McpServerSession left the pendingResponses entry in place when the downstream timeout fired or the caller cancelled, so each timed-out request leaked its entry forever (request IDs are unique). Mirror the doOnError/doOnCancel cleanup the streamable session variant already performs after .timeout(requestTimeout). Fixes #1133
1 parent 4186ca1 commit b25d2cc

3 files changed

Lines changed: 27 additions & 2 deletions

File tree

‎mcp-core/src/main/java/io/modelcontextprotocol/spec/McpClientSession.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -264,7 +264,10 @@ public <T> Mono<T> sendRequest(String method, Object requestParams, TypeRef<T> t
264264
this.pendingResponses.remove(requestId);
265265
pendingResponseSink.error(error);
266266
});
267-
})).timeout(this.requestTimeout).handle((jsonRpcResponse, deliveredResponseSink) -> {
267+
})).timeout(this.requestTimeout)
268+
.doOnError(e -> this.pendingResponses.remove(requestId))
269+
.doOnCancel(() -> this.pendingResponses.remove(requestId))
270+
.handle((jsonRpcResponse, deliveredResponseSink) -> {
268271
if (jsonRpcResponse.error() != null) {
269272
logger.info("Server returned a JSON-RPC error when calling method {}: {}", method,
270273
jsonRpcResponse.error());

‎mcp-core/src/main/java/io/modelcontextprotocol/spec/McpServerSession.java‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -184,7 +184,10 @@ public <T> Mono<T> sendRequest(String method, Object requestParams, TypeRef<T> t
184184
this.pendingResponses.remove(requestId);
185185
sink.error(error);
186186
});
187-
}).timeout(requestTimeout).handle((jsonRpcResponse, sink) -> {
187+
}).timeout(requestTimeout)
188+
.doOnError(e -> this.pendingResponses.remove(requestId))
189+
.doOnCancel(() -> this.pendingResponses.remove(requestId))
190+
.handle((jsonRpcResponse, sink) -> {
188191
if (jsonRpcResponse.error() != null) {
189192
sink.error(new McpError(jsonRpcResponse.error()));
190193
}

‎mcp-core/src/test/java/io/modelcontextprotocol/spec/McpClientSessionTests.java‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -303,4 +303,23 @@ void testGracefulShutdown() {
303303
StepVerifier.create(session.closeGracefully()).verifyComplete();
304304
}
305305

306+
@Test
307+
void testRequestTimeoutRemovesPendingResponse() throws Exception {
308+
var transport = new MockMcpClientTransport();
309+
var session = new McpClientSession(Duration.ofMillis(50), transport, Map.of(), Map.of(),
310+
Function.identity());
311+
312+
Mono<String> responseMono = session.sendRequest(TEST_METHOD, "test", responseType);
313+
314+
StepVerifier.create(responseMono).expectError(java.util.concurrent.TimeoutException.class).verify();
315+
316+
var field = McpClientSession.class.getDeclaredField("pendingResponses");
317+
field.setAccessible(true);
318+
@SuppressWarnings("unchecked")
319+
var pendingResponses = (java.util.Map<Object, ?>) field.get(session);
320+
assertThat(pendingResponses).isEmpty();
321+
322+
session.close();
323+
}
324+
306325
}

0 commit comments

Comments
 (0)