Skip to content

Commit 6479960

Browse files
committed
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.
1 parent cca65fc commit 6479960

3 files changed

Lines changed: 98 additions & 3 deletions

File tree

.changeset/keepalive-lifecycle-hardening.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,4 +2,4 @@
22
'@modelcontextprotocol/sdk': patch
33
---
44

5-
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.
5+
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.

src/server/webStandardStreamableHttp.ts

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -382,6 +382,10 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
382382
* Returns a Response object (Web Standard)
383383
*/
384384
async handleRequest(req: Request, options?: HandleRequestOptions): Promise<Response> {
385+
if (this._closed) {
386+
return this.createJsonErrorResponse(404, -32001, 'Session not found');
387+
}
388+
385389
// In stateless mode (no sessionIdGenerator), each request must use a fresh transport.
386390
// Reusing a stateless transport causes message ID collisions between clients.
387391
if (!this.sessionIdGenerator && this._hasHandledRequest) {
@@ -779,6 +783,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
779783
}
780784
}
781785

786+
// Request parsing and session initialization may await user/runtime
787+
// work. Do not register or dispatch after close() has swept state.
788+
if (this._closed) {
789+
return this.createJsonErrorResponse(404, -32001, 'Session not found');
790+
}
791+
782792
// check if it contains requests
783793
const hasRequests = messages.some(isJSONRPCRequest);
784794

@@ -841,7 +851,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
841851

842852
const headers: Record<string, string> = {
843853
'Content-Type': 'text/event-stream',
844-
'Cache-Control': 'no-cache',
854+
'Cache-Control': 'no-cache, no-transform',
845855
Connection: 'keep-alive',
846856
'X-Accel-Buffering': 'no'
847857
};
@@ -875,7 +885,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
875885
reclaimSseBookkeeping = () => {
876886
this._streamMapping.get(streamId)?.cleanup();
877887
for (const message of messages) {
878-
if (isJSONRPCRequest(message)) {
888+
if (isJSONRPCRequest(message) && this._requestToStreamMapping.get(message.id) === streamId) {
879889
this._requestToStreamMapping.delete(message.id);
880890
}
881891
}

test/server/streamableHttp.test.ts

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -270,6 +270,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => {
270270

271271
expect(response.status).toBe(200);
272272
expect(response.headers.get('content-type')).toBe('text/event-stream');
273+
expect(response.headers.get('cache-control')).toBe('no-cache, no-transform');
273274
expect(response.headers.get('x-accel-buffering')).toBe('no');
274275
expect(response.headers.get('mcp-session-id')).toBeDefined();
275276
});
@@ -487,6 +488,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => {
487488

488489
expect(sseResponse.status).toBe(200);
489490
expect(sseResponse.headers.get('content-type')).toBe('text/event-stream');
491+
expect(sseResponse.headers.get('cache-control')).toBe('no-cache, no-transform');
490492
expect(sseResponse.headers.get('x-accel-buffering')).toBe('no');
491493

492494
// Send a notification (server-initiated message) that should appear on SSE stream
@@ -1446,6 +1448,7 @@ describe.each(zodTestMatrix)('$zodVersionLabel', (entry: ZodMatrixEntry) => {
14461448
});
14471449

14481450
expect(reconnectResponse.status).toBe(200);
1451+
expect(reconnectResponse.headers.get('cache-control')).toBe('no-cache, no-transform');
14491452
expect(reconnectResponse.headers.get('x-accel-buffering')).toBe('no');
14501453

14511454
// Read the replayed notification
@@ -3567,4 +3570,86 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => {
35673570

35683571
await transport.close();
35693572
});
3573+
3574+
it('should not reclaim a successor mapping that reused the failed POST request id', async () => {
3575+
let failPriming = false;
3576+
let rejectPriming: (() => void) | undefined;
3577+
const eventStore: EventStore = {
3578+
async storeEvent(): Promise<EventId> {
3579+
if (failPriming) {
3580+
return new Promise<EventId>((_resolve, reject) => {
3581+
rejectPriming = () => reject(new Error('event store unavailable'));
3582+
});
3583+
}
3584+
return `evt-${randomUUID()}`;
3585+
},
3586+
async replayEventsAfter(): Promise<StreamId> {
3587+
return 'stream-1';
3588+
}
3589+
};
3590+
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore });
3591+
const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' });
3592+
mcpServer.registerTool('noop', { description: 'noop' }, async () => ({ content: [] }));
3593+
await mcpServer.connect(transport);
3594+
3595+
const initResponse = await transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize }));
3596+
const sessionId = initResponse.headers.get('mcp-session-id') as string;
3597+
await vi.advanceTimersByTimeAsync(0);
3598+
failPriming = true;
3599+
3600+
const pending = transport.handleRequest(
3601+
req('POST', {
3602+
body: { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'same-id' },
3603+
headers: withSession(sessionId)
3604+
})
3605+
);
3606+
await vi.advanceTimersByTimeAsync(0);
3607+
expect(rejectPriming).toBeDefined();
3608+
3609+
const internals = transport as unknown as {
3610+
_streamMapping: Map<string, unknown>;
3611+
_requestToStreamMapping: Map<unknown, string>;
3612+
};
3613+
const failedStreamId = internals._requestToStreamMapping.get('same-id');
3614+
expect(failedStreamId).toBeDefined();
3615+
internals._requestToStreamMapping.set('same-id', 'successor-stream');
3616+
3617+
rejectPriming?.();
3618+
await pending;
3619+
3620+
expect(internals._streamMapping.has(failedStreamId!)).toBe(false);
3621+
expect(internals._requestToStreamMapping.get('same-id')).toBe('successor-stream');
3622+
internals._requestToStreamMapping.delete('same-id');
3623+
await transport.close();
3624+
});
3625+
3626+
it('should not register a POST stream after close races session initialization', async () => {
3627+
let releaseInitialization: (() => void) | undefined;
3628+
const transport = new WebStandardStreamableHTTPServerTransport({
3629+
sessionIdGenerator: () => randomUUID(),
3630+
onsessioninitialized: async () => {
3631+
await new Promise<void>(resolve => {
3632+
releaseInitialization = resolve;
3633+
});
3634+
}
3635+
});
3636+
await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport);
3637+
3638+
const pendingInit = transport.handleRequest(req('POST', { body: TEST_MESSAGES.initialize }));
3639+
await vi.advanceTimersByTimeAsync(0);
3640+
expect(releaseInitialization).toBeDefined();
3641+
3642+
await transport.close();
3643+
releaseInitialization?.();
3644+
const response = await pendingInit;
3645+
3646+
expect(response.status).toBe(404);
3647+
expect(vi.getTimerCount()).toBe(0);
3648+
const internals = transport as unknown as {
3649+
_streamMapping: Map<string, unknown>;
3650+
_requestToStreamMapping: Map<unknown, string>;
3651+
};
3652+
expect(internals._streamMapping.size).toBe(0);
3653+
expect(internals._requestToStreamMapping.size).toBe(0);
3654+
});
35703655
});

0 commit comments

Comments
 (0)