Skip to content

Commit 6d20a1f

Browse files
committed
fix: address keep-alive lifecycle review findings
- startKeepAlive is a no-op after transport close: a deferred arm (e.g. resuming after an event-store replay await that straddled close()) would otherwise create a timer that close()'s sweep can never clear. - The POST SSE path arms keep-alive after the fallible awaits (priming event write, message dispatch) instead of before them, so an error path that discards the Response cannot leak a permanently-firing timer against a stream nothing can cancel. - createMcpHandler's keepAliveMs now also reaches the legacy stateless fallback's per-request transport, instead of governing only subscriptions/listen streams while the legacy leg silently used the transport default. - Documents the keep-alive behavior and the 'SSE stream disconnected: TypeError: terminated' symptom in docs/troubleshooting.md.
1 parent 9505f93 commit 6d20a1f

4 files changed

Lines changed: 127 additions & 8 deletions

File tree

docs/troubleshooting.md

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -154,6 +154,14 @@ Rewrite the imports:
154154

155155
The Resource Server helpers did not move there: `requireBearerAuth`, `mcpAuthMetadataRouter` and `OAuthTokenVerifier` are first-class in `@modelcontextprotocol/express` — see [Authorization](./serving/authorization.md). `@modelcontextprotocol/server-legacy` is frozen and receives no new features; serve new code over [Streamable HTTP](./serving/http.md), which still reaches 2025-era clients through [legacy client support](./serving/legacy-clients.md). A client limited to the HTTP+SSE transport is the one case that still needs the frozen `@modelcontextprotocol/server-legacy/sse` import above.
156156

157+
## `SSE stream disconnected: TypeError: terminated`
158+
159+
An idle SSE stream was killed by an intermediary or an idle-connection timeout — Node's `server.requestTimeout` defaults to 300 seconds, and reverse proxies and cloud load balancers have similar watchdogs. The client observes the dropped socket as this error (typically every ~5 minutes) and reconnects in a loop.
160+
161+
`WebStandardStreamableHTTPServerTransport` prevents this by writing an SSE comment frame (`: keepalive`) to every open SSE stream every 15 seconds by default. Comment frames are dropped by SSE parsers before event dispatch, so they never surface as protocol messages. Tune or disable the interval with the transport's `keepAliveMs` option (`0` disables); `createMcpHandler`'s `keepAliveMs` option covers both its `subscriptions/listen` streams and the legacy fallback's per-request transport.
162+
163+
If you still see this error, either keep-alive is disabled (`keepAliveMs: 0`) or an intermediary between client and server buffers or strips SSE data — check for proxies that buffer streaming responses (e.g. nginx without `proxy_buffering off`).
164+
157165
## Recap
158166

159167
- Every heading on this page is the exact message you searched for.

packages/server/src/server/createMcpHandler.ts

Lines changed: 18 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -194,8 +194,9 @@ export interface CreateMcpHandlerOptions {
194194
*/
195195
maxSubscriptions?: number;
196196
/**
197-
* SSE comment-frame keepalive interval for `subscriptions/listen` streams,
198-
* in milliseconds. Set to `0` to disable.
197+
* SSE comment-frame keepalive interval, in milliseconds, applied to
198+
* `subscriptions/listen` streams and to the SSE streams of the legacy
199+
* stateless fallback's per-request transport. Set to `0` to disable.
199200
* @default 15000
200201
*/
201202
keepAliveMs?: number;
@@ -306,7 +307,11 @@ function internalServerErrorResponse(id: RequestId | null = null): Response {
306307
* The entry passes its own `onerror` here when expanding the default, so
307308
* legacy-leg failures are never silently swallowed.
308309
*/
309-
export function legacyStatelessFallback(factory: McpServerFactory, onerror?: (error: Error) => void): LegacyHttpHandler {
310+
export function legacyStatelessFallback(
311+
factory: McpServerFactory,
312+
onerror?: (error: Error) => void,
313+
transportOptions?: { keepAliveMs?: number }
314+
): LegacyHttpHandler {
310315
return async (request, options) => {
311316
if (request.method.toUpperCase() !== 'POST') {
312317
return jsonRpcErrorResponse(405, -32_000, 'Method not allowed.');
@@ -317,7 +322,10 @@ export function legacyStatelessFallback(factory: McpServerFactory, onerror?: (er
317322
...(options?.authInfo !== undefined && { authInfo: options.authInfo }),
318323
requestInfo: request
319324
});
320-
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: undefined });
325+
const transport = new WebStandardStreamableHTTPServerTransport({
326+
sessionIdGenerator: undefined,
327+
...(transportOptions?.keepAliveMs !== undefined && { keepAliveMs: transportOptions.keepAliveMs })
328+
});
321329
await product.connect(transport);
322330

323331
const teardown = () => {
@@ -632,7 +640,12 @@ export function createMcpHandler(factory: McpServerFactory, options: CreateMcpHa
632640

633641
// The default posture is the stateless fallback; 'reject' is the only way
634642
// to turn legacy serving off (modern-only strict).
635-
const legacyHandler: LegacyHttpHandler | undefined = legacy === 'reject' ? undefined : legacyStatelessFallback(factory, reportError);
643+
const legacyHandler: LegacyHttpHandler | undefined =
644+
legacy === 'reject'
645+
? undefined
646+
: legacyStatelessFallback(factory, reportError, {
647+
...(options.keepAliveMs !== undefined && { keepAliveMs: options.keepAliveMs })
648+
});
636649

637650
async function serveModern(route: InboundModernRoute, request: Request, authInfo: AuthInfo | undefined): Promise<Response> {
638651
const claimedRevision = route.classification.revision;

packages/server/src/server/streamableHttp.ts

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -298,7 +298,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
298298
controller: ReadableStreamDefaultController<Uint8Array>,
299299
encoder: InstanceType<typeof TextEncoder>
300300
): void {
301-
if (this._keepAliveMs <= 0) {
301+
// A deferred arm (e.g. after an event-store await) must not outlive the
302+
// transport: close()'s timer sweep has already run and never runs again.
303+
if (this._keepAliveMs <= 0 || this._closed) {
302304
return;
303305
}
304306
this.stopKeepAlive(streamId);
@@ -939,8 +941,6 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
939941
}
940942
}
941943

942-
this.startKeepAlive(streamId, streamController!, encoder);
943-
944944
// Write priming event if event store is configured (after mapping is set up)
945945
await this.writePrimingEvent(streamController!, encoder, streamId, clientProtocolVersion);
946946

@@ -966,6 +966,14 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
966966
// The server SHOULD NOT close the SSE stream before sending all JSON-RPC responses
967967
// This will be handled by the send() method when responses are ready
968968

969+
// Arm keep-alive only after the fallible awaits above — an error
970+
// path returning 400 discards the Response, so nothing could ever
971+
// cancel the stream and clear an already-armed timer. Skip if the
972+
// responses already completed and cleaned the stream up.
973+
if (this._streamMapping.get(streamId)?.controller === streamController!) {
974+
this.startKeepAlive(streamId, streamController!, encoder);
975+
}
976+
969977
return new Response(readable, { status: 200, headers });
970978
} catch (error) {
971979
// return JSON-RPC formatted error

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

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1522,3 +1522,93 @@ describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => {
15221522
await transport.close();
15231523
});
15241524
});
1525+
1526+
describe('WebStandardStreamableHTTPServerTransport SSE keep-alive lifecycle', () => {
1527+
beforeEach(() => {
1528+
vi.useFakeTimers();
1529+
});
1530+
1531+
afterEach(() => {
1532+
vi.useRealTimers();
1533+
});
1534+
1535+
it('should not arm keep-alive when the transport closes during an event-store replay await', async () => {
1536+
let releaseReplay: (() => void) | undefined;
1537+
const eventStore: EventStore = {
1538+
async storeEvent(): Promise<EventId> {
1539+
return 'evt-1';
1540+
},
1541+
async replayEventsAfter(): Promise<StreamId> {
1542+
await new Promise<void>(resolve => {
1543+
releaseReplay = resolve;
1544+
});
1545+
return 'stream-1';
1546+
}
1547+
};
1548+
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore });
1549+
await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport);
1550+
const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize));
1551+
const sessionId = initResponse.headers.get('mcp-session-id') as string;
1552+
1553+
// Enter replayEvents and park on the replayEventsAfter await
1554+
const pendingGet = transport.handleRequest(
1555+
createRequest('GET', undefined, { sessionId, extraHeaders: { 'Last-Event-ID': 'evt-1' } })
1556+
);
1557+
await vi.advanceTimersByTimeAsync(0);
1558+
expect(releaseReplay).toBeDefined();
1559+
1560+
// Close the transport mid-await, then let the replay continuation run
1561+
await transport.close();
1562+
releaseReplay?.();
1563+
await pendingGet;
1564+
1565+
// The deferred continuation must not have armed a timer close() can never sweep
1566+
expect(vi.getTimerCount()).toBe(0);
1567+
});
1568+
1569+
it('should not leak a keep-alive timer when the priming event write fails on a POST SSE stream', async () => {
1570+
// Healthy during initialization, then the store starts failing — the
1571+
// tool call's priming event write must reject inside handlePostRequest
1572+
let storeFails = false;
1573+
const eventStore: EventStore = {
1574+
async storeEvent(): Promise<EventId> {
1575+
if (storeFails) {
1576+
throw new Error('event store unavailable');
1577+
}
1578+
return `evt-${randomUUID()}`;
1579+
},
1580+
async replayEventsAfter(): Promise<StreamId> {
1581+
return 'stream-1';
1582+
}
1583+
};
1584+
const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore });
1585+
const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' });
1586+
mcpServer.registerTool('noop', { description: 'noop' }, async (): Promise<CallToolResult> => ({ content: [] }));
1587+
await mcpServer.connect(transport);
1588+
1589+
const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize));
1590+
expect(initResponse.status).toBe(200);
1591+
const sessionId = initResponse.headers.get('mcp-session-id') as string;
1592+
// Let the init response finish sending (its send() stores an event and
1593+
// then cleans up the init stream's keep-alive) before failing the store
1594+
await vi.advanceTimersByTimeAsync(0);
1595+
expect(vi.getTimerCount()).toBe(0);
1596+
storeFails = true;
1597+
1598+
const response = await transport.handleRequest(
1599+
createRequest(
1600+
'POST',
1601+
{ jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'call-1' } as JSONRPCMessage,
1602+
{
1603+
sessionId
1604+
}
1605+
)
1606+
);
1607+
expect(response.status).toBe(400);
1608+
1609+
// The discarded stream must not carry a permanently-firing timer
1610+
expect(vi.getTimerCount()).toBe(0);
1611+
1612+
await transport.close();
1613+
});
1614+
});

0 commit comments

Comments
 (0)