From 4a79af3fb7d4613d0932084db1bea391a731a4d4 Mon Sep 17 00:00:00 2001 From: Isaac Sanchez-Hawkins <266845420+isanchez404@users.noreply.github.com> Date: Fri, 8 May 2026 08:14:39 -0400 Subject: [PATCH] fix(event-stream): cancel unavailable upstream bodies (#1142) Co-authored-by: Isaac Sanchez --- .../web/server/lib/event-stream/upstream-reader.js | 7 +++++++ .../server/lib/event-stream/upstream-reader.test.js | 12 +++++++++++- 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/packages/web/server/lib/event-stream/upstream-reader.js b/packages/web/server/lib/event-stream/upstream-reader.js index 0acff519..1e4d115a 100644 --- a/packages/web/server/lib/event-stream/upstream-reader.js +++ b/packages/web/server/lib/event-stream/upstream-reader.js @@ -31,6 +31,12 @@ function normalizeHeaders(headers) { return { ...headers }; } +async function cancelResponseBody(response) { + if (response?.body && typeof response.body.cancel === 'function') { + await response.body.cancel().catch(() => {}); + } +} + export function createUpstreamSseReader({ buildUrl, getHeaders = () => ({}), @@ -116,6 +122,7 @@ export function createUpstreamSseReader({ status: response?.status ?? 0, response, }); + await cancelResponseBody(response); await waitForReconnectDelay(reconnectDelayMs, signal); continue; } 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 616b5b96..714eb6a2 100644 --- a/packages/web/server/lib/event-stream/upstream-reader.test.js +++ b/packages/web/server/lib/event-stream/upstream-reader.test.js @@ -165,6 +165,7 @@ describe('createUpstreamSseReader', () => { it('reports unavailable upstream responses and continues reconnecting until stopped', async () => { const errors = []; let attempt = 0; + let unavailableBodyCanceled = false; let reader; reader = createUpstreamSseReader({ @@ -173,7 +174,15 @@ describe('createUpstreamSseReader', () => { fetchImpl: async (_url, options) => { attempt += 1; if (attempt === 1) { - return { ok: false, status: 503, body: null }; + return { + ok: false, + status: 503, + body: { + cancel: async () => { + unavailableBodyCanceled = true; + }, + }, + }; } return createSseResponse({ @@ -199,6 +208,7 @@ describe('createUpstreamSseReader', () => { status: 503, }), ]); + expect(unavailableBodyCanceled).toBe(true); expect(attempt).toBe(2); }); });