Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/close-replayed-request-streams.md
Original file line number Diff line number Diff line change
@@ -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.
23 changes: 22 additions & 1 deletion src/server/webStandardStreamableHttp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
76 changes: 74 additions & 2 deletions test/server/streamableHttp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3420,7 +3420,7 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => {
return 'evt-1';
},
async replayEventsAfter(): Promise<StreamId> {
return 'stream-1';
return '_GET_stream';
}
};
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore });
Expand All @@ -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);

Expand Down Expand Up @@ -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<EventId> {
const id = `${streamId}#${counter++}`;
events.push({ id, streamId, message });
return id;
},
async getStreamIdForEventId(eventId: EventId): Promise<StreamId | undefined> {
return events.find(e => e.id === eventId)?.streamId;
},
async replayEventsAfter(lastEventId: EventId, { send }): Promise<StreamId> {
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;
Expand Down
Loading