Skip to content

Commit 2090bb9

Browse files
committed
fix: replay request SSE responses after closeSSEStream
1 parent 5fc42e9 commit 2090bb9

3 files changed

Lines changed: 167 additions & 12 deletions

File tree

packages/middleware/node/test/streamableHttp.test.ts

Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1954,6 +1954,90 @@ describe('Zod v4', () => {
19541954
toolResolve!();
19551955
});
19561956

1957+
it('should replay the terminal POST SSE response after ctx.http?.closeSSE closes the request stream', async () => {
1958+
const result = await createTestServer({
1959+
sessionIdGenerator: () => randomUUID(),
1960+
eventStore: createEventStore(),
1961+
retryInterval: 1000
1962+
});
1963+
server = result.server;
1964+
transport = result.transport;
1965+
baseUrl = result.baseUrl;
1966+
mcpServer = result.mcpServer;
1967+
1968+
mcpServer.registerTool('close-and-complete', { description: 'Closes request stream and completes later' }, async ctx => {
1969+
ctx.http?.closeSSE?.();
1970+
return {
1971+
content: [{ type: 'text', text: 'Done after reconnect' }]
1972+
};
1973+
});
1974+
1975+
const initResponse = await sendPostRequest(baseUrl, TEST_MESSAGES.initialize);
1976+
sessionId = initResponse.headers.get('mcp-session-id') as string;
1977+
expect(sessionId).toBeDefined();
1978+
1979+
const toolCallRequest: JSONRPCMessage = {
1980+
jsonrpc: '2.0',
1981+
id: 101,
1982+
method: 'tools/call',
1983+
params: { name: 'close-and-complete', arguments: {} }
1984+
};
1985+
1986+
const postResponse = await fetch(baseUrl, {
1987+
method: 'POST',
1988+
headers: {
1989+
'Content-Type': 'application/json',
1990+
Accept: 'text/event-stream, application/json',
1991+
'mcp-session-id': sessionId,
1992+
'mcp-protocol-version': '2025-11-25'
1993+
},
1994+
body: JSON.stringify(toolCallRequest)
1995+
});
1996+
1997+
expect(postResponse.status).toBe(200);
1998+
1999+
const reader = postResponse.body?.getReader();
2000+
const { value } = await reader!.read();
2001+
const text = new TextDecoder().decode(value);
2002+
const idMatch = text.match(/id: ([^\n]+)/);
2003+
expect(idMatch).toBeTruthy();
2004+
const lastEventId = idMatch![1]!;
2005+
2006+
const closedRead = reader!.read();
2007+
const closedTimeout = new Promise<{ done: boolean; value: undefined }>((_, reject) =>
2008+
setTimeout(() => reject(new Error('POST SSE stream did not close in time')), 1000)
2009+
);
2010+
const { done } = await Promise.race([closedRead, closedTimeout]);
2011+
expect(done).toBe(true);
2012+
2013+
const reconnectResponse = await fetch(baseUrl, {
2014+
method: 'GET',
2015+
headers: {
2016+
Accept: 'text/event-stream',
2017+
'mcp-session-id': sessionId,
2018+
'mcp-protocol-version': '2025-11-25',
2019+
'last-event-id': lastEventId
2020+
}
2021+
});
2022+
expect(reconnectResponse.status).toBe(200);
2023+
2024+
const reconnectReader = reconnectResponse.body?.getReader();
2025+
let replayedText = '';
2026+
const replayTimeout = setTimeout(() => reconnectReader!.cancel(), 5000);
2027+
try {
2028+
while (!replayedText.includes('Done after reconnect')) {
2029+
const { value, done } = await reconnectReader!.read();
2030+
if (done) break;
2031+
replayedText += new TextDecoder().decode(value);
2032+
}
2033+
} finally {
2034+
clearTimeout(replayTimeout);
2035+
}
2036+
2037+
expect(replayedText).toContain('Done after reconnect');
2038+
expect(replayedText).toContain('"id":101');
2039+
});
2040+
19572041
it('should provide closeSSEStream callback in ctx when eventStore is configured', async () => {
19582042
const result = await createTestServer({
19592043
sessionIdGenerator: () => randomUUID(),

packages/server/src/server/streamableHttp.ts

Lines changed: 15 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -985,14 +985,15 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
985985

986986
const stream = this._streamMapping.get(streamId);
987987

988-
if (!this._enableJsonResponse && stream?.controller && stream?.encoder) {
989-
// For SSE responses, generate event ID if event store is provided
990-
let eventId: string | undefined;
988+
let eventId: string | undefined;
989+
if (!this._enableJsonResponse && this._eventStore) {
990+
// Persist request-scoped SSE events even while the controller is
991+
// temporarily disconnected so they can be replayed on reconnect.
992+
eventId = await this._eventStore.storeEvent(streamId, message);
993+
}
991994

992-
if (this._eventStore) {
993-
eventId = await this._eventStore.storeEvent(streamId, message);
994-
}
995-
// Write the event to the response stream
995+
if (!this._enableJsonResponse && stream?.controller && stream?.encoder) {
996+
// Write the event to the active response stream when available.
996997
this.writeSSEEvent(stream.controller, stream.encoder, message, eventId);
997998
}
998999

@@ -1004,10 +1005,10 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
10041005
const allResponsesReady = relatedIds.every(id => this._requestResponseMap.has(id));
10051006

10061007
if (allResponsesReady) {
1007-
if (!stream) {
1008-
throw new Error(`No connection established for request ID: ${String(requestId)}`);
1009-
}
1010-
if (this._enableJsonResponse && stream.resolveJson) {
1008+
if (this._enableJsonResponse) {
1009+
if (!stream?.resolveJson) {
1010+
throw new Error(`No connection established for request ID: ${String(requestId)}`);
1011+
}
10111012
// All responses ready, send as JSON
10121013
const headers: Record<string, string> = {
10131014
'Content-Type': 'application/json'
@@ -1023,9 +1024,11 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
10231024
} else {
10241025
stream.resolveJson(Response.json(responses, { status: 200, headers }));
10251026
}
1026-
} else {
1027+
} else if (stream) {
10271028
// End the SSE stream
10281029
stream.cleanup();
1030+
} else if (!this._eventStore) {
1031+
throw new Error(`No connection established for request ID: ${String(requestId)}`);
10291032
}
10301033
// Clean up
10311034
for (const id of relatedIds) {

packages/server/test/server/streamableHttp.test.ts

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -705,6 +705,74 @@ describe('Zod v4', () => {
705705
// Should have id: field in the SSE event
706706
expect(text).toContain('id:');
707707
});
708+
709+
it('should replay a terminal POST SSE response after closeSSEStream disconnects the request stream', async () => {
710+
sessionId = await initializeServer();
711+
712+
mcpServer.registerTool('close-and-complete', { description: 'Disconnects request SSE and completes later' }, async ctx => {
713+
ctx.http?.closeSSE?.();
714+
return {
715+
content: [{ type: 'text', text: 'Done after reconnect' }]
716+
};
717+
});
718+
719+
const request = createRequest(
720+
'POST',
721+
{
722+
jsonrpc: '2.0',
723+
id: 'tool-close-1',
724+
method: 'tools/call',
725+
params: { name: 'close-and-complete', arguments: {} }
726+
} as JSONRPCMessage,
727+
{ sessionId }
728+
);
729+
const response = await transport.handleRequest(request);
730+
731+
expect(response.status).toBe(200);
732+
733+
const reader = response.body?.getReader();
734+
const { value } = await reader!.read();
735+
const text = new TextDecoder().decode(value);
736+
const idMatch = text.match(/id: ([^\n]+)/);
737+
expect(idMatch).toBeTruthy();
738+
const lastEventId = idMatch![1]!;
739+
740+
const closedRead = reader!.read();
741+
const closedTimeout = new Promise<{ done: boolean; value: undefined }>((_, reject) =>
742+
setTimeout(() => reject(new Error('POST SSE stream did not close in time')), 1000)
743+
);
744+
const { done } = await Promise.race([closedRead, closedTimeout]);
745+
expect(done).toBe(true);
746+
747+
const reconnectResponse = await transport.handleRequest(
748+
createRequest('GET', undefined, {
749+
sessionId,
750+
extraHeaders: {
751+
'last-event-id': lastEventId
752+
}
753+
})
754+
);
755+
expect(reconnectResponse.status).toBe(200);
756+
757+
const reconnectReader = reconnectResponse.body?.getReader();
758+
let replayedText = '';
759+
const replayTimeout = setTimeout(() => reconnectReader!.cancel(), 2000);
760+
try {
761+
while (!replayedText.includes('Done after reconnect')) {
762+
const { value, done } = await reconnectReader!.read();
763+
if (done) break;
764+
replayedText += new TextDecoder().decode(value);
765+
}
766+
} finally {
767+
clearTimeout(replayTimeout);
768+
}
769+
770+
expect(replayedText).toContain('Done after reconnect');
771+
expect(parseSSEData(replayedText)).toMatchObject({
772+
jsonrpc: '2.0',
773+
id: 'tool-close-1'
774+
});
775+
});
708776
});
709777

710778
describe('HTTPServerTransport - Protocol Version Validation', () => {

0 commit comments

Comments
 (0)