From 7f383aa4336fe4c32981af8f6d6980e0d890ebcd Mon Sep 17 00:00:00 2001 From: Brennan Butler <64561607+brennanbutler01@users.noreply.github.com> Date: Fri, 18 Sep 2026 16:52:51 +0700 Subject: [PATCH 1/2] fix: stop dispatching messages after close Check connection state before dispatching each parsed event so closing from a listener suppresses the remaining messages in the same chunk. Add real-server regressions for ordinary and named events, open-connection controls, and a patch changeset. --- .changeset/quiet-stream-close.md | 5 ++++ src/EventSource.ts | 5 ++++ test/client.test.ts | 50 ++++++++++++++++++++++++++++++++ test/helpers/server.ts | 8 +++++ 4 files changed, 68 insertions(+) create mode 100644 .changeset/quiet-stream-close.md diff --git a/.changeset/quiet-stream-close.md b/.changeset/quiet-stream-close.md new file mode 100644 index 0000000..3b7ae06 --- /dev/null +++ b/.changeset/quiet-stream-close.md @@ -0,0 +1,5 @@ +--- +'eventsource': patch +--- + +Stop dispatching buffered messages when an event listener closes the connection. diff --git a/src/EventSource.ts b/src/EventSource.ts index 7588060..d31fa86 100644 --- a/src/EventSource.ts +++ b/src/EventSource.ts @@ -601,6 +601,11 @@ class EventSourceImpl extends EventTarget implements EventSource { * @internal */ #onEvent = (event: EventSourceMessage) => { + // A listener can close the connection while the parser is still processing this chunk. + if (this.#readyState === this.CLOSED) { + return + } + const origin = this.#redirectUrl ? this.#redirectUrl.origin : this.#url.origin // [spec] The `lastEventId` attribute is the last event ID string of the event // source, i.e. the persisted buffer (`#lastEventId`) - not the current event's `id`. diff --git a/test/client.test.ts b/test/client.test.ts index 7eb0fe8..b0414e3 100644 --- a/test/client.test.ts +++ b/test/client.test.ts @@ -44,6 +44,56 @@ const xOriginRedirectTest = suite === 'happy-dom' ? test.fails : test */ const onHandlerTest = suite === 'workerd' ? test.fails : test +test.each(['message', 'notice'])( + 'stops dispatching buffered %s events when a listener closes the connection', + async (eventType) => { + const seen: string[] = [] + const onMessage = getCallCounter({name: 'first message'}) + const es = new OurEventSource(`${serverUrl}/message-burst?event=${eventType}`, { + async fetch(url, init) { + const response = await request(url, init) + // Keep the real HTTP exchange, but make the chunk boundary deterministic. + const body = await response.arrayBuffer() + return new Response(body, {status: response.status, headers: response.headers}) + }, + }) + + es.addEventListener(eventType, (event) => { + seen.push(event.data) + es.close() + onMessage.listener(event) + }) + + try { + await onMessage.waitForCallCount(1) + expect(seen).toEqual(['first']) + expect(es.readyState).toBe(OurEventSource.CLOSED) + } finally { + es.close() + } + }, +) + +test.each(['message', 'notice'])( + 'dispatches every buffered %s event while open', + async (eventType) => { + const seen: string[] = [] + const onMessage = getCallCounter({name: 'messages'}) + const es = new OurEventSource(`${serverUrl}/message-burst?event=${eventType}`, esInit) + es.addEventListener(eventType, (event) => { + seen.push(event.data) + onMessage.listener(event) + }) + + try { + await onMessage.waitForCallCount(3) + expect(seen).toEqual(['first', 'second', 'third']) + } finally { + es.close() + } + }, +) + test('can connect, receive message, manually disconnect', async () => { const onMessage = getCallCounter({name: 'onMessage'}) const es = new OurEventSource(new URL(`${serverUrl}/`)) diff --git a/test/helpers/server.ts b/test/helpers/server.ts index cddddf3..9b2c223 100644 --- a/test/helpers/server.ts +++ b/test/helpers/server.ts @@ -61,6 +61,8 @@ export function handleRequest( return writeCounter(req, res) case '/mixed-ids': return writeMixedIds(req, res) + case '/message-burst': + return writeMessageBurst(req, res) case '/id-only': return writeIdOnly(req, res) case '/identified': @@ -145,6 +147,12 @@ async function writeCounter(req: IncomingMessage, res: ServerResponse) { res.end() } +function writeMessageBurst(req: IncomingMessage, res: ServerResponse) { + const event = new URL(req.url || '/', 'http://localhost').searchParams.get('event') || 'message' + res.writeHead(200, {'Content-Type': 'text/event-stream'}) + res.end(['first', 'second', 'third'].map((data) => encode({event, data})).join('')) +} + /** * Writes two messages: one with an `id` field, then one without. Per the spec, the second * event's `lastEventId` must still be `'1'`: the last event ID buffer is only updated by an From 172ff372c112831699adfd3c01cbef9a95189d16 Mon Sep 17 00:00:00 2001 From: Espen Hovlandsdal Date: Mon, 21 Sep 2026 10:14:45 -0700 Subject: [PATCH 2/2] test: add messageless event test to burst tests --- test/client.test.ts | 22 ++++++++++++++-------- test/helpers/server.ts | 2 +- 2 files changed, 15 insertions(+), 9 deletions(-) diff --git a/test/client.test.ts b/test/client.test.ts index b0414e3..fb609f7 100644 --- a/test/client.test.ts +++ b/test/client.test.ts @@ -44,9 +44,15 @@ const xOriginRedirectTest = suite === 'happy-dom' ? test.fails : test */ const onHandlerTest = suite === 'workerd' ? test.fails : test -test.each(['message', 'notice'])( - 'stops dispatching buffered %s events when a listener closes the connection', - async (eventType) => { +const bufferedEventTypes = [ + {name: 'nameless', eventType: ''}, + {name: 'message', eventType: 'message'}, + {name: 'notice', eventType: 'notice'}, +] + +test.each(bufferedEventTypes)( + 'stops dispatching buffered $name events when a listener closes the connection', + async ({eventType}) => { const seen: string[] = [] const onMessage = getCallCounter({name: 'first message'}) const es = new OurEventSource(`${serverUrl}/message-burst?event=${eventType}`, { @@ -58,7 +64,7 @@ test.each(['message', 'notice'])( }, }) - es.addEventListener(eventType, (event) => { + es.addEventListener(eventType || 'message', (event) => { seen.push(event.data) es.close() onMessage.listener(event) @@ -74,13 +80,13 @@ test.each(['message', 'notice'])( }, ) -test.each(['message', 'notice'])( - 'dispatches every buffered %s event while open', - async (eventType) => { +test.each(bufferedEventTypes)( + 'dispatches every buffered $name event while open', + async ({eventType}) => { const seen: string[] = [] const onMessage = getCallCounter({name: 'messages'}) const es = new OurEventSource(`${serverUrl}/message-burst?event=${eventType}`, esInit) - es.addEventListener(eventType, (event) => { + es.addEventListener(eventType || 'message', (event) => { seen.push(event.data) onMessage.listener(event) }) diff --git a/test/helpers/server.ts b/test/helpers/server.ts index 9b2c223..c7044d4 100644 --- a/test/helpers/server.ts +++ b/test/helpers/server.ts @@ -148,7 +148,7 @@ async function writeCounter(req: IncomingMessage, res: ServerResponse) { } function writeMessageBurst(req: IncomingMessage, res: ServerResponse) { - const event = new URL(req.url || '/', 'http://localhost').searchParams.get('event') || 'message' + const event = new URL(req.url || '/', 'http://localhost').searchParams.get('event') ?? 'message' res.writeHead(200, {'Content-Type': 'text/event-stream'}) res.end(['first', 'second', 'third'].map((data) => encode({event, data})).join('')) }