Skip to content

Commit 45f93da

Browse files
committed
Fix SSE event parsing and error signalling in HttpClient transports
The SSE parser now follows the SSE spec: an event's type resets after each event and defaults to "message", and an empty id clears the last event id. Before, the type carried over from the previous event, so a message following an "endpoint" event was taken for another endpoint. Behaviour changes compared to main: - Errors caused by the server (non-2xx responses, unknown content types, malformed messages) are McpTransportException instead of plain RuntimeException. Error messages no longer include the response body. - A 400 on the streamable GET stream no longer invalidates the session. - A non-2xx SSE connect with an empty body now fails instead of hanging, and sendMessage() completes when the server closes an SSE response without sending any events. - Errors that happen after connect() or sendMessage() has completed now reach the exception handler and stop showing up as Reactor onErrorDropped logs. Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
1 parent 6a01d9b commit 45f93da

4 files changed

Lines changed: 120 additions & 58 deletions

File tree

‎mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransport.java‎

Lines changed: 22 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import java.net.http.HttpResponse;
1212
import java.time.Duration;
1313
import java.util.List;
14+
import java.util.concurrent.atomic.AtomicBoolean;
1415
import java.util.concurrent.atomic.AtomicReference;
1516
import java.util.function.Consumer;
1617
import java.util.function.Function;
@@ -388,6 +389,14 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
388389
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
389390
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
390391
}).flatMap(requestBuilder -> Mono.create(sink -> {
392+
// Once connect() has completed, a later failure can no longer be reported
393+
// through its sink: signalling it there would only have Reactor drop it.
394+
AtomicBoolean connected = new AtomicBoolean();
395+
Runnable markConnected = () -> {
396+
if (connected.compareAndSet(false, true)) {
397+
sink.success();
398+
}
399+
};
391400
Disposable connection = Mono
392401
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
393402
HttpResponse.BodyHandlers.ofPublisher()))
@@ -407,7 +416,7 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
407416
}
408417
else {
409418
return ResponseBodyHandlers.drainThenError(response.body(), this.maxResponseSize,
410-
new RuntimeException("Failed to connect to SSE stream: " + statusCode));
419+
new McpTransportException("Failed to connect to SSE stream: " + statusCode));
411420
}
412421
})
413422
.flatMap(sseEvent -> {
@@ -418,44 +427,44 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
418427
messageEndpointValidator.validate(uri, messageEndpointUri);
419428
}
420429
catch (InvalidSseMessageEndpointException e) {
421-
sink.error(e);
422430
this.messageEndpointSink.tryEmitError(e);
423431
return Flux.error(e);
424432
}
425433
if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) {
426-
sink.success();
434+
markConnected.run();
427435
return Flux.empty(); // No further processing needed
428436
}
429-
else {
430-
sink.error(new RuntimeException("Failed to handle SSE endpoint event"));
431-
}
437+
return Flux.error(new McpTransportException("Failed to handle SSE endpoint event"));
432438
}
433439
else if (MESSAGE_EVENT_TYPE.equals(sseEvent.event())) {
434440
String data = sseEvent.data();
435441
if (data == null || data.isBlank()) {
436442
logger.debug("Skipping SSE event with empty data (stream primer)");
437-
sink.success();
443+
markConnected.run();
438444
return Flux.<McpSchema.JSONRPCMessage>empty();
439445
}
440446
JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, data);
441-
sink.success();
447+
markConnected.run();
442448
return Flux.just(message);
443449
}
444450
else {
445451
logger.debug("Received unrecognized SSE event type: {}", sseEvent);
446-
sink.success();
452+
markConnected.run();
453+
return Flux.<McpSchema.JSONRPCMessage>empty();
447454
}
448455
}
449456
catch (IOException e) {
450-
sink.error(new McpTransportException("Error processing SSE event", e));
457+
return Flux.<McpSchema.JSONRPCMessage>error(
458+
new McpTransportException("Error processing SSE event", e));
451459
}
452-
return Flux.<McpSchema.JSONRPCMessage>empty();
453460
})
454461
.flatMap(jsonRpcMessage -> handler.apply(Mono.just(jsonRpcMessage)))
455462
.onErrorComplete(t -> {
456463
if (!isClosing) {
457464
logger.warn("SSE stream observed an error", t);
458-
sink.error(t);
465+
if (connected.compareAndSet(false, true)) {
466+
sink.error(t);
467+
}
459468
}
460469
return true;
461470
})
@@ -540,7 +549,7 @@ private Mono<Void> sendHttpPost(final String endpoint, final String body) {
540549
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize).then();
541550
}
542551
return ResponseBodyHandlers.decodeAggregateResponse(response.body(), this.maxResponseSize)
543-
.flatMap(text -> Mono.error(new RuntimeException(
552+
.flatMap(text -> Mono.error(new McpTransportException(
544553
"Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
545554
});
546555
});

‎mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java‎

Lines changed: 27 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -403,15 +403,11 @@ else if (statusCode == NOT_FOUND) {
403403
}
404404
}
405405
else if (statusCode == BAD_REQUEST) {
406-
// if the session id was set, but the session no longer
407-
// exists, some servers can return 400 instead of 404
408-
if (maybeSessionId.isPresent()) {
409-
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
410-
exception = new McpTransportSessionNotFoundException(sessionIdRepresentation);
411-
}
412-
else {
413-
exception = new McpTransportException("Bad Request. Status code:" + statusCode);
414-
}
406+
// Unlike a POST, a GET is not treated as a session-not-found
407+
// signal on 400: servers also reject the listening stream
408+
// itself with 400, which must not cost the client its
409+
// session.
410+
exception = new McpTransportException("Bad Request. Status code:" + statusCode);
415411
}
416412
else if (statusCode >= 200 && statusCode < 300) {
417413
String contentType = httpResponse.headers()
@@ -521,6 +517,15 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
521517

522518
final AtomicReference<Disposable> disposableRef = new AtomicReference<>();
523519

520+
// Once sendMessage() has completed, a later failure can no longer be reported
521+
// through its sink: signalling it there would only have Reactor drop it.
522+
final AtomicBoolean delivered = new AtomicBoolean();
523+
final Runnable markDelivered = () -> {
524+
if (delivered.compareAndSet(false, true)) {
525+
deliveredSink.success();
526+
}
527+
};
528+
524529
Disposable connection = Mono.deferContextual(ctx -> {
525530
HttpRequest.Builder requestBuilder = this.requestBuilder.copy();
526531

@@ -579,22 +584,17 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
579584

580585
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
581586
logger.debug("No body returned for POST in session {}", sessionRepresentation);
582-
deliveredSink.success();
587+
markDelivered.run();
583588
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
584589
}
585590
else if (contentType.contains(TEXT_EVENT_STREAM)) {
586-
AtomicBoolean delivered = new AtomicBoolean();
587-
return consumeSseStream(httpResponse.body(), null, () -> {
588-
if (delivered.compareAndSet(false, true)) {
589-
deliveredSink.success();
590-
}
591-
});
591+
return consumeSseStream(httpResponse.body(), null, markDelivered);
592592
}
593593
else if (contentType.contains(APPLICATION_JSON)) {
594594
return ResponseBodyHandlers
595595
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
596596
.flatMapMany(data -> {
597-
deliveredSink.success();
597+
markDelivered.run();
598598
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
599599
logger.warn("Notification: {} received non-compliant response: {}",
600600
sentMessage, Utils.hasText(data) ? data : "[empty]");
@@ -613,7 +613,7 @@ else if (contentType.contains(APPLICATION_JSON)) {
613613
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
614614
sessionRepresentation);
615615
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
616-
new RuntimeException("Unknown media type returned: " + contentType));
616+
new McpTransportException("Unknown media type returned: " + contentType));
617617
}
618618
else if (statusCode == NOT_FOUND) {
619619
if (maybeSessionId.isPresent()) {
@@ -640,7 +640,7 @@ else if (statusCode >= 400 && statusCode < 500) {
640640
}
641641

642642
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
643-
new RuntimeException("Failed to send message, status code: " + statusCode));
643+
new McpTransportException("Failed to send message, status code: " + statusCode));
644644
})
645645
.onErrorMap(CompletionException.class, Throwable::getCause))
646646
.retryWhen(authorizationErrorRetrySpec())
@@ -660,10 +660,16 @@ else if (statusCode >= 400 && statusCode < 500) {
660660
catch (Exception e) {
661661
logger.error("Error handling exception {}", t.getMessage(), e);
662662
}
663-
// inform the caller of sendMessage
664-
deliveredSink.error(t);
663+
// inform the caller of sendMessage, unless it has already completed
664+
if (delivered.compareAndSet(false, true)) {
665+
deliveredSink.error(t);
666+
}
665667
return true;
666668
})
669+
// An exchange can end without anything having signalled delivery, e.g.
670+
// an SSE response closed before its first event. The server accepted
671+
// the message all the same, so sendMessage() must not be left pending.
672+
.doOnComplete(markDelivered)
667673
.contextWrite(deliveredSink.contextView())
668674
.subscribe();
669675

‎mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java‎

Lines changed: 34 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,11 @@ class ResponseBodyHandlers {
4747
*/
4848
private static final int SSE_FRAMING_OVERHEAD = "event: ".length();
4949

50+
/**
51+
* The type of an SSE event that does not name one with an {@code event:} field.
52+
*/
53+
private static final String DEFAULT_EVENT_TYPE = "message";
54+
5055
record SseEvent(String id, String event, String data) {
5156
}
5257

@@ -368,11 +373,17 @@ private int indexOfLineTerminator(int from) {
368373

369374
/**
370375
* Stateful SSE line parser. Accumulates {@code data:}, {@code id:} and {@code event:}
371-
* fields until a blank line dispatches the event. Per the SSE spec, {@code id} and
372-
* {@code event} persist across events until re-set; {@code data} is reset after each
373-
* dispatch, and a blank line dispatches only when a {@code data:} field was seen,
374-
* whether or not it carried a value. Comments and fields the parser does not handle,
375-
* such as {@code retry:}, are ignored as the spec requires.
376+
* fields until a blank line dispatches the event. Per the SSE spec, {@code id} is the
377+
* last event ID and persists across events until re-set, with an empty value clearing
378+
* it; {@code event} and {@code data} are reset by every blank line, so an event that
379+
* does not name its type is a {@code message} event whatever preceded it. A blank
380+
* line dispatches only when a {@code data:} field was seen, whether or not it carried
381+
* a value. Comments and fields the parser does not handle, such as {@code retry:},
382+
* are ignored as the spec requires.
383+
*
384+
* @see <a href=
385+
* "https://html.spec.whatwg.org/multipage/server-sent-events.html#event-stream-interpretation">Interpreting
386+
* an event stream</a>
376387
*/
377388
static final class SseEventParser {
378389

@@ -399,12 +410,7 @@ static final class SseEventParser {
399410

400411
Optional<SseEvent> feed(String line) {
401412
if (line.isEmpty()) {
402-
if (data.length() == 0) {
403-
return Optional.empty();
404-
}
405-
SseEvent result = new SseEvent(id, event, data.toString().trim());
406-
data.setLength(0);
407-
return Optional.of(result);
413+
return flush();
408414
}
409415
if (line.startsWith("data:")) {
410416
// Every data field appends its value followed by a separator, so a
@@ -422,16 +428,16 @@ Optional<SseEvent> feed(String line) {
422428
data.append(value).append('\n');
423429
}
424430
else if (line.startsWith("id:")) {
425-
String rest = line.substring(3);
426-
if (!rest.isEmpty()) {
427-
id = rest.trim();
431+
String value = line.substring(3).trim();
432+
// The spec ignores an id carrying a NULL, and an empty id resets the last
433+
// event ID, which leaves nothing to resume from.
434+
if (value.indexOf('\0') == -1) {
435+
id = value.isEmpty() ? null : value;
428436
}
429437
}
430438
else if (line.startsWith("event:")) {
431-
String rest = line.substring(6);
432-
if (!rest.isEmpty()) {
433-
event = rest.trim();
434-
}
439+
String value = line.substring(6).trim();
440+
event = value.isEmpty() ? null : value;
435441
}
436442
else if (line.startsWith(":")) {
437443
logger.debug("Ignoring comment line: {}", line);
@@ -444,11 +450,19 @@ else if (line.startsWith(":")) {
444450
return Optional.empty();
445451
}
446452

453+
/**
454+
* Emits the pending event, if a {@code data:} field was seen, and resets the
455+
* per-event state. The event type is reset even when nothing is dispatched, as
456+
* the spec requires, while the id is the last event ID and so survives. An event
457+
* that did not name its type is emitted as a {@code message} event.
458+
*/
447459
Optional<SseEvent> flush() {
448-
if (data.length() == 0) {
460+
String type = this.event;
461+
this.event = null;
462+
if (data.isEmpty()) {
449463
return Optional.empty();
450464
}
451-
SseEvent result = new SseEvent(id, event, data.toString().trim());
465+
SseEvent result = new SseEvent(id, type != null ? type : DEFAULT_EVENT_TYPE, data.toString().trim());
452466
data.setLength(0);
453467
return Optional.of(result);
454468
}

‎mcp-core/src/test/java/io/modelcontextprotocol/client/transport/SseEventParserTests.java‎

Lines changed: 37 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ void simpleDataEvent() {
2323
assertThat(event).isPresent();
2424
assertThat(event.get().data()).isEqualTo("hello");
2525
assertThat(event.get().id()).isNull();
26-
assertThat(event.get().event()).isNull();
26+
assertThat(event.get().event()).isEqualTo("message");
2727
}
2828

2929
@Test
@@ -50,22 +50,55 @@ void idAndEventFieldsCaptured() {
5050
}
5151

5252
@Test
53-
void idAndEventPersistAcrossEvents() {
53+
void idPersistsAcrossEventsButEventTypeDoesNot() {
5454
SseEventParser p = new SseEventParser(Integer.MAX_VALUE);
5555
p.feed("id: 1");
56-
p.feed("event: message");
56+
p.feed("event: endpoint");
5757
p.feed("data: one");
5858
SseEvent first = p.feed("").orElseThrow();
5959
assertThat(first.id()).isEqualTo("1");
60-
assertThat(first.event()).isEqualTo("message");
60+
assertThat(first.event()).isEqualTo("endpoint");
6161

62+
// An event that does not name its type is a message event, not another
63+
// endpoint event
6264
p.feed("data: two");
6365
SseEvent second = p.feed("").orElseThrow();
6466
assertThat(second.id()).isEqualTo("1");
6567
assertThat(second.event()).isEqualTo("message");
6668
assertThat(second.data()).isEqualTo("two");
6769
}
6870

71+
@Test
72+
void blankLineWithNoDataStillResetsEventType() {
73+
SseEventParser p = new SseEventParser(Integer.MAX_VALUE);
74+
p.feed("event: endpoint");
75+
assertThat(p.feed("")).isEmpty();
76+
p.feed("data: payload");
77+
SseEvent event = p.feed("").orElseThrow();
78+
assertThat(event.event()).isEqualTo("message");
79+
}
80+
81+
@Test
82+
void emptyIdClearsLastEventId() {
83+
SseEventParser p = new SseEventParser(Integer.MAX_VALUE);
84+
p.feed("id: 1");
85+
p.feed("data: one");
86+
assertThat(p.feed("").orElseThrow().id()).isEqualTo("1");
87+
88+
p.feed("id:");
89+
p.feed("data: two");
90+
assertThat(p.feed("").orElseThrow().id()).isNull();
91+
}
92+
93+
@Test
94+
void idContainingNullIsIgnored() {
95+
SseEventParser p = new SseEventParser(Integer.MAX_VALUE);
96+
p.feed("id: 1");
97+
p.feed("id: 2\u0000" + "3");
98+
p.feed("data: payload");
99+
assertThat(p.feed("").orElseThrow().id()).isEqualTo("1");
100+
}
101+
69102
@Test
70103
void commentLineIgnored() {
71104
SseEventParser p = new SseEventParser(Integer.MAX_VALUE);

0 commit comments

Comments
 (0)