From 8d104e171f6ef4a83400b4093850bd66fe9a7ce8 Mon Sep 17 00:00:00 2001 From: Robert Dailey Date: Sat, 21 Mar 2026 12:59:41 -0500 Subject: [PATCH 1/7] feat(server): add keepAliveInterval to standalone GET SSE stream Adds an opt-in keepAliveInterval option to WebStandardStreamableHTTPServerTransportOptions that sends periodic SSE comments (`: keepalive`) on the standalone GET SSE stream. Reverse proxies commonly close connections that are idle for 30-60s. With no server-initiated messages, the GET SSE stream has no traffic during quiet periods, causing silent disconnections. This option lets operators send harmless SSE comments at a configurable cadence to keep the connection alive. The timer is cleared on close(), closeStandaloneSSEStream(), and on stream cancellation. Disabled by default; no behavior change for existing deployments. Addresses upstream #28, #876. --- .changeset/add-sse-keepalive.md | 5 ++ packages/server/src/server/streamableHttp.ts | 35 ++++++++ .../server/test/server/streamableHttp.test.ts | 89 +++++++++++++++++++ 3 files changed, 129 insertions(+) create mode 100644 .changeset/add-sse-keepalive.md diff --git a/.changeset/add-sse-keepalive.md b/.changeset/add-sse-keepalive.md new file mode 100644 index 0000000000..78654d9d9a --- /dev/null +++ b/.changeset/add-sse-keepalive.md @@ -0,0 +1,5 @@ +--- +'@modelcontextprotocol/server': minor +--- + +Add optional `keepAliveInterval` to `WebStandardStreamableHTTPServerTransportOptions` that sends periodic SSE comments on the standalone GET stream to prevent reverse proxy idle timeout disconnections. diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 7da5fb853c..1f99de9159 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -148,6 +148,15 @@ export interface WebStandardStreamableHTTPServerTransportOptions { */ retryInterval?: number; + /** + * Interval in milliseconds for sending SSE keepalive comments on the standalone + * GET SSE stream. When set, the transport sends periodic SSE comments + * (`: keepalive`) to prevent reverse proxies from closing idle connections. + * + * Disabled by default (no keepalive comments are sent). + */ + keepAliveInterval?: number; + /** * List of protocol versions that this transport will accept. * Used to validate the `mcp-protocol-version` header in incoming requests. @@ -246,6 +255,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private _allowedOrigins?: string[]; private _enableDnsRebindingProtection: boolean; private _retryInterval?: number; + private _keepAliveInterval?: number; + private _keepAliveTimer?: ReturnType; private _supportedProtocolVersions: string[]; sessionId?: string; @@ -263,6 +274,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { this._allowedOrigins = options.allowedOrigins; this._enableDnsRebindingProtection = options.enableDnsRebindingProtection ?? false; this._retryInterval = options.retryInterval; + this._keepAliveInterval = options.keepAliveInterval; this._supportedProtocolVersions = options.supportedProtocolVersions ?? SUPPORTED_PROTOCOL_VERSIONS; } @@ -473,6 +485,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // it still points at THIS controller — a stale cancel must not // delete a successor stream registered by a later GET/resume. if (this._streamMapping.get(this._standaloneSseStreamId)?.controller === streamController) { + this._clearKeepAliveTimer(); this._streamMapping.delete(this._standaloneSseStreamId); } } @@ -503,6 +516,19 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); + // Start keepalive timer to send periodic SSE comments that prevent + // reverse proxies from closing the connection due to idle timeouts + if (this._keepAliveInterval !== undefined) { + this._keepAliveTimer = setInterval(() => { + try { + streamController!.enqueue(encoder.encode(': keepalive\n\n')); + } catch { + // Controller is closed or errored, stop sending keepalives + this._clearKeepAliveTimer(); + } + }, this._keepAliveInterval); + } + return new Response(readable, { headers }); } @@ -974,11 +1000,19 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { return undefined; } + private _clearKeepAliveTimer(): void { + if (this._keepAliveTimer !== undefined) { + clearInterval(this._keepAliveTimer); + this._keepAliveTimer = undefined; + } + } + async close(): Promise { if (this._closed) { return; } this._closed = true; + this._clearKeepAliveTimer(); // Close all SSE connections for (const { cleanup } of this._streamMapping.values()) { @@ -1011,6 +1045,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Use this to implement polling behavior for server-initiated notifications. */ closeStandaloneSSEStream(): void { + this._clearKeepAliveTimer(); const stream = this._streamMapping.get(this._standaloneSseStreamId); if (stream) { stream.cleanup(); diff --git a/packages/server/test/server/streamableHttp.test.ts b/packages/server/test/server/streamableHttp.test.ts index beca451113..066203a293 100644 --- a/packages/server/test/server/streamableHttp.test.ts +++ b/packages/server/test/server/streamableHttp.test.ts @@ -1406,4 +1406,93 @@ describe('Zod v4', () => { expect(cleanupCalls).toEqual(['stream-1']); }); }); + + describe('HTTPServerTransport - keepAliveInterval', () => { + let transport: WebStandardStreamableHTTPServerTransport; + let mcpServer: McpServer; + + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(async () => { + vi.useRealTimers(); + await transport.close(); + }); + + async function setupTransport(keepAliveInterval?: number): Promise { + mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + + transport = new WebStandardStreamableHTTPServerTransport({ + sessionIdGenerator: () => randomUUID(), + keepAliveInterval + }); + + await mcpServer.connect(transport); + + const initReq = createRequest('POST', TEST_MESSAGES.initialize); + const initRes = await transport.handleRequest(initReq); + return initRes.headers.get('mcp-session-id') as string; + } + + it('should send SSE keepalive comments periodically when keepAliveInterval is set', async () => { + const sessionId = await setupTransport(50); + + const getReq = createRequest('GET', undefined, { sessionId }); + const getRes = await transport.handleRequest(getReq); + + expect(getRes.status).toBe(200); + expect(getRes.body).not.toBeNull(); + + const reader = getRes.body!.getReader(); + + // Advance past two intervals to accumulate keepalive comments + vi.advanceTimersByTime(120); + + const { value } = await reader.read(); + const text = new TextDecoder().decode(value); + expect(text).toContain(': keepalive'); + }); + + it('should not send SSE comments when keepAliveInterval is not set', async () => { + const sessionId = await setupTransport(undefined); + + const getReq = createRequest('GET', undefined, { sessionId }); + const getRes = await transport.handleRequest(getReq); + + expect(getRes.status).toBe(200); + expect(getRes.body).not.toBeNull(); + + const reader = getRes.body!.getReader(); + + // Advance time; no keepalive should be enqueued + vi.advanceTimersByTime(200); + + // Close the transport to end the stream, then read whatever was buffered + await transport.close(); + + const chunks: string[] = []; + for (let result = await reader.read(); !result.done; result = await reader.read()) { + chunks.push(new TextDecoder().decode(result.value)); + } + + const allText = chunks.join(''); + expect(allText).not.toContain(': keepalive'); + }); + + it('should clear the keepalive interval when the transport is closed', async () => { + const sessionId = await setupTransport(50); + + const getReq = createRequest('GET', undefined, { sessionId }); + const getRes = await transport.handleRequest(getReq); + + expect(getRes.status).toBe(200); + + // Close the transport, which should clear the interval + await transport.close(); + + // Advancing timers after close should not throw + vi.advanceTimersByTime(200); + }); + }); }); From 939df127cb70d7a4426d6b54f028694ed816700f Mon Sep 17 00:00:00 2001 From: Robert Dailey Date: Tue, 31 Mar 2026 18:29:49 -0500 Subject: [PATCH 2/7] fix(server): add keepalive timer in replayEvents() and fix test The replayEvents() code path (client reconnects with Last-Event-ID) was missing keepalive timer setup, so reconnecting clients would lose keepalive protection and get dropped again at the next proxy idle timeout. Also fixes the cleanup test to actually prove that close() clears the timer by asserting vi.getTimerCount() drops to 0, instead of relying on the catch fallback which would self-clear anyway. Addresses PR review feedback from @felixweinberger on PR #1726. --- packages/server/src/server/streamableHttp.ts | 13 +++++++++++++ packages/server/test/server/streamableHttp.test.ts | 4 ++-- 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 1f99de9159..724b7b1204 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -590,6 +590,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // a stale cancel from an earlier resume must not delete a // successor resumed stream a re-poll has since registered. if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) { + this._clearKeepAliveTimer(); this._streamMapping.delete(replayedStreamId); } } @@ -644,6 +645,18 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + // Start keepalive timer for the replayed stream so reconnecting + // clients remain protected from proxy idle timeouts + if (this._keepAliveInterval !== undefined) { + this._keepAliveTimer = setInterval(() => { + try { + streamController!.enqueue(encoder.encode(': keepalive\n\n')); + } catch { + this._clearKeepAliveTimer(); + } + }, this._keepAliveInterval); + } + return new Response(readable, { headers }); } catch (error) { this.onerror?.(error as Error); diff --git a/packages/server/test/server/streamableHttp.test.ts b/packages/server/test/server/streamableHttp.test.ts index 066203a293..98d65fcf5a 100644 --- a/packages/server/test/server/streamableHttp.test.ts +++ b/packages/server/test/server/streamableHttp.test.ts @@ -1487,12 +1487,12 @@ describe('Zod v4', () => { const getRes = await transport.handleRequest(getReq); expect(getRes.status).toBe(200); + expect(vi.getTimerCount()).toBe(1); // Close the transport, which should clear the interval await transport.close(); - // Advancing timers after close should not throw - vi.advanceTimersByTime(200); + expect(vi.getTimerCount()).toBe(0); }); }); }); From c3903075a632e530a97219ca53dbf879c94ed27e Mon Sep 17 00:00:00 2001 From: Robert Dailey Date: Wed, 1 Apr 2026 07:39:02 -0500 Subject: [PATCH 3/7] fix(server): clear keepalive timer before reassignment in replayEvents Prevent timer handle leaks when concurrent reconnect requests bypass the conflict check. This edge case is possible when EventStore omits the optional getStreamIdForEventId method, allowing duplicate replay attempts to start new timers without cleaning up existing ones. --- packages/server/src/server/streamableHttp.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 724b7b1204..4153d908ce 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -648,6 +648,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // Start keepalive timer for the replayed stream so reconnecting // clients remain protected from proxy idle timeouts if (this._keepAliveInterval !== undefined) { + this._clearKeepAliveTimer(); this._keepAliveTimer = setInterval(() => { try { streamController!.enqueue(encoder.encode(': keepalive\n\n')); From a2df5c106a426723066ea9a0182ae3443788c92b Mon Sep 17 00:00:00 2001 From: Robert Dailey Date: Wed, 1 Apr 2026 20:09:01 -0500 Subject: [PATCH 4/7] fix(server): prevent duplicate and accidental keepalive timer clears Adds defensive _clearKeepAliveTimer() before setInterval in handleGetRequest() to prevent duplicate timers when the GET stream is reinitialized. Moves _clearKeepAliveTimer() inside the if(stream) guard in closeStandaloneSSEStream() so it only clears when the standalone stream is being torn down, preventing accidental clearing of a replay stream's timer. Closes #1726 --- packages/server/src/server/streamableHttp.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 4153d908ce..cd2b7e580a 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -519,6 +519,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // Start keepalive timer to send periodic SSE comments that prevent // reverse proxies from closing the connection due to idle timeouts if (this._keepAliveInterval !== undefined) { + this._clearKeepAliveTimer(); this._keepAliveTimer = setInterval(() => { try { streamController!.enqueue(encoder.encode(': keepalive\n\n')); @@ -1059,9 +1060,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Use this to implement polling behavior for server-initiated notifications. */ closeStandaloneSSEStream(): void { - this._clearKeepAliveTimer(); const stream = this._streamMapping.get(this._standaloneSseStreamId); if (stream) { + this._clearKeepAliveTimer(); stream.cleanup(); } } From 0487cff7ac6c1864199e1978af44858b4ee94b09 Mon Sep 17 00:00:00 2001 From: Felix Weinberger Date: Thu, 2 Apr 2026 15:32:29 +0000 Subject: [PATCH 5/7] Move keepAliveTimer to per-stream ownership and add concurrent-stream regression test --- packages/server/src/server/streamableHttp.ts | 93 +++++++++---------- .../server/test/server/streamableHttp.test.ts | 52 +++++++++++ 2 files changed, 96 insertions(+), 49 deletions(-) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index cd2b7e580a..3dc2064192 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -72,6 +72,8 @@ interface StreamMapping { replayedEventIds?: Set; /** Cleanup function to close stream and remove mapping */ cleanup: () => void; + /** Per-stream keepalive timer; cleared by this stream's cleanup/cancel */ + keepAliveTimer?: ReturnType; } /** @@ -256,7 +258,6 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private _enableDnsRebindingProtection: boolean; private _retryInterval?: number; private _keepAliveInterval?: number; - private _keepAliveTimer?: ReturnType; private _supportedProtocolVersions: string[]; sessionId?: string; @@ -473,19 +474,33 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } const encoder = new TextEncoder(); - let streamController: ReadableStreamDefaultController; + let streamController!: ReadableStreamDefaultController; + + const mapping: StreamMapping = { + encoder, + cleanup: () => { + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); + this._streamMapping.delete(this._standaloneSseStreamId); + try { + streamController.close(); + } catch { + // Controller might already be closed + } + } + }; // Create a ReadableStream with a controller we can use to push SSE events const readable = new ReadableStream({ start: controller => { streamController = controller; + mapping.controller = controller; }, cancel: () => { // Stream was cancelled by client. Only drop the mapping when // it still points at THIS controller — a stale cancel must not // delete a successor stream registered by a later GET/resume. if (this._streamMapping.get(this._standaloneSseStreamId)?.controller === streamController) { - this._clearKeepAliveTimer(); + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); this._streamMapping.delete(this._standaloneSseStreamId); } } @@ -502,30 +517,17 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { headers['mcp-session-id'] = this.sessionId; } - // Store the stream mapping with the controller for pushing data - this._streamMapping.set(this._standaloneSseStreamId, { - controller: streamController!, - encoder, - cleanup: () => { - this._streamMapping.delete(this._standaloneSseStreamId); - try { - streamController!.close(); - } catch { - // Controller might already be closed - } - } - }); + this._streamMapping.set(this._standaloneSseStreamId, mapping); // Start keepalive timer to send periodic SSE comments that prevent // reverse proxies from closing the connection due to idle timeouts if (this._keepAliveInterval !== undefined) { - this._clearKeepAliveTimer(); - this._keepAliveTimer = setInterval(() => { + mapping.keepAliveTimer = setInterval(() => { try { - streamController!.enqueue(encoder.encode(': keepalive\n\n')); + streamController.enqueue(encoder.encode(': keepalive\n\n')); } catch { // Controller is closed or errored, stop sending keepalives - this._clearKeepAliveTimer(); + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); } }, this._keepAliveInterval); } @@ -579,9 +581,23 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // eslint-disable-next-line prefer-const let replayedStreamId: string | undefined; + const mapping: StreamMapping = { + encoder, + cleanup: () => { + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); + this._streamMapping.delete(replayedStreamId!); + try { + streamController!.close(); + } catch { + // Controller might already be closed + } + } + }; + const readable = new ReadableStream({ start: controller => { streamController = controller; + mapping.controller = controller; }, cancel: () => { // Stream was cancelled by client — drop the mapping so a @@ -591,7 +607,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // a stale cancel from an earlier resume must not delete a // successor resumed stream a re-poll has since registered. if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) { - this._clearKeepAliveTimer(); + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); this._streamMapping.delete(replayedStreamId); } } @@ -605,7 +621,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const success = this.writeSSEEvent(streamController!, encoder, message, eventId); if (!success) { try { - streamController!.close(); + streamController.close(); } catch { // Controller might already be closed } @@ -613,19 +629,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); - this._streamMapping.set(replayedStreamId, { - controller: streamController!, - encoder, - replayedEventIds, - cleanup: () => { - this._streamMapping.delete(replayedStreamId!); - try { - streamController!.close(); - } catch { - // Controller might already be closed - } - } - }); + mapping.replayedEventIds = replayedEventIds; + this._streamMapping.set(replayedStreamId!, mapping); // If this is a per-request stream and no in-flight request still // targets this streamId, the request was already retired by the @@ -649,12 +654,11 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // Start keepalive timer for the replayed stream so reconnecting // clients remain protected from proxy idle timeouts if (this._keepAliveInterval !== undefined) { - this._clearKeepAliveTimer(); - this._keepAliveTimer = setInterval(() => { + mapping.keepAliveTimer = setInterval(() => { try { - streamController!.enqueue(encoder.encode(': keepalive\n\n')); + streamController.enqueue(encoder.encode(': keepalive\n\n')); } catch { - this._clearKeepAliveTimer(); + if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); } }, this._keepAliveInterval); } @@ -1015,21 +1019,13 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { return undefined; } - private _clearKeepAliveTimer(): void { - if (this._keepAliveTimer !== undefined) { - clearInterval(this._keepAliveTimer); - this._keepAliveTimer = undefined; - } - } - async close(): Promise { if (this._closed) { return; } this._closed = true; - this._clearKeepAliveTimer(); - // Close all SSE connections + // Close all SSE connections (each cleanup() also clears its own keepAliveTimer) for (const { cleanup } of this._streamMapping.values()) { cleanup(); } @@ -1062,7 +1058,6 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { closeStandaloneSSEStream(): void { const stream = this._streamMapping.get(this._standaloneSseStreamId); if (stream) { - this._clearKeepAliveTimer(); stream.cleanup(); } } diff --git a/packages/server/test/server/streamableHttp.test.ts b/packages/server/test/server/streamableHttp.test.ts index 98d65fcf5a..23be0487e4 100644 --- a/packages/server/test/server/streamableHttp.test.ts +++ b/packages/server/test/server/streamableHttp.test.ts @@ -1494,5 +1494,57 @@ describe('Zod v4', () => { expect(vi.getTimerCount()).toBe(0); }); + + it('should maintain independent keepalive timers per concurrent stream', async () => { + // Minimal event store so a GET with last-event-id triggers replayEvents() + // and creates a second StreamMapping entry concurrently with the standalone GET. + const eventStore: EventStore = { + async storeEvent(streamId) { + return `${streamId}_evt`; + }, + async getStreamIdForEventId() { + return 'replay-stream'; + }, + async replayEventsAfter() { + return 'replay-stream'; + } + }; + + mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + transport = new WebStandardStreamableHTTPServerTransport({ + sessionIdGenerator: () => randomUUID(), + keepAliveInterval: 50, + eventStore + }); + await mcpServer.connect(transport); + const initRes = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + const sessionId = initRes.headers.get('mcp-session-id') as string; + + // Stream A: standalone GET + const resA = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + expect(resA.status).toBe(200); + const readerA = resA.body!.getReader(); + expect(vi.getTimerCount()).toBe(1); + + // Stream B: GET with last-event-id -> replayEvents path, separate mapping key + const resB = await transport.handleRequest( + createRequest('GET', undefined, { sessionId, extraHeaders: { 'last-event-id': 'evt-1' } }) + ); + expect(resB.status).toBe(200); + const readerB = resB.body!.getReader(); + expect(vi.getTimerCount()).toBe(2); + + // Cancel stream B; its keepalive timer must be cleared without affecting A's + await readerB.cancel(); + expect(vi.getTimerCount()).toBe(1); + + // Stream A still receives keepalives after B is cancelled + vi.advanceTimersByTime(60); + const { value } = await readerA.read(); + expect(new TextDecoder().decode(value)).toContain(': keepalive'); + + await readerA.cancel(); + expect(vi.getTimerCount()).toBe(0); + }); }); }); From be9c114a9aa32ef642985a4b52d368280dd9f7d2 Mon Sep 17 00:00:00 2001 From: Felix Weinberger Date: Thu, 2 Apr 2026 15:35:00 +0000 Subject: [PATCH 6/7] Fix prefer-const lint: use closure-local keepAliveTimer in replayEvents --- packages/server/src/server/streamableHttp.ts | 39 ++++++++++---------- 1 file changed, 20 insertions(+), 19 deletions(-) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 3dc2064192..881a003f16 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -580,24 +580,11 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // replayEventsAfter resolves) — must be `let`. // eslint-disable-next-line prefer-const let replayedStreamId: string | undefined; - - const mapping: StreamMapping = { - encoder, - cleanup: () => { - if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); - this._streamMapping.delete(replayedStreamId!); - try { - streamController!.close(); - } catch { - // Controller might already be closed - } - } - }; + let keepAliveTimer: ReturnType | undefined; const readable = new ReadableStream({ start: controller => { streamController = controller; - mapping.controller = controller; }, cancel: () => { // Stream was cancelled by client — drop the mapping so a @@ -607,7 +594,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // a stale cancel from an earlier resume must not delete a // successor resumed stream a re-poll has since registered. if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) { - if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); + if (keepAliveTimer) clearInterval(keepAliveTimer); this._streamMapping.delete(replayedStreamId); } } @@ -629,7 +616,20 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); - mapping.replayedEventIds = replayedEventIds; + const mapping: StreamMapping = { + controller: streamController!, + encoder, + replayedEventIds, + cleanup: () => { + if (keepAliveTimer) clearInterval(keepAliveTimer); + this._streamMapping.delete(replayedStreamId!); + try { + streamController!.close(); + } catch { + // Controller might already be closed + } + } + }; this._streamMapping.set(replayedStreamId!, mapping); // If this is a per-request stream and no in-flight request still @@ -654,13 +654,14 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // Start keepalive timer for the replayed stream so reconnecting // clients remain protected from proxy idle timeouts if (this._keepAliveInterval !== undefined) { - mapping.keepAliveTimer = setInterval(() => { + keepAliveTimer = setInterval(() => { try { - streamController.enqueue(encoder.encode(': keepalive\n\n')); + streamController!.enqueue(encoder.encode(': keepalive\n\n')); } catch { - if (mapping.keepAliveTimer) clearInterval(mapping.keepAliveTimer); + if (keepAliveTimer) clearInterval(keepAliveTimer); } }, this._keepAliveInterval); + mapping.keepAliveTimer = keepAliveTimer; } return new Response(readable, { headers }); From 82ee1a609eaf81ae1d20edf0121d4c54d80af924 Mon Sep 17 00:00:00 2001 From: Robert Dailey Date: Sat, 18 Jul 2026 09:25:04 -0500 Subject: [PATCH 7/7] fix(server): always clear keepalive timer on replay stream cancel The early-close block for completed replay requests deletes the stream mapping entry before the keepalive timer is started. When the client subsequently closes the stream, the cancel callback's stale-guard check (which gates timer cleanup on the mapping still pointing at this controller) cannot find the mapping, leaving the timer running. - Move clearInterval above the stale-guard check so it always fires for the closure-local timer - Guard keepalive timer startup with a mapping-presence check to skip setup when early-close already removed the entry - Simulate an in-flight request in the concurrent-stream test so the replay stream stays open past the early-close block --- packages/server/src/server/streamableHttp.ts | 20 ++++++++++--------- .../server/test/server/streamableHttp.test.ts | 5 ++++- 2 files changed, 15 insertions(+), 10 deletions(-) diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 881a003f16..4f1c4d51f6 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -587,14 +587,15 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { streamController = controller; }, cancel: () => { - // Stream was cancelled by client — drop the mapping so a - // subsequent reconnect with the same Last-Event-ID is not - // refused with 409 by the conflict check above. Only delete - // when the mapped entry is still THIS closure's controller: - // a stale cancel from an earlier resume must not delete a - // successor resumed stream a re-poll has since registered. + // Always clear the closure-local keepalive timer; the + // mapping may already have been removed by the early-close + // block for completed requests, but the timer is still ours. + if (keepAliveTimer) clearInterval(keepAliveTimer); + // Drop the mapping so a subsequent reconnect with the same + // Last-Event-ID is not refused with 409. Only delete when + // the mapped entry is still THIS closure's controller: a + // stale cancel must not delete a successor resumed stream. if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) { - if (keepAliveTimer) clearInterval(keepAliveTimer); this._streamMapping.delete(replayedStreamId); } } @@ -652,8 +653,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } // Start keepalive timer for the replayed stream so reconnecting - // clients remain protected from proxy idle timeouts - if (this._keepAliveInterval !== undefined) { + // clients remain protected from proxy idle timeouts. + // Skip if the early-close block above already removed the mapping. + if (this._keepAliveInterval !== undefined && this._streamMapping.has(replayedStreamId!)) { keepAliveTimer = setInterval(() => { try { streamController!.enqueue(encoder.encode(': keepalive\n\n')); diff --git a/packages/server/test/server/streamableHttp.test.ts b/packages/server/test/server/streamableHttp.test.ts index 23be0487e4..01e0c1972c 100644 --- a/packages/server/test/server/streamableHttp.test.ts +++ b/packages/server/test/server/streamableHttp.test.ts @@ -1526,7 +1526,10 @@ describe('Zod v4', () => { const readerA = resA.body!.getReader(); expect(vi.getTimerCount()).toBe(1); - // Stream B: GET with last-event-id -> replayEvents path, separate mapping key + // Stream B: GET with last-event-id -> replayEvents path, separate mapping key. + // Simulate an in-flight request for the replay stream so the + // early-close-for-completed-requests block keeps the stream open. + (transport as any)._requestToStreamMapping.set('fake-req', 'replay-stream'); const resB = await transport.handleRequest( createRequest('GET', undefined, { sessionId, extraHeaders: { 'last-event-id': 'evt-1' } }) );