From 2b5408406c1d5ace0567bc9980a72e1694fd719f Mon Sep 17 00:00:00 2001 From: Isaac Sanchez-Hawkins <266845420+isanchez404@users.noreply.github.com> Date: Sat, 23 May 2026 13:48:16 -0400 Subject: [PATCH] fix: clean up upstream reader abort listener (#1378) Co-authored-by: Isaac Sanchez --- .../lib/event-stream/upstream-reader.js | 22 +++++++++++++++---- .../lib/event-stream/upstream-reader.test.js | 5 ++--- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/packages/web/server/lib/event-stream/upstream-reader.js b/packages/web/server/lib/event-stream/upstream-reader.js index 77b37849..bfa4d69a 100644 --- a/packages/web/server/lib/event-stream/upstream-reader.js +++ b/packages/web/server/lib/event-stream/upstream-reader.js @@ -63,21 +63,34 @@ export function createUpstreamSseReader({ let stopped = false; let activeController = null; let lastEventId = typeof initialLastEventId === 'string' ? initialLastEventId : ''; + let stopListenerAttached = false; - const stop = () => { + function detachStopListener() { + if (!stopListenerAttached) return; + signal?.removeEventListener('abort', stop); + stopListenerAttached = false; + } + + function attachStopListener() { + if (!signal || signal.aborted || stopListenerAttached) return; + signal.addEventListener('abort', stop, { once: true }); + stopListenerAttached = true; + } + + function stop() { stopped = true; + detachStopListener(); if (activeController && !activeController.signal.aborted) { activeController.abort(); } - }; - - signal?.addEventListener('abort', stop, { once: true }); + } const start = () => { if (running) { return running; } + attachStopListener(); stopped = false; running = (async () => { while (!stopped && !signal?.aborted) { @@ -210,6 +223,7 @@ export function createUpstreamSseReader({ } } })().finally(() => { + detachStopListener(); running = null; }); diff --git a/packages/web/server/lib/event-stream/upstream-reader.test.js b/packages/web/server/lib/event-stream/upstream-reader.test.js index 1d7c8fc3..f6dd6217 100644 --- a/packages/web/server/lib/event-stream/upstream-reader.test.js +++ b/packages/web/server/lib/event-stream/upstream-reader.test.js @@ -234,7 +234,7 @@ describe('createUpstreamSseReader', () => { expect(attempt).toBe(2); }); - it('removes reconnect delay abort listeners after normal timeout completion', async () => { + it('removes abort listeners after stop', async () => { const tracked = createTrackedSignal(); let attempt = 0; let reader; @@ -270,7 +270,6 @@ describe('createUpstreamSseReader', () => { await reader.start(); expect(attempt).toBe(2); - // The top-level stop listener remains; the reconnect-delay listener should be removed. - expect(tracked.getListenerCount()).toBe(1); + expect(tracked.getListenerCount()).toBe(0); }); });