From c5d22b29305aab1743544cbe51e6e623d51975a0 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 11:24:07 +0100 Subject: [PATCH 1/6] fix: harden SSE keep-alive lifecycle (backport of the v2 review fixes) Backports to v1.x the keep-alive hardening that landed on main after the original v1 fix (#2538) merged: - startKeepAlive is a no-op after transport close, so a deferred arm (e.g. an event-store replayEventsAfter await straddling close()) cannot create a timer the close sweep already missed. - The POST SSE path arms keep-alive after the fallible awaits (priming event write, message dispatch) instead of before them, so an error path that discards the Response cannot leak a permanently-firing timer against a stream nothing can cancel. - The guard uses the !(keepAliveMs > 0) polarity, so a non-finite value (e.g. parseInt of an unset env var) disables keep-alive instead of arming a Node-clamped ~1ms interval that floods every stream. Each fix has a fake-timer regression test. --- src/server/webStandardStreamableHttp.ts | 20 +++++- test/server/streamableHttp.test.ts | 96 +++++++++++++++++++++++++ 2 files changed, 113 insertions(+), 3 deletions(-) diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index 932ad56600..d1f7794979 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -227,6 +227,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // when sessionId is not set (undefined), it means the transport is in stateless mode private sessionIdGenerator: (() => string) | undefined; private _started: boolean = false; + private _closed: boolean = false; private _hasHandledRequest: boolean = false; private _streamMapping: Map = new Map(); private _requestToStreamMapping: Map = new Map(); @@ -271,7 +272,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * clears itself if a write fails (stream already closed/cancelled). */ private startKeepAlive(streamId: string, controller: ReadableStreamDefaultController, encoder: TextEncoder): void { - if (this._keepAliveMs <= 0) { + // A deferred arm (e.g. after an event-store await that straddled + // close()) must not outlive the transport: close()'s timer sweep has + // already run. The `> 0` polarity disables keep-alive for non-finite + // values (NaN fails every comparison) instead of arming a + // Node-clamped ~1ms interval. + if (!(this._keepAliveMs > 0) || this._closed) { return; } this.stopKeepAlive(streamId); @@ -842,8 +848,6 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } - this.startKeepAlive(streamId, streamController!, encoder); - // Write priming event if event store is configured (after mapping is set up) await this.writePrimingEvent(streamController!, encoder, streamId, clientProtocolVersion); @@ -869,6 +873,14 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // The server SHOULD NOT close the SSE stream before sending all JSON-RPC responses // This will be handled by the send() method when responses are ready + // Arm keep-alive only after the fallible awaits above — an error + // path returning 400 discards the Response, so nothing could ever + // cancel the stream and clear an already-armed timer. Skip if the + // responses already completed and cleaned the stream up. + if (this._streamMapping.get(streamId)?.controller === streamController!) { + this.startKeepAlive(streamId, streamController!, encoder); + } + return new Response(readable, { status: 200, headers }); } catch (error) { // return JSON-RPC formatted error @@ -961,6 +973,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } async close(): Promise { + this._closed = true; + // Close all SSE connections this._streamMapping.forEach(({ cleanup }) => { cleanup(); diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index 2046a71697..d61ff623c8 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -3290,6 +3290,10 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { }); } + function withSession(sessionId: string, extra?: Record): Record { + return { 'mcp-session-id': sessionId, 'mcp-protocol-version': '2025-11-25', ...extra }; + } + async function createTransport(options?: { keepAliveMs?: number }): Promise<{ transport: WebStandardStreamableHTTPServerTransport; sessionId: string; @@ -3447,4 +3451,96 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { await transport.close(); }); + + it('should disable keep-alive for a non-finite keepAliveMs instead of arming a clamped interval', async () => { + const { transport, sessionId } = await createTransport({ keepAliveMs: Number.NaN }); + + const response = await transport.handleRequest(req('GET', { headers: withSession(sessionId) })); + expect(response.status).toBe(200); + + // No timer may be armed: setInterval(fn, NaN) would be clamped by + // Node to ~1ms and flood the stream with keep-alive frames. + expect(vi.getTimerCount()).toBe(0); + const reader = response.body!.getReader(); + await vi.advanceTimersByTimeAsync(60000); + const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); + expect(raced).toBe('pending'); + + await transport.close(); + }); + + it('should not arm keep-alive when the transport closes during an event-store replay await', async () => { + let releaseReplay: (() => void) | undefined; + const eventStore: EventStore = { + async storeEvent(): Promise { + return 'evt-1'; + }, + async replayEventsAfter(): Promise { + await new Promise(resolve => { + releaseReplay = resolve; + }); + return '_GET_stream'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport); + const initResponse = await transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize })); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + + // Enter replayEvents and park on the replayEventsAfter await + const pendingGet = transport.handleRequest(req('GET', { headers: { ...withSession(sessionId), 'Last-Event-ID': 'evt-1' } })); + await vi.advanceTimersByTimeAsync(0); + expect(releaseReplay).toBeDefined(); + + // Close the transport mid-await, then let the replay continuation run + await transport.close(); + releaseReplay?.(); + await pendingGet; + + // The deferred continuation must not have armed a timer close() can never sweep + expect(vi.getTimerCount()).toBe(0); + }); + + it('should not leak a keep-alive timer when the priming event write fails on a POST SSE stream', async () => { + // Healthy during initialization, then the store starts failing — the + // tool call's priming event write must reject inside handlePostRequest + let storeFails = false; + const eventStore: EventStore = { + async storeEvent(): Promise { + if (storeFails) { + throw new Error('event store unavailable'); + } + return `evt-${randomUUID()}`; + }, + async replayEventsAfter(): Promise { + return 'stream-1'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + mcpServer.registerTool('noop', { description: 'noop' }, async () => ({ content: [] })); + await mcpServer.connect(transport); + + const initResponse = await transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize })); + expect(initResponse.status).toBe(200); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + // Let the init response finish sending (its send() stores an event and + // then cleans up the init stream's keep-alive) before failing the store + await vi.advanceTimersByTimeAsync(0); + expect(vi.getTimerCount()).toBe(0); + storeFails = true; + + const response = await transport.handleRequest( + req('POST', { + body: { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'call-1' }, + headers: withSession(sessionId) + }) + ); + expect(response.status).toBe(400); + + // The discarded stream must not carry a permanently-firing timer + expect(vi.getTimerCount()).toBe(0); + + await transport.close(); + }); }); From a333f8ba72de88e1b8f92868f2db560e640d0fee Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 11:25:09 +0100 Subject: [PATCH 2/6] chore: add changeset --- .changeset/keepalive-lifecycle-hardening.md | 5 +++++ 1 file changed, 5 insertions(+) create mode 100644 .changeset/keepalive-lifecycle-hardening.md diff --git a/.changeset/keepalive-lifecycle-hardening.md b/.changeset/keepalive-lifecycle-hardening.md new file mode 100644 index 0000000000..39fd2d6278 --- /dev/null +++ b/.changeset/keepalive-lifecycle-hardening.md @@ -0,0 +1,5 @@ +--- +'@modelcontextprotocol/sdk': patch +--- + +Hardens the Streamable HTTP server transport's SSE keep-alive lifecycle: keep-alive timers can no longer be armed after the transport closes or leak when a priming event write fails, and a non-finite `keepAliveMs` disables keep-alive instead of arming a clamped ~1ms interval. From 739c0095dd220cf5a6ef966cbe2c1517a27d325a Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 11:55:20 +0100 Subject: [PATCH 3/6] fix: complete the keep-alive error-path hardening Addresses the remaining review findings on the same paths: - The POST SSE error path now reclaims the stream bookkeeping too, not just the timer: a failure after registration (failed priming event write, throwing handler) runs the registered cleanup and retires the request mappings, so repeated failures during an event-store outage no longer accrete orphaned stream entries for the session's lifetime. - The replay continuation checks _closed after the replayEventsAfter await: when close() raced the await, the stream is ended instead of re-registered on the dead transport, so the client observes stream end (rather than hanging) and a later resume isn't 409-blocked by a stale entry. - The keep-alive guard rejects all non-finite intervals: Infinity (and any out-of-range delay) is clamped by setInterval to ~1ms just like NaN, so the guard now requires Number.isFinite as well as > 0. Tests extended accordingly: the non-finite case is parameterized over NaN and Infinity, and the error-path tests assert the internal stream maps are empty, not just the timer count. --- src/server/webStandardStreamableHttp.ts | 39 +++++++++++++++++-- test/server/streamableHttp.test.ts | 50 ++++++++++++++++++------- 2 files changed, 71 insertions(+), 18 deletions(-) diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index d1f7794979..ca0843f796 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -274,10 +274,10 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private startKeepAlive(streamId: string, controller: ReadableStreamDefaultController, encoder: TextEncoder): void { // A deferred arm (e.g. after an event-store await that straddled // close()) must not outlive the transport: close()'s timer sweep has - // already run. The `> 0` polarity disables keep-alive for non-finite - // values (NaN fails every comparison) instead of arming a - // Node-clamped ~1ms interval. - if (!(this._keepAliveMs > 0) || this._closed) { + // already run. Non-finite intervals (NaN, Infinity) and non-positive + // ones disable keep-alive: setInterval would clamp them to ~1ms and + // flood every stream with comment frames. + if (!Number.isFinite(this._keepAliveMs) || this._keepAliveMs <= 0 || this._closed) { return; } this.stopKeepAlive(streamId); @@ -589,6 +589,21 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); + // The transport may have closed while the replay await was parked: + // its cleanup sweep ran before this stream was registered, so + // registering now would strand a mapping (and an open controller) + // on a dead transport, hang the client on a stream that never + // ends, and 409-block a later resume of this stream id. End the + // stream instead so the client observes termination. + if (this._closed) { + try { + streamController!.close(); + } catch { + // Controller might already be closed + } + return new Response(readable, { headers }); + } + this._streamMapping.set(replayedStreamId, { controller: streamController!, encoder, @@ -664,6 +679,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Handles POST requests containing JSON-RPC messages */ private async handlePostRequest(req: Request, options?: HandleRequestOptions): Promise { + // Set once the SSE stream bookkeeping has been registered, so the + // catch below can reclaim it: an error after registration (a failed + // priming event write, a throwing message handler) returns 400 and + // discards the Response, leaving nothing that could ever cancel the + // stream or retire the request mappings. + let reclaimSseBookkeeping: (() => void) | undefined; try { // Validate the Accept header const acceptHeader = req.headers.get('accept'); @@ -848,6 +869,15 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + reclaimSseBookkeeping = () => { + this._streamMapping.get(streamId)?.cleanup(); + for (const message of messages) { + if (isJSONRPCRequest(message)) { + this._requestToStreamMapping.delete(message.id); + } + } + }; + // Write priming event if event store is configured (after mapping is set up) await this.writePrimingEvent(streamController!, encoder, streamId, clientProtocolVersion); @@ -885,6 +915,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } catch (error) { // return JSON-RPC formatted error this.onerror?.(error as Error); + reclaimSseBookkeeping?.(); return this.createJsonErrorResponse(400, -32700, 'Parse error', { data: String(error) }); } } diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index d61ff623c8..85752e1079 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -3452,22 +3452,26 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { await transport.close(); }); - it('should disable keep-alive for a non-finite keepAliveMs instead of arming a clamped interval', async () => { - const { transport, sessionId } = await createTransport({ keepAliveMs: Number.NaN }); + it.each([Number.NaN, Number.POSITIVE_INFINITY])( + 'should disable keep-alive for non-finite keepAliveMs %s instead of arming a clamped interval', + async keepAliveMs => { + const { transport, sessionId } = await createTransport({ keepAliveMs }); - const response = await transport.handleRequest(req('GET', { headers: withSession(sessionId) })); - expect(response.status).toBe(200); + const response = await transport.handleRequest(req('GET', { headers: withSession(sessionId) })); + expect(response.status).toBe(200); - // No timer may be armed: setInterval(fn, NaN) would be clamped by - // Node to ~1ms and flood the stream with keep-alive frames. - expect(vi.getTimerCount()).toBe(0); - const reader = response.body!.getReader(); - await vi.advanceTimersByTimeAsync(60000); - const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); - expect(raced).toBe('pending'); + // No timer may be armed: setInterval with a NaN/out-of-range delay + // is clamped by Node to ~1ms and would flood the stream with + // keep-alive frames. + expect(vi.getTimerCount()).toBe(0); + const reader = response.body!.getReader(); + await vi.advanceTimersByTimeAsync(60000); + const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); + expect(raced).toBe('pending'); - await transport.close(); - }); + await transport.close(); + } + ); it('should not arm keep-alive when the transport closes during an event-store replay await', async () => { let releaseReplay: (() => void) | undefined; @@ -3495,10 +3499,18 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { // Close the transport mid-await, then let the replay continuation run await transport.close(); releaseReplay?.(); - await pendingGet; + const replayResponse = await pendingGet; // The deferred continuation must not have armed a timer close() can never sweep expect(vi.getTimerCount()).toBe(0); + + // The continuation must not re-register the stream on the closed + // transport: the client observes stream end instead of hanging on a + // dead session, and a later resume isn't 409-blocked by a stale entry. + const { done } = await replayResponse.body!.getReader().read(); + expect(done).toBe(true); + const internals = transport as unknown as { _streamMapping: Map }; + expect(internals._streamMapping.size).toBe(0); }); it('should not leak a keep-alive timer when the priming event write fails on a POST SSE stream', async () => { @@ -3541,6 +3553,16 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { // The discarded stream must not carry a permanently-firing timer expect(vi.getTimerCount()).toBe(0); + // The stream bookkeeping registered before the failed priming write + // must be reclaimed too: repeated failures during an event-store + // outage must not accrete orphaned stream entries or request mappings. + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + expect(internals._streamMapping.size).toBe(0); + expect(internals._requestToStreamMapping.size).toBe(0); + await transport.close(); }); }); From 15a556f3a7dd2f24c5068111ce90e96d5c017688 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 12:53:33 +0100 Subject: [PATCH 4/6] fix: reject overflowing keep-alive intervals Reject timer delays above Node's maximum in addition to NaN and Infinity, and keep the lifecycle regression focused on cleanup rather than the transport's pre-existing error-status mapping. --- src/server/webStandardStreamableHttp.ts | 11 +++++------ test/server/streamableHttp.test.ts | 7 +++---- 2 files changed, 8 insertions(+), 10 deletions(-) diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index ca0843f796..12bf25561b 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -274,10 +274,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private startKeepAlive(streamId: string, controller: ReadableStreamDefaultController, encoder: TextEncoder): void { // A deferred arm (e.g. after an event-store await that straddled // close()) must not outlive the transport: close()'s timer sweep has - // already run. Non-finite intervals (NaN, Infinity) and non-positive - // ones disable keep-alive: setInterval would clamp them to ~1ms and - // flood every stream with comment frames. - if (!Number.isFinite(this._keepAliveMs) || this._keepAliveMs <= 0 || this._closed) { + // already run. Invalid timer delays disable keep-alive rather than + // letting setInterval clamp them to ~1ms and flood every stream. + if (!Number.isFinite(this._keepAliveMs) || this._keepAliveMs <= 0 || this._keepAliveMs > 2_147_483_647 || this._closed) { return; } this.stopKeepAlive(streamId); @@ -681,8 +680,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private async handlePostRequest(req: Request, options?: HandleRequestOptions): Promise { // Set once the SSE stream bookkeeping has been registered, so the // catch below can reclaim it: an error after registration (a failed - // priming event write, a throwing message handler) returns 400 and - // discards the Response, leaving nothing that could ever cancel the + // priming event write, a throwing message handler) returns an error + // response, leaving nothing that could ever cancel the discarded // stream or retire the request mappings. let reclaimSseBookkeeping: (() => void) | undefined; try { diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index 85752e1079..3f217ec6db 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -3452,8 +3452,8 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { await transport.close(); }); - it.each([Number.NaN, Number.POSITIVE_INFINITY])( - 'should disable keep-alive for non-finite keepAliveMs %s instead of arming a clamped interval', + it.each([Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( + 'should disable keep-alive for invalid keepAliveMs %s instead of arming a clamped interval', async keepAliveMs => { const { transport, sessionId } = await createTransport({ keepAliveMs }); @@ -3542,13 +3542,12 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { expect(vi.getTimerCount()).toBe(0); storeFails = true; - const response = await transport.handleRequest( + await transport.handleRequest( req('POST', { body: { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'call-1' }, headers: withSession(sessionId) }) ); - expect(response.status).toBe(400); // The discarded stream must not carry a permanently-firing timer expect(vi.getTimerCount()).toBe(0); From cca65fc7feb2d2883c1dea51539ae172b1af0730 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 15:37:38 +0100 Subject: [PATCH 5/6] fix: align SSE keep-alive delivery safeguards Reject sub-millisecond intervals that Node clamps to 1ms and disable nginx-style buffering on GET, POST, and replay SSE responses so heartbeat frames reach clients. --- .changeset/keepalive-lifecycle-hardening.md | 2 +- src/server/webStandardStreamableHttp.ts | 14 +++++++++----- test/server/streamableHttp.test.ts | 5 ++++- 3 files changed, 14 insertions(+), 7 deletions(-) diff --git a/.changeset/keepalive-lifecycle-hardening.md b/.changeset/keepalive-lifecycle-hardening.md index 39fd2d6278..10563be215 100644 --- a/.changeset/keepalive-lifecycle-hardening.md +++ b/.changeset/keepalive-lifecycle-hardening.md @@ -2,4 +2,4 @@ '@modelcontextprotocol/sdk': patch --- -Hardens the Streamable HTTP server transport's SSE keep-alive lifecycle: keep-alive timers can no longer be armed after the transport closes or leak when a priming event write fails, and a non-finite `keepAliveMs` disables keep-alive instead of arming a clamped ~1ms interval. +Hardens the Streamable HTTP server transport's SSE keep-alive lifecycle: timers can no longer be armed after transport close or leak when a priming event write fails, invalid timer delays safely disable keep-alive, and SSE responses disable nginx-style proxy buffering. diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index 12bf25561b..db6b841dcd 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -156,7 +156,8 @@ export interface WebStandardStreamableHTTPServerTransportOptions { * * Comment frames are ignored by SSE parsers and never surface as messages. * Defaults to 15000 (per the WHATWG SSE spec recommendation of roughly every - * 15 seconds). Set to 0 to disable keep-alive frames. + * 15 seconds). Set to 0 to disable keep-alive frames; values below 1, above + * 2147483647, or non-finite values also disable the timer. */ keepAliveMs?: number; } @@ -276,7 +277,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // close()) must not outlive the transport: close()'s timer sweep has // already run. Invalid timer delays disable keep-alive rather than // letting setInterval clamp them to ~1ms and flood every stream. - if (!Number.isFinite(this._keepAliveMs) || this._keepAliveMs <= 0 || this._keepAliveMs > 2_147_483_647 || this._closed) { + if (!Number.isFinite(this._keepAliveMs) || this._keepAliveMs < 1 || this._keepAliveMs > 2_147_483_647 || this._closed) { return; } this.stopKeepAlive(streamId); @@ -493,7 +494,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', - Connection: 'keep-alive' + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; // After initialization, always include the session ID if we have one @@ -552,7 +554,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', - Connection: 'keep-alive' + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; if (this.sessionId !== undefined) { @@ -839,7 +842,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', - Connection: 'keep-alive' + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; // After initialization, always include the session ID if we have one diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index 3f217ec6db..392cb7400f 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -270,6 +270,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { expect(response.status).toBe(200); expect(response.headers.get('content-type')).toBe('text/event-stream'); + expect(response.headers.get('x-accel-buffering')).toBe('no'); expect(response.headers.get('mcp-session-id')).toBeDefined(); }); @@ -486,6 +487,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { expect(sseResponse.status).toBe(200); expect(sseResponse.headers.get('content-type')).toBe('text/event-stream'); + expect(sseResponse.headers.get('x-accel-buffering')).toBe('no'); // Send a notification (server-initiated message) that should appear on SSE stream const notification: JSONRPCMessage = { @@ -1444,6 +1446,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { }); expect(reconnectResponse.status).toBe(200); + expect(reconnectResponse.headers.get('x-accel-buffering')).toBe('no'); // Read the replayed notification const reconnectReader = reconnectResponse.body?.getReader(); @@ -3452,7 +3455,7 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { await transport.close(); }); - it.each([Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( + it.each([0.5, Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( 'should disable keep-alive for invalid keepAliveMs %s instead of arming a clamped interval', async keepAliveMs => { const { transport, sessionId } = await createTransport({ keepAliveMs }); From 64799602e60ab407fda70802abbb597ec1ddd446 Mon Sep 17 00:00:00 2001 From: Matt Carey Date: Fri, 24 Jul 2026 15:58:08 +0100 Subject: [PATCH 6/6] fix: preserve POST lifecycle ownership Reject work resuming after transport close, ownership-guard request-map cleanup, and align no-transform response headers across GET, POST, and replay streams. --- .changeset/keepalive-lifecycle-hardening.md | 2 +- src/server/webStandardStreamableHttp.ts | 14 +++- test/server/streamableHttp.test.ts | 85 +++++++++++++++++++++ 3 files changed, 98 insertions(+), 3 deletions(-) diff --git a/.changeset/keepalive-lifecycle-hardening.md b/.changeset/keepalive-lifecycle-hardening.md index 10563be215..c7cf2940fe 100644 --- a/.changeset/keepalive-lifecycle-hardening.md +++ b/.changeset/keepalive-lifecycle-hardening.md @@ -2,4 +2,4 @@ '@modelcontextprotocol/sdk': patch --- -Hardens the Streamable HTTP server transport's SSE keep-alive lifecycle: timers can no longer be armed after transport close or leak when a priming event write fails, invalid timer delays safely disable keep-alive, and SSE responses disable nginx-style proxy buffering. +Hardens the Streamable HTTP server transport's SSE lifecycle: deferred work cannot register streams after transport close, error cleanup preserves successor request mappings, invalid timer delays safely disable keep-alive, and SSE responses disable proxy buffering. diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index db6b841dcd..f8946f8778 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -382,6 +382,10 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Returns a Response object (Web Standard) */ async handleRequest(req: Request, options?: HandleRequestOptions): Promise { + if (this._closed) { + return this.createJsonErrorResponse(404, -32001, 'Session not found'); + } + // In stateless mode (no sessionIdGenerator), each request must use a fresh transport. // Reusing a stateless transport causes message ID collisions between clients. if (!this.sessionIdGenerator && this._hasHandledRequest) { @@ -779,6 +783,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + // Request parsing and session initialization may await user/runtime + // work. Do not register or dispatch after close() has swept state. + if (this._closed) { + return this.createJsonErrorResponse(404, -32001, 'Session not found'); + } + // check if it contains requests const hasRequests = messages.some(isJSONRPCRequest); @@ -841,7 +851,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', - 'Cache-Control': 'no-cache', + 'Cache-Control': 'no-cache, no-transform', Connection: 'keep-alive', 'X-Accel-Buffering': 'no' }; @@ -875,7 +885,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { reclaimSseBookkeeping = () => { this._streamMapping.get(streamId)?.cleanup(); for (const message of messages) { - if (isJSONRPCRequest(message)) { + if (isJSONRPCRequest(message) && this._requestToStreamMapping.get(message.id) === streamId) { this._requestToStreamMapping.delete(message.id); } } diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index 392cb7400f..ac4a457248 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -270,6 +270,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { expect(response.status).toBe(200); expect(response.headers.get('content-type')).toBe('text/event-stream'); + expect(response.headers.get('cache-control')).toBe('no-cache, no-transform'); expect(response.headers.get('x-accel-buffering')).toBe('no'); expect(response.headers.get('mcp-session-id')).toBeDefined(); }); @@ -487,6 +488,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { expect(sseResponse.status).toBe(200); expect(sseResponse.headers.get('content-type')).toBe('text/event-stream'); + expect(sseResponse.headers.get('cache-control')).toBe('no-cache, no-transform'); expect(sseResponse.headers.get('x-accel-buffering')).toBe('no'); // Send a notification (server-initiated message) that should appear on SSE stream @@ -1446,6 +1448,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => { }); expect(reconnectResponse.status).toBe(200); + expect(reconnectResponse.headers.get('cache-control')).toBe('no-cache, no-transform'); expect(reconnectResponse.headers.get('x-accel-buffering')).toBe('no'); // Read the replayed notification @@ -3567,4 +3570,86 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { await transport.close(); }); + + it('should not reclaim a successor mapping that reused the failed POST request id', async () => { + let failPriming = false; + let rejectPriming: (() => void) | undefined; + const eventStore: EventStore = { + async storeEvent(): Promise { + if (failPriming) { + return new Promise((_resolve, reject) => { + rejectPriming = () => reject(new Error('event store unavailable')); + }); + } + return `evt-${randomUUID()}`; + }, + async replayEventsAfter(): Promise { + return 'stream-1'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + mcpServer.registerTool('noop', { description: 'noop' }, async () => ({ content: [] })); + await mcpServer.connect(transport); + + const initResponse = await transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize })); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + await vi.advanceTimersByTimeAsync(0); + failPriming = true; + + const pending = transport.handleRequest( + req('POST', { + body: { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'same-id' }, + headers: withSession(sessionId) + }) + ); + await vi.advanceTimersByTimeAsync(0); + expect(rejectPriming).toBeDefined(); + + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + const failedStreamId = internals._requestToStreamMapping.get('same-id'); + expect(failedStreamId).toBeDefined(); + internals._requestToStreamMapping.set('same-id', 'successor-stream'); + + rejectPriming?.(); + await pending; + + expect(internals._streamMapping.has(failedStreamId!)).toBe(false); + expect(internals._requestToStreamMapping.get('same-id')).toBe('successor-stream'); + internals._requestToStreamMapping.delete('same-id'); + await transport.close(); + }); + + it('should not register a POST stream after close races session initialization', async () => { + let releaseInitialization: (() => void) | undefined; + const transport = new WebStandardStreamableHTTPServerTransport({ + sessionIdGenerator: () => randomUUID(), + onsessioninitialized: async () => { + await new Promise(resolve => { + releaseInitialization = resolve; + }); + } + }); + await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport); + + const pendingInit = transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize })); + await vi.advanceTimersByTimeAsync(0); + expect(releaseInitialization).toBeDefined(); + + await transport.close(); + releaseInitialization?.(); + const response = await pendingInit; + + expect(response.status).toBe(404); + expect(vi.getTimerCount()).toBe(0); + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + expect(internals._streamMapping.size).toBe(0); + expect(internals._requestToStreamMapping.size).toBe(0); + }); });