From 8d4dcb5111d7aa6db2d597d55fe6a483235958c3 Mon Sep 17 00:00:00 2001 From: meraklbz Date: Sun, 16 Aug 2026 23:06:03 +0800 Subject: [PATCH] fix: reset SSE event type and id buffers after dispatch --- .../client/transport/ResponseSubscribers.java | 8 ++ .../transport/ResponseSubscribersTests.java | 76 +++++++++++++++++++ 2 files changed, 84 insertions(+) create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/ResponseSubscribersTests.java diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java index 29dc23c35..5044e8f75 100644 --- a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java +++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java @@ -149,6 +149,11 @@ protected void hookOnNext(String line) { this.sink.next(new SseResponseEvent(responseInfo, sseEvent)); this.eventBuilder.setLength(0); + // Per the WHATWG SSE spec, both the data buffer and the event type + // buffer are reset when an event is dispatched, so a stale event type + // must not leak into the following events. + this.currentEventId.set(null); + this.currentEventType.set(null); } } else { @@ -194,6 +199,9 @@ protected void hookOnComplete() { String eventData = this.eventBuilder.toString(); SseEvent sseEvent = new SseEvent(currentEventId.get(), currentEventType.get(), eventData.trim()); this.sink.next(new SseResponseEvent(responseInfo, sseEvent)); + this.eventBuilder.setLength(0); + this.currentEventId.set(null); + this.currentEventType.set(null); } this.sink.complete(); } diff --git a/mcp-core/src/test/java/io/modelcontextprotocol/client/transport/ResponseSubscribersTests.java b/mcp-core/src/test/java/io/modelcontextprotocol/client/transport/ResponseSubscribersTests.java new file mode 100644 index 000000000..f8ae07a3f --- /dev/null +++ b/mcp-core/src/test/java/io/modelcontextprotocol/client/transport/ResponseSubscribersTests.java @@ -0,0 +1,76 @@ +/* + * Copyright 2024-2026 the original author or authors. + */ + +package io.modelcontextprotocol.client.transport; + +import java.net.http.HttpResponse.ResponseInfo; +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.reactivestreams.Subscription; + +import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent; +import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseLineSubscriber; +import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseResponseEvent; +import reactor.core.publisher.Flux; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; + +/** + * Unit tests for {@link ResponseSubscribers.SseLineSubscriber} event buffer handling. + * + *

+ * Per the WHATWG + * HTML Living Standard ยง9.2.6, both the data buffer and the event type buffer must be + * reset when an event is dispatched, so that a stale event type or id does not leak into + * subsequent events. + */ +class ResponseSubscribersTests { + + @Test + void eventTypeIsResetAfterEventDispatch() { + List events = subscribeAndFeed("event: ping", "data: {\"jsonrpc\":\"2.0\",\"method\":\"ping\"}", + "", "data: {\"jsonrpc\":\"2.0\",\"method\":\"notifications/tools/list_changed\"}", ""); + + assertThat(events).hasSize(2); + assertThat(events.get(0)).isInstanceOfSatisfying(SseResponseEvent.class, + event -> assertThat(event.sseEvent().event()).as("the first event carries its explicit event type") + .isEqualTo("ping")); + assertThat(events.get(1)).isInstanceOfSatisfying(SseResponseEvent.class, + event -> assertThat(event.sseEvent().event()) + .as("the event type buffer must be reset after dispatch, so a bare data event has no stale type") + .isNull()); + } + + @Test + void eventIdIsResetAfterEventDispatch() { + List events = subscribeAndFeed("id: 1", "data: {\"jsonrpc\":\"2.0\",\"method\":\"ping\"}", "", + "data: {\"jsonrpc\":\"2.0\",\"method\":\"notifications/tools/list_changed\"}", ""); + + assertThat(events).hasSize(2); + assertThat(events.get(0)).isInstanceOfSatisfying(SseResponseEvent.class, + event -> assertThat(event.sseEvent().id()).as("the first event carries its explicit id").isEqualTo("1")); + assertThat(events.get(1)).isInstanceOfSatisfying(SseResponseEvent.class, + event -> assertThat(event.sseEvent().id()) + .as("the id buffer must be reset after dispatch, so a bare data event has no stale id") + .isNull()); + } + + private static List subscribeAndFeed(String... lines) { + ResponseInfo responseInfo = mock(ResponseInfo.class); + Subscription subscription = mock(Subscription.class); + + return Flux.create(sink -> { + SseLineSubscriber subscriber = new SseLineSubscriber(responseInfo, sink); + subscriber.hookOnSubscribe(subscription); + for (String line : lines) { + subscriber.hookOnNext(line); + } + subscriber.hookOnComplete(); + }).collectList().block(); + } + +}