Skip to content

Commit 3ccb4ec

Browse files
committed
fix(spec): tolerate unregistered listening stream for server-initiated messages
Right after initialize, the server may push a server-initiated request (e.g. roots/list when the client declared the roots capability) before the client has opened its GET /mcp stream. The session previously delegated to a MissingMcpTransportSession that failed immediately with Stream unavailable for session <id>, surfacing as an intermittent -32603/500 on the clients POST (seen in CI HttpServletStreamableIntegrationTests.testRootsNotificationWithEmptyRootsList). Changes: introduce package-private MissingListeningStreamException (subclass of IllegalStateException, same message) thrown by MissingMcpTransportSession; McpStreamableServerSession.sendRequest/sendNotification retry missing-stream failures with exponential backoff capped at 12 attempts (~5s); accept(notification) falls back to the retrying session facade when no stream is registered yet so handler-initiated pushes are covered too. Related: #952, #1061
1 parent a7bfddc commit 3ccb4ec

4 files changed

Lines changed: 260 additions & 10 deletions

File tree

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

Lines changed: 42 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import reactor.core.publisher.Flux;
2828
import reactor.core.publisher.Mono;
2929
import reactor.core.publisher.MonoSink;
30+
import reactor.util.retry.Retry;
3031

3132
/**
3233
* Representation of a Streamable HTTP server session that keeps track of mapping
@@ -41,6 +42,24 @@ public class McpStreamableServerSession implements McpLoggableSession {
4142

4243
private static final Logger logger = LoggerFactory.getLogger(McpStreamableServerSession.class);
4344

45+
/**
46+
* Delay before the first retry of a server-initiated message whose listening stream
47+
* is not registered yet.
48+
*/
49+
private static final Duration MISSING_STREAM_RETRY_MIN_BACKOFF = Duration.ofMillis(50);
50+
51+
/**
52+
* Upper bound for the exponential backoff between retries.
53+
*/
54+
private static final Duration MISSING_STREAM_RETRY_MAX_BACKOFF = Duration.ofMillis(500);
55+
56+
/**
57+
* Number of retries tolerated for a server-initiated message whose listening stream
58+
* is not registered yet. With the backoff above this yields roughly a 5s grace
59+
* window, after which the last failure propagates to the caller.
60+
*/
61+
private static final int MISSING_STREAM_RETRY_MAX_ATTEMPTS = 12;
62+
4463
private final ConcurrentHashMap<Object, McpStreamableServerSessionStream> requestIdToStream = new ConcurrentHashMap<>();
4564

4665
private final String id;
@@ -156,18 +175,26 @@ private String generateRequestId() {
156175

157176
@Override
158177
public <T> Mono<T> sendRequest(String method, Object requestParams, TypeRef<T> typeRef) {
159-
return Mono.defer(() -> {
160-
McpLoggableSession listeningStream = this.listeningStreamRef.get();
161-
return listeningStream.sendRequest(method, requestParams, typeRef);
162-
});
178+
return Mono.defer(() -> this.listeningStreamRef.get().sendRequest(method, requestParams, typeRef))
179+
.retryWhen(missingStreamRetry());
163180
}
164181

165182
@Override
166183
public Mono<Void> sendNotification(String method, Object params) {
167-
return Mono.defer(() -> {
168-
McpLoggableSession listeningStream = this.listeningStreamRef.get();
169-
return listeningStream.sendNotification(method, params);
170-
});
184+
return Mono.defer(() -> this.listeningStreamRef.get().sendNotification(method, params))
185+
.retryWhen(missingStreamRetry());
186+
}
187+
188+
/**
189+
* Retry specification tolerating a not-yet-registered listening stream: retries are
190+
* bounded by {@link #MISSING_STREAM_GRACE_PERIOD}, after which the last failure
191+
* propagates to the caller. Any other failure type propagates immediately.
192+
*/
193+
private Retry missingStreamRetry() {
194+
return Retry.backoff(MISSING_STREAM_RETRY_MAX_ATTEMPTS, MISSING_STREAM_RETRY_MIN_BACKOFF)
195+
.maxBackoff(MISSING_STREAM_RETRY_MAX_BACKOFF)
196+
.filter(MissingListeningStreamException.class::isInstance)
197+
.transientErrors(true);
171198
}
172199

173200
public Mono<Void> delete() {
@@ -254,6 +281,13 @@ public Mono<Void> accept(McpSchema.JSONRPCNotification notification) {
254281
return Mono.empty();
255282
}
256283
McpLoggableSession listeningStream = this.listeningStreamRef.get();
284+
if (listeningStream == this.missingMcpTransportSession) {
285+
// The listening stream may not be registered yet (e.g. the client is
286+
// still opening GET /mcp right after initialize). Delegate to this
287+
// session so that server-initiated requests triggered by the handler
288+
// retry until the stream shows up instead of failing immediately.
289+
listeningStream = this;
290+
}
257291
return notificationHandler.handle(new McpAsyncServerExchange(this.id, listeningStream,
258292
this.clientCapabilities.get(), this.clientInfo.get(), transportContext, this.jsonSchemaValidator),
259293
notification.params());
Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
1+
/*
2+
* Copyright 2024-2025 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.spec;
6+
7+
/**
8+
* Signals that a session has no listening stream registered to carry server-initiated
9+
* requests and notifications to the client. Extends {@link IllegalStateException} for
10+
* backwards compatibility with existing handling of the previously thrown plain
11+
* instances.
12+
*/
13+
@SuppressWarnings("serial")
14+
class MissingListeningStreamException extends IllegalStateException {
15+
16+
MissingListeningStreamException(String sessionId) {
17+
super("Stream unavailable for session " + sessionId);
18+
}
19+
20+
}

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,12 +32,12 @@ public MissingMcpTransportSession(String sessionId) {
3232

3333
@Override
3434
public <T> Mono<T> sendRequest(String method, Object requestParams, TypeRef<T> typeRef) {
35-
return Mono.error(new IllegalStateException("Stream unavailable for session " + this.sessionId));
35+
return Mono.error(new MissingListeningStreamException(this.sessionId));
3636
}
3737

3838
@Override
3939
public Mono<Void> sendNotification(String method, Object params) {
40-
return Mono.error(new IllegalStateException("Stream unavailable for session " + this.sessionId));
40+
return Mono.error(new MissingListeningStreamException(this.sessionId));
4141
}
4242

4343
@Override
Lines changed: 196 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,196 @@
1+
/*
2+
* Copyright 2024-2025 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.spec;
6+
7+
import java.time.Duration;
8+
import java.util.List;
9+
import java.util.Map;
10+
import java.util.concurrent.CountDownLatch;
11+
import java.util.concurrent.CopyOnWriteArrayList;
12+
import java.util.concurrent.TimeUnit;
13+
import java.util.concurrent.atomic.AtomicReference;
14+
15+
import io.modelcontextprotocol.json.TypeRef;
16+
import io.modelcontextprotocol.server.McpNotificationHandler;
17+
import io.modelcontextprotocol.server.McpRequestHandler;
18+
import io.modelcontextprotocol.spec.McpSchema.JSONRPCMessage;
19+
import io.modelcontextprotocol.spec.McpSchema.JSONRPCRequest;
20+
import io.modelcontextprotocol.spec.McpSchema.JSONRPCResponse;
21+
import org.junit.jupiter.api.Test;
22+
import reactor.core.publisher.Mono;
23+
import reactor.test.StepVerifier;
24+
25+
import static org.assertj.core.api.Assertions.assertThat;
26+
27+
/**
28+
* Verifies that server-initiated requests and notifications tolerate a not-yet registered
29+
* listening stream: they are retried until the client opens the GET /mcp stream and fail
30+
* after the grace period if it never appears. This is the race behind intermittent
31+
* {@code Stream unavailable for session} failures right after {@code initialize}.
32+
*/
33+
class McpStreamableServerSessionMissingStreamTests {
34+
35+
private static final TypeRef<Map<String, Object>> MAP_TYPE = new TypeRef<Map<String, Object>>() {
36+
};
37+
38+
private final RecordingTransport transport = new RecordingTransport();
39+
40+
@Test
41+
void notificationIsRetriedUntilListeningStreamRegisters() throws Exception {
42+
var session = newSession(Map.of(), Map.of());
43+
var delivered = new CountDownLatch(1);
44+
45+
var disposable = session.sendNotification("notifications/test", Map.of()).subscribe(v -> {
46+
}, e -> delivered.countDown(), delivered::countDown);
47+
try {
48+
assertThat(this.transport.sent).isEmpty();
49+
assertThat(delivered.getCount()).isEqualTo(1);
50+
51+
session.listeningStream(this.transport);
52+
assertThat(delivered.await(2, TimeUnit.SECONDS)).isTrue();
53+
}
54+
finally {
55+
disposable.dispose();
56+
}
57+
58+
assertThat(this.transport.sent).hasSize(1);
59+
}
60+
61+
@Test
62+
void requestIsRetriedAndResponseCorrelatedAfterRegistration() throws Exception {
63+
var session = newSession(Map.of(), Map.of());
64+
var response = new AtomicReference<Map<String, Object>>();
65+
var error = new AtomicReference<Throwable>();
66+
var completed = new CountDownLatch(1);
67+
68+
var disposable = session.sendRequest("sampling/createMessage", Map.of("messages", List.of()), MAP_TYPE)
69+
.subscribe(response::set, e -> {
70+
error.set(e);
71+
completed.countDown();
72+
}, completed::countDown);
73+
try {
74+
assertThat(this.transport.sent).isEmpty();
75+
session.listeningStream(this.transport);
76+
assertThat(this.transport.firstMessage.await(2, TimeUnit.SECONDS)).isTrue();
77+
78+
var request = (JSONRPCRequest) this.transport.sent.get(0);
79+
session.accept(JSONRPCResponse.result(request.id(), Map.of("stopReason", "endTurn")))
80+
.block(Duration.ofSeconds(1));
81+
82+
assertThat(completed.await(2, TimeUnit.SECONDS)).isTrue();
83+
if (error.get() != null) {
84+
throw new AssertionError("Request did not complete successfully", error.get());
85+
}
86+
}
87+
finally {
88+
disposable.dispose();
89+
}
90+
91+
assertThat(response.get()).containsEntry("stopReason", "endTurn");
92+
}
93+
94+
@Test
95+
void sendRequestFailsAfterGracePeriodWithoutListeningStream() {
96+
var session = newSession(Map.of(), Map.of());
97+
98+
StepVerifier.withVirtualTime(() -> session.sendRequest("sampling/createMessage", Map.of(), MAP_TYPE))
99+
.thenAwait(Duration.ofSeconds(6))
100+
.expectErrorSatisfies(e -> assertThat(e).hasRootCauseMessage("Stream unavailable for session test-session"))
101+
.verify(Duration.ofSeconds(2));
102+
}
103+
104+
@Test
105+
void notificationHandlerCanPushRequestBeforeStreamRegistered() throws Exception {
106+
var delivered = new AtomicReference<McpSchema.ListRootsResult>();
107+
McpNotificationHandler initializedHandler = (exchange,
108+
params) -> exchange.listRoots().doOnNext(delivered::set).then();
109+
110+
var session = newSession(Map.of(), Map.of("notifications/initialized", initializedHandler));
111+
var handled = new CountDownLatch(1);
112+
var error = new AtomicReference<Throwable>();
113+
114+
var disposable = session.accept(new McpSchema.JSONRPCNotification("notifications/initialized", Map.of()))
115+
.subscribe(v -> {
116+
}, e -> {
117+
error.set(e);
118+
handled.countDown();
119+
}, handled::countDown);
120+
try {
121+
assertThat(this.transport.sent).isEmpty();
122+
session.listeningStream(this.transport);
123+
assertThat(this.transport.firstMessage.await(2, TimeUnit.SECONDS)).isTrue();
124+
125+
var request = this.transport.sent.stream()
126+
.filter(JSONRPCRequest.class::isInstance)
127+
.map(JSONRPCRequest.class::cast)
128+
.findFirst()
129+
.orElseThrow();
130+
session
131+
.accept(JSONRPCResponse.result(request.id(),
132+
Map.of("roots", List.of(Map.of("name", "workspace", "uri", "file:///ws")))))
133+
.block(Duration.ofSeconds(1));
134+
135+
assertThat(handled.await(2, TimeUnit.SECONDS)).isTrue();
136+
if (error.get() != null) {
137+
throw new AssertionError("Handler did not complete successfully", error.get());
138+
}
139+
}
140+
finally {
141+
disposable.dispose();
142+
}
143+
144+
assertThat(delivered.get().roots()).singleElement().satisfies(root -> {
145+
assertThat(root.name()).isEqualTo("workspace");
146+
assertThat(root.uri()).isEqualTo("file:///ws");
147+
});
148+
}
149+
150+
private McpStreamableServerSession newSession(Map<String, McpRequestHandler<?>> requests,
151+
Map<String, McpNotificationHandler> notifications) {
152+
return new McpStreamableServerSession("test-session", McpSchema.ClientCapabilities.builder().build(),
153+
new McpSchema.Implementation("test-client", "1.0"), Duration.ofSeconds(10), requests, notifications);
154+
}
155+
156+
private static final class RecordingTransport implements McpStreamableServerTransport {
157+
158+
final CopyOnWriteArrayList<JSONRPCMessage> sent = new CopyOnWriteArrayList<>();
159+
160+
final CountDownLatch firstMessage = new CountDownLatch(1);
161+
162+
@Override
163+
public Mono<Void> sendMessage(JSONRPCMessage message, String messageId) {
164+
this.sent.add(message);
165+
this.firstMessage.countDown();
166+
return Mono.empty();
167+
}
168+
169+
@Override
170+
public Mono<Void> sendMessage(JSONRPCMessage message) {
171+
return sendMessage(message, "ignored");
172+
}
173+
174+
@Override
175+
public Mono<Void> closeGracefully() {
176+
return Mono.empty();
177+
}
178+
179+
@Override
180+
@SuppressWarnings("unchecked")
181+
public <T> T unmarshalFrom(Object data, TypeRef<T> typeRef) {
182+
// Minimal hand-rolled conversion: mcp-core tests run without a JSON binding
183+
// module on the classpath, so McpJsonDefaults is unavailable here.
184+
if (typeRef.getType() == McpSchema.ListRootsResult.class) {
185+
var map = (Map<String, Object>) data;
186+
var roots = ((List<Map<String, Object>>) map.get("roots")).stream()
187+
.map(root -> new McpSchema.Root((String) root.get("uri"), (String) root.get("name"), null))
188+
.toList();
189+
return (T) new McpSchema.ListRootsResult(roots, (String) map.get("nextCursor"), null);
190+
}
191+
return (T) data;
192+
}
193+
194+
}
195+
196+
}

0 commit comments

Comments
 (0)