From 3b11713e794f9c49d1aebd5cef6390a42b3546bf Mon Sep 17 00:00:00 2001 From: Ashwin Mridul Date: Thu, 10 Sep 2026 18:13:43 +0530 Subject: [PATCH] fix: close replayed request SSE streams when no request remains in flight After Last-Event-ID replay, keep-alive was holding the successor stream open even when the correlated request had already retired, so a later reconnect could be refused with 409. Close and unregister per-request streams in that case; leave the standalone GET stream open. Co-authored-by: Cursor --- .changeset/close-replayed-request-streams.md | 5 ++ src/server/webStandardStreamableHttp.ts | 23 +++++- test/server/streamableHttp.test.ts | 76 +++++++++++++++++++- 3 files changed, 101 insertions(+), 3 deletions(-) create mode 100644 .changeset/close-replayed-request-streams.md diff --git a/.changeset/close-replayed-request-streams.md b/.changeset/close-replayed-request-streams.md new file mode 100644 index 0000000000..6ceaff8c5a --- /dev/null +++ b/.changeset/close-replayed-request-streams.md @@ -0,0 +1,5 @@ +--- +'@modelcontextprotocol/sdk': patch +--- + +Close a Streamable HTTP request stream after Last-Event-ID replay when no in-flight request still maps to it, so resume polling converges instead of holding keep-alive open and blocking a later reconnect with 409. diff --git a/src/server/webStandardStreamableHttp.ts b/src/server/webStandardStreamableHttp.ts index cffca7bc0e..7159a1c4a9 100644 --- a/src/server/webStandardStreamableHttp.ts +++ b/src/server/webStandardStreamableHttp.ts @@ -653,7 +653,28 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // The resume itself proves the client holds a Last-Event-ID cursor this._resumableStreams.add(replayedStreamId); - keepAliveTimer = this.startKeepAlive(streamController!, encoder); + // If this is a per-request stream and no in-flight request still + // targets this streamId, the request was already retired by the + // clean-return path while disconnected and the replay above just + // delivered the final response. Per the spec the server SHOULD + // close the SSE stream after the JSON-RPC response — close and + // unregister so a later reconnect isn't refused with 409. The + // standalone GET stream is never request-scoped and stays open. + if (replayedStreamId !== this._standaloneSseStreamId) { + const hasInFlightRequest = [...this._requestToStreamMapping.values()].includes(replayedStreamId); + if (!hasInFlightRequest) { + this._streamMapping.delete(replayedStreamId); + try { + streamController!.close(); + } catch { + // Controller might already be closed + } + } + } + + if (this._streamMapping.get(replayedStreamId)?.controller === streamController!) { + keepAliveTimer = this.startKeepAlive(streamController!, encoder); + } return new Response(readable, { headers }); } catch (error) { diff --git a/test/server/streamableHttp.test.ts b/test/server/streamableHttp.test.ts index f402abeff5..ced27430f5 100644 --- a/test/server/streamableHttp.test.ts +++ b/test/server/streamableHttp.test.ts @@ -3420,7 +3420,7 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { return 'evt-1'; }, async replayEventsAfter(): Promise { - return 'stream-1'; + return '_GET_stream'; } }; const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); @@ -3432,7 +3432,7 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { const first = await transport.handleRequest(req('GET', { headers: replayHeaders })); expect(first.status).toBe(200); - // Reconnect with the same Last-Event-ID — re-registers 'stream-1' + // Reconnect with the same Last-Event-ID — re-registers '_GET_stream' const second = await transport.handleRequest(req('GET', { headers: replayHeaders })); expect(second.status).toBe(200); @@ -3982,6 +3982,78 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive lifecycle', () expect(vi.getTimerCount()).toBe(0); }); + it('should close and unregister the resumed request stream when reconnecting after the request was already retired', async () => { + // Retire-then-reconnect: closeSSEStream → tool result. With no live + // writer the final response is stored and the request id is retired. + // A subsequent Last-Event-ID reconnect must replay that response AND + // close the resumed stream so a second reconnect is not refused with 409. + const events: { id: string; streamId: string; message: JSONRPCMessage }[] = []; + let counter = 0; + const eventStore: EventStore = { + async storeEvent(streamId: StreamId, message: JSONRPCMessage): Promise { + const id = `${streamId}#${counter++}`; + events.push({ id, streamId, message }); + return id; + }, + async getStreamIdForEventId(eventId: EventId): Promise { + return events.find(e => e.id === eventId)?.streamId; + }, + async replayEventsAfter(lastEventId: EventId, { send }): Promise { + const index = events.findIndex(e => e.id === lastEventId); + const streamId = events[index]?.streamId ?? '_GET_stream'; + for (const event of events.slice(index + 1).filter(e => e.streamId === streamId)) { + await send(event.id, event.message); + } + return streamId; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + mcpServer.tool('retire', 'closeSSE then emit then return', {}, async (_args, extra) => { + extra.closeSSEStream?.(); + await extra.sendNotification({ + method: 'notifications/progress', + params: { progressToken: 'retire-1', progress: 75 } + }); + return { content: [{ type: 'text', text: 'done' }] }; + }); + 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; + + const original = await transport.handleRequest( + req('POST', { + body: { jsonrpc: '2.0', method: 'tools/call', params: { name: 'retire', arguments: {} }, id: 'retire-1' }, + headers: { 'mcp-session-id': sessionId, 'mcp-protocol-version': '2025-11-25' } + }) + ); + const primingText = await original.text().catch(() => ''); + const primingEventId = /^id: (.+)$/m.exec(primingText)?.[1]; + expect(primingEventId).toBeDefined(); + + await vi.advanceTimersByTimeAsync(0); + // eslint-disable-next-line @typescript-eslint/no-explicit-any + expect((transport as any)._requestToStreamMapping.has('retire-1')).toBe(false); + + const reconnect = await transport.handleRequest( + req('GET', { headers: { 'mcp-session-id': sessionId, 'mcp-protocol-version': '2025-11-25', 'Last-Event-ID': primingEventId! } }) + ); + expect(reconnect.status).toBe(200); + const replayed = await reconnect.text(); + expect(replayed).toContain('notifications/progress'); + expect(replayed).toContain('"id":"retire-1"'); + expect(replayed).toContain('"result"'); + + const reconnect2 = await transport.handleRequest( + req('GET', { headers: { 'mcp-session-id': sessionId, 'mcp-protocol-version': '2025-11-25', 'Last-Event-ID': primingEventId! } }) + ); + expect(reconnect2.status).toBe(200); + await reconnect2.body?.cancel(); + + await transport.close(); + expect(vi.getTimerCount()).toBe(0); + }); + it('should not re-deliver replayed server notifications to a resumed standalone stream', async () => { const events: { id: string; streamId: string; message: JSONRPCMessage }[] = []; let counter = 0;