diff --git a/.changeset/streamable-http-sse-keepalive.md b/.changeset/streamable-http-sse-keepalive.md new file mode 100644 index 0000000000..f9c107a17f --- /dev/null +++ b/.changeset/streamable-http-sse-keepalive.md @@ -0,0 +1,5 @@ +--- +'@modelcontextprotocol/server': minor +--- + +SSE streams served by the SDK now emit keep-alive comment frames (`: keepalive`) so idle connections are not killed by client body-idle timeouts or intermediaries. `WebStandardStreamableHTTPServerTransport` and `PerRequestHTTPServerTransport` gain a `keepAliveMs` option (default 15000; set 0 to disable), and `createMcpHandler`'s existing `keepAliveMs` now covers modern per-request exchange streams and the legacy stateless fallback in addition to `subscriptions/listen` streams. Streaming responses disable proxy transformations and nginx-style buffering, and deferred work cannot register streams after transport close. diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index 1fc325b7e4..f832246607 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -154,6 +154,14 @@ Rewrite the imports: 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. +## `SSE stream disconnected: TypeError: terminated` + +An idle SSE response went too long without delivering body bytes. A client-side body-idle timeout (for example undici's `bodyTimeout`) or an intermediary such as a reverse proxy or cloud load balancer may then terminate the stream. The client observes this as `TypeError: terminated` and reconnects in a loop. + +The SDK's HTTP serving prevents this by writing an SSE comment frame (`: keepalive`) to every open SSE stream every 15 seconds by default — `WebStandardStreamableHTTPServerTransport` on all of its streams, and `createMcpHandler` on `subscriptions/listen` streams, modern per-request exchange streams, and the legacy fallback's per-request transport. Comment frames are dropped by SSE parsers before event dispatch, so they never surface as protocol messages. Tune or disable the interval with the `keepAliveMs` option on the transport or handler (`0` disables). + +If you still see this error, keep-alive may be disabled (`keepAliveMs: 0`), the client timeout may be shorter than the configured interval, or an intermediary may buffer or strip SSE data. Check for proxies that buffer streaming responses (for example nginx without `proxy_buffering off`). + ## Recap - Every heading on this page is the exact message you searched for. diff --git a/packages/server/src/server/createMcpHandler.ts b/packages/server/src/server/createMcpHandler.ts index f17deaa860..360b4f066f 100644 --- a/packages/server/src/server/createMcpHandler.ts +++ b/packages/server/src/server/createMcpHandler.ts @@ -145,8 +145,9 @@ export interface CreateMcpHandlerOptions { * - `'stateless'` (the default, also when the option is omitted) — * old-school stateless serving: each legacy request is answered by a * fresh instance from the same factory over a streamable HTTP transport - * constructed with only `sessionIdGenerator: undefined` (the established - * stateless idiom). Because serving is per-request and stateless, GET and + * constructed with `sessionIdGenerator: undefined` (the established + * stateless idiom), plus the handler's `keepAliveMs` when provided. + * Because serving is per-request and stateless, GET and * DELETE (2025 session operations) are answered with `405` / * `Method not allowed.`. * - `'reject'` — modern-only strict: legacy-classified requests are @@ -194,8 +195,12 @@ export interface CreateMcpHandlerOptions { */ maxSubscriptions?: number; /** - * SSE comment-frame keepalive interval for `subscriptions/listen` streams, - * in milliseconds. Set to `0` to disable. + * SSE comment-frame keepalive interval, in milliseconds, applied to every + * SSE stream this handler serves: `subscriptions/listen` streams, modern + * per-request exchange streams, and the legacy stateless fallback's + * per-request transport. In modern `auto` mode, the timer starts only after + * the exchange upgrades to SSE; use `responseMode: 'sse'` when a silent + * long-running handler needs heartbeat bytes. Set to `0` to disable. * @default 15000 */ keepAliveMs?: number; @@ -295,18 +300,26 @@ function internalServerErrorResponse(id: RequestId | null = null): Response { * strict modern endpoint). * * Each POST is served by a fresh instance from the factory connected to a - * fresh streamable HTTP transport constructed with only - * `sessionIdGenerator: undefined` — the established stateless idiom, unchanged. - * Because serving is per-request and stateless, GET and DELETE (2025 session - * operations) are answered with `405` / `Method not allowed.`, exactly like the - * canonical stateless example. + * fresh streamable HTTP transport constructed with + * `sessionIdGenerator: undefined` (the established stateless idiom) plus any + * `transportOptions`. Because serving is per-request and stateless, GET and + * DELETE (2025 session operations) are answered with `405` / + * `Method not allowed.`, exactly like the canonical stateless example. * * The optional `onerror` callback receives factory and serving failures on * this leg (reporting only — the response stays the 500 internal-error body). * The entry passes its own `onerror` here when expanding the default, so * legacy-leg failures are never silently swallowed. + * + * The optional `transportOptions` are threaded into each per-request + * transport; currently just `keepAliveMs`, the SSE keep-alive comment-frame + * interval (the entry forwards its own `keepAliveMs` option here). */ -export function legacyStatelessFallback(factory: McpServerFactory, onerror?: (error: Error) => void): LegacyHttpHandler { +export function legacyStatelessFallback( + factory: McpServerFactory, + onerror?: (error: Error) => void, + transportOptions?: { keepAliveMs?: number } +): LegacyHttpHandler { return async (request, options) => { if (request.method.toUpperCase() !== 'POST') { return jsonRpcErrorResponse(405, -32_000, 'Method not allowed.'); @@ -317,7 +330,10 @@ export function legacyStatelessFallback(factory: McpServerFactory, onerror?: (er ...(options?.authInfo !== undefined && { authInfo: options.authInfo }), requestInfo: request }); - const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: undefined }); + const transport = new WebStandardStreamableHTTPServerTransport({ + sessionIdGenerator: undefined, + ...(transportOptions?.keepAliveMs !== undefined && { keepAliveMs: transportOptions.keepAliveMs }) + }); await product.connect(transport); const teardown = () => { @@ -632,7 +648,12 @@ export function createMcpHandler(factory: McpServerFactory, options: CreateMcpHa // The default posture is the stateless fallback; 'reject' is the only way // to turn legacy serving off (modern-only strict). - const legacyHandler: LegacyHttpHandler | undefined = legacy === 'reject' ? undefined : legacyStatelessFallback(factory, reportError); + const legacyHandler: LegacyHttpHandler | undefined = + legacy === 'reject' + ? undefined + : legacyStatelessFallback(factory, reportError, { + ...(options.keepAliveMs !== undefined && { keepAliveMs: options.keepAliveMs }) + }); async function serveModern(route: InboundModernRoute, request: Request, authInfo: AuthInfo | undefined): Promise { const claimedRevision = route.classification.revision; @@ -778,7 +799,8 @@ export function createMcpHandler(factory: McpServerFactory, options: CreateMcpHa classification: route.classification, request, ...(authInfo !== undefined && { authInfo }), - ...(responseMode !== undefined && { responseMode }) + ...(responseMode !== undefined && { responseMode }), + ...(options.keepAliveMs !== undefined && { keepAliveMs: options.keepAliveMs }) }); if (route.messageKind === 'notification') { // Notification exchanges have no terminal response to ride the diff --git a/packages/server/src/server/invoke.ts b/packages/server/src/server/invoke.ts index 6966968604..97b5b17553 100644 --- a/packages/server/src/server/invoke.ts +++ b/packages/server/src/server/invoke.ts @@ -35,6 +35,12 @@ export interface InvokeContext { authInfo?: AuthInfo; /** Response shaping for the exchange; defaults to `auto` (lazy SSE upgrade). */ responseMode?: PerRequestResponseMode; + /** + * SSE keep-alive comment-frame interval for the exchange's stream, in + * milliseconds; passed through to the per-request transport. `0` disables. + * @default 15000 + */ + keepAliveMs?: number; } /** @@ -58,7 +64,8 @@ export async function invoke( ): Promise { const transport = new PerRequestHTTPServerTransport({ classification: ctx.classification, - ...(ctx.responseMode !== undefined && { responseMode: ctx.responseMode }) + ...(ctx.responseMode !== undefined && { responseMode: ctx.responseMode }), + ...(ctx.keepAliveMs !== undefined && { keepAliveMs: ctx.keepAliveMs }) }); await server.connect(transport); return transport.handleMessage(message, { diff --git a/packages/server/src/server/listenRouter.ts b/packages/server/src/server/listenRouter.ts index 96dcb16beb..7ccce8dadd 100644 --- a/packages/server/src/server/listenRouter.ts +++ b/packages/server/src/server/listenRouter.ts @@ -34,9 +34,10 @@ import { codecForVersion, MODERN_WIRE_REVISION, SERVER_INFO_META_KEY, SUBSCRIPTI import type { ServerEventBus } from './serverEventBus'; import { honoredSubset, listenFilterAccepts, serverEventToNotification } from './serverEventBus'; +import { armSseKeepAlive, DEFAULT_SSE_KEEP_ALIVE_MS } from './sseKeepAlive'; /** Default SSE comment-frame keepalive interval for listen streams. */ -export const DEFAULT_LISTEN_KEEPALIVE_MS = 15_000; +export const DEFAULT_LISTEN_KEEPALIVE_MS = DEFAULT_SSE_KEEP_ALIVE_MS; /** Default capacity guard: refuse a new subscription when this many are already open. */ export const DEFAULT_MAX_SUBSCRIPTIONS = 1024; @@ -218,14 +219,7 @@ export function createListenRouter(options: ListenRouterOptions): ListenRouter { writeNotification(note.method, note.params); }); - if (keepAliveMs > 0) { - keepAliveTimer = setInterval(() => writeFrame(': keepalive\n\n'), keepAliveMs); - // Do not hold the event loop open on idle subscriptions. Node's - // setInterval returns a Timeout with .unref(); browsers/Workers - // return a number — the cast is an environment shim, not a - // workaround for SDK typing. - (keepAliveTimer as { unref?: () => void }).unref?.(); - } + keepAliveTimer = armSseKeepAlive(keepAliveMs, () => writeFrame(': keepalive\n\n')); open.add(teardown); }, @@ -251,7 +245,7 @@ export function createListenRouter(options: ListenRouterOptions): ListenRouter { status: 200, headers: { 'Content-Type': 'text/event-stream', - 'Cache-Control': 'no-cache', + 'Cache-Control': 'no-cache, no-transform', Connection: 'keep-alive', 'X-Accel-Buffering': 'no' } diff --git a/packages/server/src/server/perRequestTransport.ts b/packages/server/src/server/perRequestTransport.ts index 5003946404..53400a05f5 100644 --- a/packages/server/src/server/perRequestTransport.ts +++ b/packages/server/src/server/perRequestTransport.ts @@ -58,6 +58,8 @@ import { SdkErrorCode } from '@modelcontextprotocol/core-internal'; +import { armSseKeepAlive, DEFAULT_SSE_KEEP_ALIVE_MS } from './sseKeepAlive'; + /** * How the transport shapes its HTTP response for a request: * @@ -79,6 +81,15 @@ export interface PerRequestHTTPServerTransportOptions { classification: MessageClassification; /** Response shaping for the exchange; defaults to `auto`. */ responseMode?: PerRequestResponseMode; + /** + * Interval in milliseconds between SSE keep-alive comment frames + * (`: keepalive`) written while the exchange's SSE stream is open. With + * `responseMode: 'sse'`, this also protects long-running handlers that emit + * no mid-call output. Set to `0` to disable; values below `1`, above + * `2147483647`, or non-finite values also disable the timer. + * @default 15000 + */ + keepAliveMs?: number; } /** Per-exchange context handed to {@linkcode PerRequestHTTPServerTransport.handleMessage}. */ @@ -140,10 +151,13 @@ export class PerRequestHTTPServerTransport implements Transport { private _deferredResponse?: DeferredResponse; private _sse?: SseSink; private _abortCleanup?: () => void; + private readonly _keepAliveMs: number; + private _keepAliveTimer?: ReturnType; constructor(options: PerRequestHTTPServerTransportOptions) { this._classification = options.classification; this._responseMode = options.responseMode ?? 'auto'; + this._keepAliveMs = options.keepAliveMs ?? DEFAULT_SSE_KEEP_ALIVE_MS; } async start(): Promise { @@ -342,6 +356,7 @@ export class PerRequestHTTPServerTransport implements Transport { this._abortCleanup?.(); this._abortCleanup = undefined; + this.stopKeepAlive(); if (this._sse !== undefined && !this._sse.closed) { this._sse.closed = true; @@ -382,13 +397,14 @@ export class PerRequestHTTPServerTransport implements Transport { } }); this._sse = { controller, encoder: new TextEncoder(), closed: false }; + this.startKeepAlive(); this.settleResponse( new Response(readable, { status: 200, headers: { 'Content-Type': 'text/event-stream', - 'Cache-Control': 'no-cache', + 'Cache-Control': 'no-cache, no-transform', Connection: 'keep-alive', // Disable proxy buffering so streamed messages are // delivered as they are written. @@ -398,7 +414,29 @@ export class PerRequestHTTPServerTransport implements Transport { ); } + /** + * Arms the exchange's keep-alive interval, writing an SSE comment frame + * every `keepAliveMs` while the stream is open. `writeCommentFrame` + * already drops frames once the exchange is closed or the stream is + * finalized, so the interval body needs no extra guards; the timer itself + * is cleared on stream finalization and transport close. + */ + private startKeepAlive(): void { + if (this._closed) { + return; + } + this._keepAliveTimer = armSseKeepAlive(this._keepAliveMs, () => this.writeCommentFrame('keepalive')); + } + + private stopKeepAlive(): void { + if (this._keepAliveTimer !== undefined) { + clearInterval(this._keepAliveTimer); + this._keepAliveTimer = undefined; + } + } + private finalizeStream(): void { + this.stopKeepAlive(); if (this._sse !== undefined && !this._sse.closed) { this._sse.closed = true; try { diff --git a/packages/server/src/server/sseKeepAlive.ts b/packages/server/src/server/sseKeepAlive.ts new file mode 100644 index 0000000000..fa43faf8ec --- /dev/null +++ b/packages/server/src/server/sseKeepAlive.ts @@ -0,0 +1,25 @@ +/** Default interval between SSE keep-alive comment frames. */ +export const DEFAULT_SSE_KEEP_ALIVE_MS = 15_000; + +/** Largest delay accepted by JavaScript timers without overflow clamping. */ +const MAX_TIMER_DELAY_MS = 2_147_483_647; + +/** + * Arms an SSE keep-alive timer when `intervalMs` is a valid timer delay. + * + * Invalid delays disable keep-alive. In Node.js, sub-millisecond, non-finite, + * and overflowing delays are clamped to roughly 1 ms, which would otherwise + * flood every open stream with comment frames. + * + * @returns The timer handle, or `undefined` when keep-alive is disabled. + */ +export function armSseKeepAlive(intervalMs: number, onTick: () => void): ReturnType | undefined { + if (!Number.isFinite(intervalMs) || intervalMs < 1 || intervalMs > MAX_TIMER_DELAY_MS) { + return undefined; + } + + const timer = setInterval(onTick, intervalMs); + // Node.js timers expose `unref`; browser and Workers timers are numbers. + (timer as { unref?: () => void }).unref?.(); + return timer; +} diff --git a/packages/server/src/server/streamableHttp.ts b/packages/server/src/server/streamableHttp.ts index 7da5fb853c..79a4c136be 100644 --- a/packages/server/src/server/streamableHttp.ts +++ b/packages/server/src/server/streamableHttp.ts @@ -19,6 +19,8 @@ import { SUPPORTED_PROTOCOL_VERSIONS } from '@modelcontextprotocol/core-internal'; +import { armSseKeepAlive, DEFAULT_SSE_KEEP_ALIVE_MS } from './sseKeepAlive'; + export type StreamId = string; export type EventId = string; @@ -148,6 +150,20 @@ export interface WebStandardStreamableHTTPServerTransportOptions { */ retryInterval?: number; + /** + * Interval in milliseconds between SSE keep-alive comment frames (`: keepalive`) + * written to open SSE streams. Keep-alive frames prevent idle streams (e.g. the + * standalone `GET` stream, or a `POST` stream during a long-running tool call) + * from being killed by intermediaries and server idle timeouts, which clients + * observe as `SSE stream disconnected: TypeError: terminated`. + * + * Comment frames are ignored by SSE parsers and never surface as messages. + * Defaults to `15000` (per the WHATWG SSE spec recommendation of roughly every + * 15 seconds). Set to `0` to disable keep-alive frames; values below `1`, above + * `2147483647`, or non-finite values also disable the timer. + */ + keepAliveMs?: number; + /** * List of protocol versions that this transport will accept. * Used to validate the `mcp-protocol-version` header in incoming requests. @@ -247,6 +263,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { private _enableDnsRebindingProtection: boolean; private _retryInterval?: number; private _supportedProtocolVersions: string[]; + private _keepAliveMs: number; + private _keepAliveTimers: Map> = new Map(); sessionId?: string; onclose?: () => void; @@ -264,6 +282,49 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { this._enableDnsRebindingProtection = options.enableDnsRebindingProtection ?? false; this._retryInterval = options.retryInterval; this._supportedProtocolVersions = options.supportedProtocolVersions ?? SUPPORTED_PROTOCOL_VERSIONS; + this._keepAliveMs = options.keepAliveMs ?? DEFAULT_SSE_KEEP_ALIVE_MS; + } + + /** + * Arms a keep-alive interval for an SSE stream that periodically writes an SSE + * comment frame so intermediaries and idle timeouts don't kill the connection. + * Replaces any timer already armed for the same stream id (a resumed stream + * re-registered under the same id supersedes its predecessor's timer). The + * timer is cleared via {@linkcode stopKeepAlive} when the stream is cleaned up, + * and clears itself if a write fails (stream already closed/cancelled). + */ + private startKeepAlive( + streamId: string, + controller: ReadableStreamDefaultController, + encoder: InstanceType + ): void { + // A deferred arm (e.g. after an event-store await) must not outlive the + // transport: close()'s timer sweep has already run and never runs again. + if (this._closed) { + return; + } + this.stopKeepAlive(streamId); + const timer = armSseKeepAlive(this._keepAliveMs, () => { + try { + controller.enqueue(encoder.encode(': keepalive\n\n')); + } catch { + this.stopKeepAlive(streamId); + } + }); + if (timer !== undefined) { + this._keepAliveTimers.set(streamId, timer); + } + } + + /** + * Clears the keep-alive interval for a stream, if one is armed. + */ + private stopKeepAlive(streamId: string): void { + const timer = this._keepAliveTimers.get(streamId); + if (timer !== undefined) { + clearInterval(timer); + this._keepAliveTimers.delete(streamId); + } } /** @@ -352,6 +413,10 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Returns a `Response` object (Web Standard) */ async handleRequest(req: Request, options?: HandleRequestOptions): Promise { + if (this._closed) { + return this.createJsonErrorResponse(404, -32_001, 'Session not found'); + } + // Validate request headers for DNS rebinding protection const validationError = this.validateRequestHeaders(req); if (validationError) { @@ -473,6 +538,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // it still points at THIS controller — a stale cancel must not // delete a successor stream registered by a later GET/resume. if (this._streamMapping.get(this._standaloneSseStreamId)?.controller === streamController) { + this.stopKeepAlive(this._standaloneSseStreamId); this._streamMapping.delete(this._standaloneSseStreamId); } } @@ -481,7 +547,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', - Connection: 'keep-alive' + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; // After initialization, always include the session ID if we have one @@ -494,6 +561,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { controller: streamController!, encoder, cleanup: () => { + this.stopKeepAlive(this._standaloneSseStreamId); this._streamMapping.delete(this._standaloneSseStreamId); try { streamController!.close(); @@ -503,6 +571,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); + this.startKeepAlive(this._standaloneSseStreamId, streamController!, encoder); + return new Response(readable, { headers }); } @@ -537,7 +607,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache, no-transform', - Connection: 'keep-alive' + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; if (this.sessionId !== undefined) { @@ -564,6 +635,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // a stale cancel from an earlier resume must not delete a // successor resumed stream a re-poll has since registered. if (replayedStreamId !== undefined && this._streamMapping.get(replayedStreamId)?.controller === streamController) { + this.stopKeepAlive(replayedStreamId); this._streamMapping.delete(replayedStreamId); } } @@ -585,11 +657,27 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } }); + // The transport may have closed while the replay await was parked: + // its cleanup sweep ran before this stream was registered, so + // registering now would strand a mapping (and an open controller) + // on a dead transport, hang the client on a stream that never + // ends, and 409-block a later resume of this stream id. End the + // stream instead so the client observes termination. + if (this._closed) { + try { + streamController!.close(); + } catch { + // Controller might already be closed + } + return new Response(readable, { headers }); + } + this._streamMapping.set(replayedStreamId, { controller: streamController!, encoder, replayedEventIds, cleanup: () => { + this.stopKeepAlive(replayedStreamId!); this._streamMapping.delete(replayedStreamId!); try { streamController!.close(); @@ -618,6 +706,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + // Only arm keep-alive if the stream is still registered — the + // no-in-flight-request path above may have already closed it. + if (this._streamMapping.get(replayedStreamId)?.controller === streamController!) { + this.startKeepAlive(replayedStreamId, streamController!, encoder); + } + return new Response(readable, { headers }); } catch (error) { this.onerror?.(error as Error); @@ -677,6 +771,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { * Handles `POST` requests containing JSON-RPC messages */ private async handlePostRequest(req: Request, options?: HandleRequestOptions): Promise { + // Set once the SSE stream bookkeeping has been registered, so the + // catch below can reclaim it: an error after registration (a failed + // priming event write, a throwing message handler) returns an error + // response, leaving nothing that could ever cancel the discarded + // stream or retire the request mappings. + let reclaimSseBookkeeping: (() => void) | undefined; try { // Validate the Accept header const acceptHeader = req.headers.get('accept'); @@ -770,6 +870,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + // Request parsing and session initialization may await user/runtime + // work. Do not register or dispatch after close() has swept state. + if (this._closed) { + return this.createJsonErrorResponse(404, -32_001, 'Session not found'); + } + // check if it contains requests const hasRequests = messages.some(element => isJSONRPCRequest(element)); @@ -830,6 +936,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // resumed stream under the same streamId) must not delete // the successor. if (this._streamMapping.get(streamId)?.controller === streamController) { + this.stopKeepAlive(streamId); this._streamMapping.delete(streamId); } } @@ -837,8 +944,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { const headers: Record = { 'Content-Type': 'text/event-stream', - 'Cache-Control': 'no-cache', - Connection: 'keep-alive' + 'Cache-Control': 'no-cache, no-transform', + Connection: 'keep-alive', + 'X-Accel-Buffering': 'no' }; // After initialization, always include the session ID if we have one @@ -854,6 +962,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { controller: streamController!, encoder, cleanup: () => { + this.stopKeepAlive(streamId); this._streamMapping.delete(streamId); try { streamController!.close(); @@ -866,6 +975,15 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } } + reclaimSseBookkeeping = () => { + this._streamMapping.get(streamId)?.cleanup(); + for (const message of messages) { + if (isJSONRPCRequest(message) && this._requestToStreamMapping.get(message.id) === streamId) { + this._requestToStreamMapping.delete(message.id); + } + } + }; + // Write priming event if event store is configured (after mapping is set up) await this.writePrimingEvent(streamController!, encoder, streamId, clientProtocolVersion); @@ -891,10 +1009,19 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { // The server SHOULD NOT close the SSE stream before sending all JSON-RPC responses // This will be handled by the send() method when responses are ready + // Arm keep-alive only after the fallible awaits above — an error + // path returning 400 discards the Response, so nothing could ever + // cancel the stream and clear an already-armed timer. Skip if the + // responses already completed and cleaned the stream up. + if (this._streamMapping.get(streamId)?.controller === streamController!) { + this.startKeepAlive(streamId, streamController!, encoder); + } + return new Response(readable, { status: 200, headers }); } catch (error) { // return JSON-RPC formatted error this.onerror?.(error as Error); + reclaimSseBookkeeping?.(); return this.createJsonErrorResponse(400, -32_700, 'Parse error', { data: String(error) }); } } @@ -986,6 +1113,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport { } this._streamMapping.clear(); + // Clear any keep-alive timers not already cleared by stream cleanup + for (const timer of this._keepAliveTimers.values()) { + clearInterval(timer); + } + this._keepAliveTimers.clear(); + // Clear any pending responses this._requestResponseMap.clear(); this.onclose?.(); diff --git a/packages/server/test/server/createMcpHandler.test.ts b/packages/server/test/server/createMcpHandler.test.ts index 232f781926..ded506e57c 100644 --- a/packages/server/test/server/createMcpHandler.test.ts +++ b/packages/server/test/server/createMcpHandler.test.ts @@ -820,3 +820,57 @@ describe('createMcpHandler — close()', () => { // Type-level pin: a zero-argument factory stays assignable to McpServerFactory unchanged. const zeroArgFactory = () => new McpServer({ name: 'zero-arg', version: '1.0.0' }); void createMcpHandler(zeroArgFactory); + +describe('createMcpHandler — keepAliveMs', () => { + function gatedFactory(): { factory: () => McpServer; release: () => void } { + let release!: () => void; + const gate = new Promise(resolve => { + release = resolve; + }); + const factory = (): McpServer => { + const s = new McpServer({ name: 'ka', version: '1.0.0' }); + s.registerTool('gated', { inputSchema: z.object({}) }, async () => { + await gate; + return { content: [{ type: 'text', text: 'done' }] }; + }); + return s; + }; + return { factory, release }; + } + + it('threads keepAliveMs into the modern per-request exchange stream', async () => { + vi.useFakeTimers(); + try { + const { factory, release } = gatedFactory(); + const handler = createMcpHandler(factory, { responseMode: 'sse', keepAliveMs: 1_000 }); + const responsePromise = handler.fetch(postRequest(modernToolsCall('gated', {}))); + await vi.advanceTimersByTimeAsync(1_000); + release(); + const response = await responsePromise; + expect(response.headers.get('content-type')).toContain('text/event-stream'); + const text = await response.text(); + expect(text).toContain(': keepalive'); + } finally { + vi.useRealTimers(); + } + }); + + it('threads keepAliveMs into the legacy stateless fallback per-request transport', async () => { + vi.useFakeTimers(); + try { + const { factory, release } = gatedFactory(); + const handler = createMcpHandler(factory, { keepAliveMs: 1_000 }); + const responsePromise = handler.fetch( + postRequest({ jsonrpc: '2.0', id: 9, method: 'tools/call', params: { name: 'gated', arguments: {} } }) + ); + await vi.advanceTimersByTimeAsync(1_000); + release(); + const response = await responsePromise; + expect(response.headers.get('content-type')).toContain('text/event-stream'); + const text = await response.text(); + expect(text).toContain(': keepalive'); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/packages/server/test/server/createMcpHandlerListen.test.ts b/packages/server/test/server/createMcpHandlerListen.test.ts index 2fa5000742..7a524d0c9f 100644 --- a/packages/server/test/server/createMcpHandlerListen.test.ts +++ b/packages/server/test/server/createMcpHandlerListen.test.ts @@ -13,7 +13,7 @@ import { PROTOCOL_VERSION_META_KEY, SUBSCRIPTION_ID_META_KEY } from '@modelcontextprotocol/core-internal'; -import { describe, expect, it } from 'vitest'; +import { describe, expect, it, vi } from 'vitest'; import { createMcpHandler } from '../../src/server/createMcpHandler'; import { McpServer } from '../../src/server/mcp'; @@ -105,6 +105,8 @@ describe('createMcpHandler — subscriptions/listen', () => { ); const response = await handler.fetch(listenRequest(1, { toolsListChanged: true })); expect(response.status).toBe(200); + expect(response.headers.get('cache-control')).toBe('no-cache, no-transform'); + expect(response.headers.get('x-accel-buffering')).toBe('no'); const [ack] = await readMessages(response, 1); // The factory is consulted exactly once (capabilities probe only); the // instance is never connected and is closed immediately after the @@ -116,6 +118,27 @@ describe('createMcpHandler — subscriptions/listen', () => { await handler.close(); }); + it.each([0.5, Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( + 'disables keep-alive for invalid keepAliveMs %s instead of arming a clamped interval', + async keepAliveMs => { + vi.useFakeTimers(); + try { + const handler = createMcpHandler(trivialFactory(), { keepAliveMs }); + const response = await handler.fetch(listenRequest(1, { toolsListChanged: true })); + const reader = response.body!.getReader(); + await reader.read(); // acknowledgement + expect(vi.getTimerCount()).toBe(0); + await vi.advanceTimersByTimeAsync(60_000); + const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); + expect(raced).toBe('pending'); + await reader.cancel(); + await handler.close(); + } finally { + vi.useRealTimers(); + } + } + ); + it('ack is the first frame, stamped with the listen id verbatim, carrying the honored subset', async () => { const handler = createMcpHandler(trivialFactory(), { keepAliveMs: 0 }); const response = await handler.fetch(listenRequest('sub-42', { toolsListChanged: true, promptsListChanged: false })); diff --git a/packages/server/test/server/perRequestStreaming.test.ts b/packages/server/test/server/perRequestStreaming.test.ts index 6d350ed2ef..9440cbdf4e 100644 --- a/packages/server/test/server/perRequestStreaming.test.ts +++ b/packages/server/test/server/perRequestStreaming.test.ts @@ -11,7 +11,7 @@ import { PROTOCOL_VERSION_META_KEY, setNegotiatedProtocolVersion } from '@modelcontextprotocol/core-internal'; -import { describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import type { PerRequestResponseMode } from '../../src/server/perRequestTransport'; import { PerRequestHTTPServerTransport } from '../../src/server/perRequestTransport'; @@ -46,14 +46,16 @@ interface StreamingSetup { async function setup( handler: (ctx: ServerContext) => Promise, - responseMode?: PerRequestResponseMode + responseMode?: PerRequestResponseMode, + keepAliveMs?: number ): Promise { const server = new Server({ name: 'streaming-test', version: '1.0.0' }, { capabilities: { tools: {} } }); server.setRequestHandler('tools/call', async (_request, ctx) => handler(ctx)); setNegotiatedProtocolVersion(server, MODERN_REVISION); const transport = new PerRequestHTTPServerTransport({ classification: MODERN, - ...(responseMode !== undefined && { responseMode }) + ...(responseMode !== undefined && { responseMode }), + ...(keepAliveMs !== undefined && { keepAliveMs }) }); await server.connect(transport); return { server, transport }; @@ -92,7 +94,7 @@ describe('lazy upgrade matrix', () => { const response = await transport.handleMessage(toolsCall()); expect(response.status).toBe(200); expect(response.headers.get('content-type')).toBe('text/event-stream'); - expect(response.headers.get('cache-control')).toBe('no-cache'); + expect(response.headers.get('cache-control')).toBe('no-cache, no-transform'); expect(response.headers.get('x-accel-buffering')).toBe('no'); const frames = await sseFrames(response); @@ -249,3 +251,105 @@ describe('disconnect is cancellation', () => { expect(observedSignal?.aborted).toBe(true); }); }); + +describe('keep-alive', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('writes keep-alive comment frames while a forced-sse exchange is streaming', async () => { + let release!: () => void; + const gate = new Promise(resolve => { + release = resolve; + }); + const { transport } = await setup(async () => { + await gate; + return { content: [] }; + }, 'sse'); + + const responsePromise = transport.handleMessage(toolsCall()); + // The stream opened at dispatch end; the handler now idles past the + // default interval with no mid-call output. + await vi.advanceTimersByTimeAsync(15_000); + release(); + const response = await responsePromise; + const frames = await sseFrames(response); + expect(frames[0]).toBe(': keepalive'); + + // The exchange completed and closed the transport: no timer survives. + expect(vi.getTimerCount()).toBe(0); + }); + + it('writes keep-alive frames after an auto exchange upgrades to SSE', async () => { + let release!: () => void; + const gate = new Promise(resolve => { + release = resolve; + }); + const { transport } = await setup(async ctx => { + await ctx.mcpReq.notify(progressNotification(1)); + await gate; + return { content: [] }; + }); + + const responsePromise = transport.handleMessage(toolsCall()); + // Let the handler run, emit the upgrading notification, then idle. + await vi.advanceTimersByTimeAsync(15_000); + release(); + const response = await responsePromise; + const frames = await sseFrames(response); + expect(frames).toContain(': keepalive'); + expect(vi.getTimerCount()).toBe(0); + }); + + it('does not write keep-alive frames when keepAliveMs is 0', async () => { + let release!: () => void; + const gate = new Promise(resolve => { + release = resolve; + }); + const { transport } = await setup( + async () => { + await gate; + return { content: [] }; + }, + 'sse', + 0 + ); + + const responsePromise = transport.handleMessage(toolsCall()); + await vi.advanceTimersByTimeAsync(60_000); + release(); + const response = await responsePromise; + const frames = await sseFrames(response); + expect(frames.some(frame => frame.startsWith(': keepalive'))).toBe(false); + }); + + it.each([0.5, Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( + 'disables keep-alive for invalid keepAliveMs %s instead of arming a clamped interval', + async keepAliveMs => { + let release!: () => void; + const gate = new Promise(resolve => { + release = resolve; + }); + const { transport } = await setup( + async () => { + await gate; + return { content: [] }; + }, + 'sse', + keepAliveMs + ); + + const responsePromise = transport.handleMessage(toolsCall()); + expect(vi.getTimerCount()).toBe(0); + await vi.advanceTimersByTimeAsync(1_000); + release(); + const response = await responsePromise; + const frames = await sseFrames(response); + expect(frames.some(frame => frame.startsWith(': keepalive'))).toBe(false); + } + ); +}); diff --git a/packages/server/test/server/streamableHttp.test.ts b/packages/server/test/server/streamableHttp.test.ts index beca451113..f54708c8f0 100644 --- a/packages/server/test/server/streamableHttp.test.ts +++ b/packages/server/test/server/streamableHttp.test.ts @@ -165,6 +165,8 @@ describe('Zod v4', () => { expect(response.status).toBe(200); expect(response.headers.get('content-type')).toBe('text/event-stream'); + expect(response.headers.get('cache-control')).toBe('no-cache, no-transform'); + expect(response.headers.get('x-accel-buffering')).toBe('no'); expect(response.headers.get('mcp-session-id')).toBeDefined(); }); @@ -357,6 +359,8 @@ describe('Zod v4', () => { expect(response.status).toBe(200); expect(response.headers.get('content-type')).toBe('text/event-stream'); + expect(response.headers.get('cache-control')).toBe('no-cache, no-transform'); + expect(response.headers.get('x-accel-buffering')).toBe('no'); expect(response.headers.get('mcp-session-id')).toBe(sessionId); }); @@ -857,6 +861,8 @@ describe('Zod v4', () => { createRequest('GET', undefined, { sessionId, extraHeaders: { 'Last-Event-ID': primingId! } }) ); expect(reconnect.status).toBe(200); + expect(reconnect.headers.get('cache-control')).toBe('no-cache, no-transform'); + expect(reconnect.headers.get('x-accel-buffering')).toBe('no'); release(); const replayed = await readSSEEvent(reconnect); expect(replayed).toContain('notifications/progress'); @@ -1407,3 +1413,333 @@ describe('Zod v4', () => { }); }); }); + +describe('WebStandardStreamableHTTPServerTransport SSE keep-alive', () => { + async function createTransport(options?: { keepAliveMs?: number }): Promise<{ + transport: WebStandardStreamableHTTPServerTransport; + sessionId: string; + }> { + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), ...options }); + await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport); + const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + expect(initResponse.status).toBe(200); + return { transport, sessionId: initResponse.headers.get('mcp-session-id') as string }; + } + + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('should write keep-alive comment frames to an idle standalone GET stream', async () => { + const { transport, sessionId } = await createTransport(); + + const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + expect(response.status).toBe(200); + + const reader = response.body!.getReader(); + await vi.advanceTimersByTimeAsync(15000); + const { value } = await reader.read(); + expect(new TextDecoder().decode(value)).toBe(': keepalive\n\n'); + + await transport.close(); + }); + + it('should honor a custom keepAliveMs interval', async () => { + const { transport, sessionId } = await createTransport({ keepAliveMs: 1000 }); + + const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + const reader = response.body!.getReader(); + + await vi.advanceTimersByTimeAsync(3000); + let received = ''; + for (let i = 0; i < 3; i++) { + const { value } = await reader.read(); + received += new TextDecoder().decode(value); + } + expect(received).toBe(': keepalive\n\n'.repeat(3)); + + await transport.close(); + }); + + it('should not write keep-alive frames when keepAliveMs is 0', async () => { + const { transport, sessionId } = await createTransport({ keepAliveMs: 0 }); + + const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + const reader = response.body!.getReader(); + + await vi.advanceTimersByTimeAsync(60000); + const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); + expect(raced).toBe('pending'); + + await transport.close(); + }); + + it.each([0.5, Number.NaN, Number.POSITIVE_INFINITY, 2_147_483_648])( + 'should disable keep-alive for invalid keepAliveMs %s instead of arming a clamped interval', + async keepAliveMs => { + const { transport, sessionId } = await createTransport({ keepAliveMs }); + + const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + expect(response.status).toBe(200); + + // No timer may be armed: setInterval with a NaN/out-of-range delay + // is clamped by Node to ~1ms and would flood the stream with + // keep-alive frames. + expect(vi.getTimerCount()).toBe(0); + const reader = response.body!.getReader(); + await vi.advanceTimersByTimeAsync(60000); + const raced = await Promise.race([reader.read(), Promise.resolve('pending')]); + expect(raced).toBe('pending'); + + await transport.close(); + } + ); + + it('should stop keep-alive frames after the stream is closed', async () => { + const { transport, sessionId } = await createTransport(); + + const response = await transport.handleRequest(createRequest('GET', undefined, { sessionId })); + const reader = response.body!.getReader(); + + await transport.close(); + const { done } = await reader.read(); + expect(done).toBe(true); + + // Advancing time after close must not throw or fire further writes + expect(vi.getTimerCount()).toBe(0); + await vi.advanceTimersByTimeAsync(60000); + }); + + it('should write keep-alive frames on a POST SSE stream while a request is pending', async () => { + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID() }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + let resolveTool: (() => void) | undefined; + mcpServer.registerTool('slow', { description: 'never resolves until released' }, async (): Promise => { + await new Promise(resolve => { + resolveTool = resolve; + }); + return { content: [{ type: 'text', text: 'done' }] }; + }); + await mcpServer.connect(transport); + + const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + + const response = await transport.handleRequest( + createRequest( + 'POST', + { jsonrpc: '2.0', method: 'tools/call', params: { name: 'slow', arguments: {} }, id: 'call-1' } as JSONRPCMessage, + { + sessionId + } + ) + ); + expect(response.status).toBe(200); + const reader = response.body!.getReader(); + + await vi.advanceTimersByTimeAsync(15000); + const { value } = await reader.read(); + expect(new TextDecoder().decode(value)).toBe(': keepalive\n\n'); + + resolveTool?.(); + await transport.close(); + }); +}); + +describe('WebStandardStreamableHTTPServerTransport SSE keep-alive lifecycle', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('should not arm keep-alive when the transport closes during an event-store replay await', async () => { + let releaseReplay: (() => void) | undefined; + const eventStore: EventStore = { + async storeEvent(): Promise { + return 'evt-1'; + }, + async replayEventsAfter(): Promise { + await new Promise(resolve => { + releaseReplay = resolve; + }); + // Resume the standalone GET stream: it skips the + // no-in-flight-request early close unconditionally, so the + // continuation genuinely reaches the keep-alive arm and only + // the closed-transport guard keeps the timer count at zero. + return '_GET_stream'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport); + const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + + // Enter replayEvents and park on the replayEventsAfter await + const pendingGet = transport.handleRequest( + createRequest('GET', undefined, { sessionId, extraHeaders: { 'Last-Event-ID': 'evt-1' } }) + ); + await vi.advanceTimersByTimeAsync(0); + expect(releaseReplay).toBeDefined(); + + // Close the transport mid-await, then let the replay continuation run + await transport.close(); + releaseReplay?.(); + const replayResponse = await pendingGet; + + // The deferred continuation must not have armed a timer close() can never sweep + expect(vi.getTimerCount()).toBe(0); + + // The continuation must not re-register the stream on the closed + // transport: the client observes stream end instead of hanging on a + // dead session, and a later resume isn't 409-blocked by a stale entry. + const { done } = await replayResponse.body!.getReader().read(); + expect(done).toBe(true); + const internals = transport as unknown as { _streamMapping: Map }; + expect(internals._streamMapping.size).toBe(0); + }); + + it('should not leak a keep-alive timer when the priming event write fails on a POST SSE stream', async () => { + // Healthy during initialization, then the store starts failing — the + // tool call's priming event write must reject inside handlePostRequest + let storeFails = false; + const eventStore: EventStore = { + async storeEvent(): Promise { + if (storeFails) { + throw new Error('event store unavailable'); + } + return `evt-${randomUUID()}`; + }, + async replayEventsAfter(): Promise { + return 'stream-1'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + mcpServer.registerTool('noop', { description: 'noop' }, async (): Promise => ({ content: [] })); + await mcpServer.connect(transport); + + const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + expect(initResponse.status).toBe(200); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + // Let the init response finish sending (its send() stores an event and + // then cleans up the init stream's keep-alive) before failing the store + await vi.advanceTimersByTimeAsync(0); + expect(vi.getTimerCount()).toBe(0); + storeFails = true; + + await transport.handleRequest( + createRequest( + 'POST', + { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'call-1' } as JSONRPCMessage, + { + sessionId + } + ) + ); + + // The discarded stream must not carry a permanently-firing timer + expect(vi.getTimerCount()).toBe(0); + + // The stream bookkeeping registered before the failed priming write + // must be reclaimed too: repeated failures during an event-store + // outage must not accrete orphaned stream entries or request mappings. + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + expect(internals._streamMapping.size).toBe(0); + expect(internals._requestToStreamMapping.size).toBe(0); + + await transport.close(); + }); + + it('should not reclaim a successor mapping that reused the failed POST request id', async () => { + let failPriming = false; + let rejectPriming: (() => void) | undefined; + const eventStore: EventStore = { + async storeEvent(): Promise { + if (failPriming) { + return new Promise((_resolve, reject) => { + rejectPriming = () => reject(new Error('event store unavailable')); + }); + } + return `evt-${randomUUID()}`; + }, + async replayEventsAfter(): Promise { + return 'stream-1'; + } + }; + const transport = new WebStandardStreamableHTTPServerTransport({ sessionIdGenerator: () => randomUUID(), eventStore }); + const mcpServer = new McpServer({ name: 'test-server', version: '1.0.0' }); + mcpServer.registerTool('noop', { description: 'noop' }, async (): Promise => ({ content: [] })); + await mcpServer.connect(transport); + + const initResponse = await transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + const sessionId = initResponse.headers.get('mcp-session-id') as string; + await vi.advanceTimersByTimeAsync(0); + failPriming = true; + + const pending = transport.handleRequest( + createRequest( + 'POST', + { jsonrpc: '2.0', method: 'tools/call', params: { name: 'noop', arguments: {} }, id: 'same-id' } as JSONRPCMessage, + { sessionId } + ) + ); + await vi.advanceTimersByTimeAsync(0); + expect(rejectPriming).toBeDefined(); + + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + const failedStreamId = internals._requestToStreamMapping.get('same-id'); + expect(failedStreamId).toBeDefined(); + internals._requestToStreamMapping.set('same-id', 'successor-stream'); + + rejectPriming?.(); + await pending; + + expect(internals._streamMapping.has(failedStreamId!)).toBe(false); + expect(internals._requestToStreamMapping.get('same-id')).toBe('successor-stream'); + internals._requestToStreamMapping.delete('same-id'); + await transport.close(); + }); + + it('should not register a POST stream after close races session initialization', async () => { + let releaseInitialization: (() => void) | undefined; + const transport = new WebStandardStreamableHTTPServerTransport({ + sessionIdGenerator: () => randomUUID(), + onsessioninitialized: async () => { + await new Promise(resolve => { + releaseInitialization = resolve; + }); + } + }); + await new McpServer({ name: 'test-server', version: '1.0.0' }).connect(transport); + + const pendingInit = transport.handleRequest(createRequest('POST', TEST_MESSAGES.initialize)); + await vi.advanceTimersByTimeAsync(0); + expect(releaseInitialization).toBeDefined(); + + await transport.close(); + releaseInitialization?.(); + const response = await pendingInit; + + expect(response.status).toBe(404); + expect(vi.getTimerCount()).toBe(0); + const internals = transport as unknown as { + _streamMapping: Map; + _requestToStreamMapping: Map; + }; + expect(internals._streamMapping.size).toBe(0); + expect(internals._requestToStreamMapping.size).toBe(0); + }); +});