Skip to content

Commit 9505f93

Browse files
committed
fix: send SSE keep-alive comment frames from WebStandardStreamableHTTPServerTransport
Idle SSE streams (the standalone GET stream in particular, but also POST response streams during long-running tool calls) are killed by intermediaries and server idle timeouts, which clients observe as "SSE stream disconnected: TypeError: terminated" followed by a reconnect loop. The transport now writes an SSE comment frame (`: keepalive`) to every open SSE stream every keepAliveMs milliseconds (default 15000, per the WHATWG SSE spec recommendation; set 0 to disable). Comment frames are dropped by SSE parsers and never surface as protocol messages. The timer is unref'd so it never holds the process open, and is cleared on stream cleanup/cancel and transport close. Naming matches the existing keepAliveMs on createMcpHandler's subscriptions/listen streams. v2 port of #2538 (v1.x); refs #1211
1 parent 1e1392e commit 9505f93

3 files changed

Lines changed: 201 additions & 0 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@modelcontextprotocol/server': patch
3+
---
4+
5+
`WebStandardStreamableHTTPServerTransport` now writes SSE keep-alive comment frames (`: keepalive`) to open SSE streams so idle connections (e.g. the standalone GET stream, or a POST stream during a long-running tool call) are not killed by intermediaries or server idle timeouts. Configurable via the new `keepAliveMs` option (default 15000; set 0 to disable).

packages/server/src/server/streamableHttp.ts

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,19 @@ export interface WebStandardStreamableHTTPServerTransportOptions {
148148
*/
149149
retryInterval?: number;
150150

151+
/**
152+
* Interval in milliseconds between SSE keep-alive comment frames (`: keepalive`)
153+
* written to open SSE streams. Keep-alive frames prevent idle streams (e.g. the
154+
* standalone `GET` stream, or a `POST` stream during a long-running tool call)
155+
* from being killed by intermediaries and server idle timeouts, which clients
156+
* observe as `SSE stream disconnected: TypeError: terminated`.
157+
*
158+
* Comment frames are ignored by SSE parsers and never surface as messages.
159+
* Defaults to `15000` (per the WHATWG SSE spec recommendation of roughly every
160+
* 15 seconds). Set to `0` to disable keep-alive frames.
161+
*/
162+
keepAliveMs?: number;
163+
151164
/**
152165
* List of protocol versions that this transport will accept.
153166
* Used to validate the `mcp-protocol-version` header in incoming requests.
@@ -161,6 +174,9 @@ export interface WebStandardStreamableHTTPServerTransportOptions {
161174
supportedProtocolVersions?: string[];
162175
}
163176

177+
/** Default interval between SSE keep-alive comment frames. */
178+
const DEFAULT_KEEP_ALIVE_MS = 15_000;
179+
164180
/**
165181
* Options for handling a request
166182
*/
@@ -247,6 +263,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
247263
private _enableDnsRebindingProtection: boolean;
248264
private _retryInterval?: number;
249265
private _supportedProtocolVersions: string[];
266+
private _keepAliveMs: number;
267+
private _keepAliveTimers: Map<string, ReturnType<typeof setInterval>> = new Map();
250268

251269
sessionId?: string;
252270
onclose?: () => void;
@@ -264,6 +282,47 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
264282
this._enableDnsRebindingProtection = options.enableDnsRebindingProtection ?? false;
265283
this._retryInterval = options.retryInterval;
266284
this._supportedProtocolVersions = options.supportedProtocolVersions ?? SUPPORTED_PROTOCOL_VERSIONS;
285+
this._keepAliveMs = options.keepAliveMs ?? DEFAULT_KEEP_ALIVE_MS;
286+
}
287+
288+
/**
289+
* Arms a keep-alive interval for an SSE stream that periodically writes an SSE
290+
* comment frame so intermediaries and idle timeouts don't kill the connection.
291+
* Replaces any timer already armed for the same stream id (a resumed stream
292+
* re-registered under the same id supersedes its predecessor's timer). The
293+
* timer is cleared via {@linkcode stopKeepAlive} when the stream is cleaned up,
294+
* and clears itself if a write fails (stream already closed/cancelled).
295+
*/
296+
private startKeepAlive(
297+
streamId: string,
298+
controller: ReadableStreamDefaultController<Uint8Array>,
299+
encoder: InstanceType<typeof TextEncoder>
300+
): void {
301+
if (this._keepAliveMs <= 0) {
302+
return;
303+
}
304+
this.stopKeepAlive(streamId);
305+
const timer = setInterval(() => {
306+
try {
307+
controller.enqueue(encoder.encode(': keepalive\n\n'));
308+
} catch {
309+
this.stopKeepAlive(streamId);
310+
}
311+
}, this._keepAliveMs);
312+
// Don't let the keep-alive timer hold the process open (Node.js only)
313+
(timer as { unref?: () => void }).unref?.();
314+
this._keepAliveTimers.set(streamId, timer);
315+
}
316+
317+
/**
318+
* Clears the keep-alive interval for a stream, if one is armed.
319+
*/
320+
private stopKeepAlive(streamId: string): void {
321+
const timer = this._keepAliveTimers.get(streamId);
322+
if (timer !== undefined) {
323+
clearInterval(timer);
324+
this._keepAliveTimers.delete(streamId);
325+
}
267326
}
268327

269328
/**
@@ -473,6 +532,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
473532
// it still points at THIS controller — a stale cancel must not
474533
// delete a successor stream registered by a later GET/resume.
475534
if (this._streamMapping.get(this._standaloneSseStreamId)?.controller === streamController) {
535+
this.stopKeepAlive(this._standaloneSseStreamId);
476536
this._streamMapping.delete(this._standaloneSseStreamId);
477537
}
478538
}
@@ -494,6 +554,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
494554
controller: streamController!,
495555
encoder,
496556
cleanup: () => {
557+
this.stopKeepAlive(this._standaloneSseStreamId);
497558
this._streamMapping.delete(this._standaloneSseStreamId);
498559
try {
499560
streamController!.close();
@@ -503,6 +564,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
503564
}
504565
});
505566

567+
this.startKeepAlive(this._standaloneSseStreamId, streamController!, encoder);
568+
506569
return new Response(readable, { headers });
507570
}
508571

@@ -564,6 +627,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
564627
// a stale cancel from an earlier resume must not delete a
565628
// successor resumed stream a re-poll has since registered.
566629
if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) {
630+
this.stopKeepAlive(replayedStreamId);
567631
this._streamMapping.delete(replayedStreamId);
568632
}
569633
}
@@ -590,6 +654,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
590654
encoder,
591655
replayedEventIds,
592656
cleanup: () => {
657+
this.stopKeepAlive(replayedStreamId!);
593658
this._streamMapping.delete(replayedStreamId!);
594659
try {
595660
streamController!.close();
@@ -618,6 +683,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
618683
}
619684
}
620685

686+
// Only arm keep-alive if the stream is still registered — the
687+
// no-in-flight-request path above may have already closed it.
688+
if (this._streamMapping.get(replayedStreamId)?.controller === streamController!) {
689+
this.startKeepAlive(replayedStreamId, streamController!, encoder);
690+
}
691+
621692
return new Response(readable, { headers });
622693
} catch (error) {
623694
this.onerror?.(error as Error);
@@ -830,6 +901,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
830901
// resumed stream under the same streamId) must not delete
831902
// the successor.
832903
if (this._streamMapping.get(streamId)?.controller === streamController) {
904+
this.stopKeepAlive(streamId);
833905
this._streamMapping.delete(streamId);
834906
}
835907
}
@@ -854,6 +926,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
854926
controller: streamController!,
855927
encoder,
856928
cleanup: () => {
929+
this.stopKeepAlive(streamId);
857930
this._streamMapping.delete(streamId);
858931
try {
859932
streamController!.close();
@@ -866,6 +939,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
866939
}
867940
}
868941

942+
this.startKeepAlive(streamId, streamController!, encoder);
943+
869944
// Write priming event if event store is configured (after mapping is set up)
870945
await this.writePrimingEvent(streamController!, encoder, streamId, clientProtocolVersion);
871946

@@ -986,6 +1061,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
9861061
}
9871062
this._streamMapping.clear();
9881063

1064+
// Clear any keep-alive timers not already cleared by stream cleanup
1065+
for (const timer of this._keepAliveTimers.values()) {
1066+
clearInterval(timer);
1067+
}
1068+
this._keepAliveTimers.clear();
1069+
9891070
// Clear any pending responses
9901071
this._requestResponseMap.clear();
9911072
this.onclose?.();

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

Lines changed: 115 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1407,3 +1407,118 @@ describe('Zod v4', () => {
14071407
});
14081408
});
14091409
});
1410+
1411+
describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => {
1412+
async function createTransport(options?: { keepAliveMs?: number }): Promise<{
1413+
transport: WebStandardStreamableHTTPServerTransport;
1414+
sessionId: string;
1415+
}> {
1416+
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), ...options });
1417+
await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport);
1418+
const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize));
1419+
expect(initResponse.status).toBe(200);
1420+
return { transport, sessionId: initResponse.headers.get('mcp-session-id') as string };
1421+
}
1422+
1423+
beforeEach(() => {
1424+
vi.useFakeTimers();
1425+
});
1426+
1427+
afterEach(() => {
1428+
vi.useRealTimers();
1429+
});
1430+
1431+
it('should write keep-alive comment frames to an idle standalone GET stream', async () => {
1432+
const { transport, sessionId } = await createTransport();
1433+
1434+
const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId }));
1435+
expect(response.status).toBe(200);
1436+
1437+
const reader = response.body!.getReader();
1438+
await vi.advanceTimersByTimeAsync(15000);
1439+
const { value } = await reader.read();
1440+
expect(new TextDecoder().decode(value)).toBe(': keepalive\n\n');
1441+
1442+
await transport.close();
1443+
});
1444+
1445+
it('should honor a custom keepAliveMs interval', async () => {
1446+
const { transport, sessionId } = await createTransport({ keepAliveMs: 1000 });
1447+
1448+
const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId }));
1449+
const reader = response.body!.getReader();
1450+
1451+
await vi.advanceTimersByTimeAsync(3000);
1452+
let received = '';
1453+
for (let i = 0; i < 3; i++) {
1454+
const { value } = await reader.read();
1455+
received += new TextDecoder().decode(value);
1456+
}
1457+
expect(received).toBe(': keepalive\n\n'.repeat(3));
1458+
1459+
await transport.close();
1460+
});
1461+
1462+
it('should not write keep-alive frames when keepAliveMs is 0', async () => {
1463+
const { transport, sessionId } = await createTransport({ keepAliveMs: 0 });
1464+
1465+
const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId }));
1466+
const reader = response.body!.getReader();
1467+
1468+
await vi.advanceTimersByTimeAsync(60000);
1469+
const raced = await Promise.race([reader.read(), Promise.resolve('pending')]);
1470+
expect(raced).toBe('pending');
1471+
1472+
await transport.close();
1473+
});
1474+
1475+
it('should stop keep-alive frames after the stream is closed', async () => {
1476+
const { transport, sessionId } = await createTransport();
1477+
1478+
const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId }));
1479+
const reader = response.body!.getReader();
1480+
1481+
await transport.close();
1482+
const { done } = await reader.read();
1483+
expect(done).toBe(true);
1484+
1485+
// Advancing time after close must not throw or fire further writes
1486+
expect(vi.getTimerCount()).toBe(0);
1487+
await vi.advanceTimersByTimeAsync(60000);
1488+
});
1489+
1490+
it('should write keep-alive frames on a POST SSE stream while a request is pending', async () => {
1491+
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID() });
1492+
const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' });
1493+
let resolveTool: (() => void) | undefined;
1494+
mcpServer.registerTool('slow', { description: 'never resolves until released' }, async (): Promise<CallToolResult> => {
1495+
await new Promise<void>(resolve => {
1496+
resolveTool = resolve;
1497+
});
1498+
return { content: [{ type: 'text', text: 'done' }] };
1499+
});
1500+
await mcpServer.connect(transport);
1501+
1502+
const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize));
1503+
const sessionId = initResponse.headers.get('mcp-session-id') as string;
1504+
1505+
const response = await transport.handleRequest(
1506+
createRequest(
1507+
'POST',
1508+
{ jsonrpc: '2.0', method: 'tools/call', params: { name: 'slow', arguments: {} }, id: 'call-1' } as JSONRPCMessage,
1509+
{
1510+
sessionId
1511+
}
1512+
)
1513+
);
1514+
expect(response.status).toBe(200);
1515+
const reader = response.body!.getReader();
1516+
1517+
await vi.advanceTimersByTimeAsync(15000);
1518+
const { value } = await reader.read();
1519+
expect(new TextDecoder().decode(value)).toBe(': keepalive\n\n');
1520+
1521+
resolveTool?.();
1522+
await transport.close();
1523+
});
1524+
});

0 commit comments

Comments
 (0)