@@ -19,6 +19,8 @@ import {
1919 SUPPORTED_PROTOCOL_VERSIONS
2020} from '@modelcontextprotocol/core-internal' ;
2121
22+ import { armSseKeepAlive , DEFAULT_SSE_KEEP_ALIVE_MS } from './sseKeepAlive' ;
23+
2224export type StreamId = string ;
2325export type EventId = string ;
2426
@@ -148,6 +150,12 @@ export interface WebStandardStreamableHTTPServerTransportOptions {
148150 */
149151 retryInterval ?: number ;
150152
153+ /**
154+ * Interval in milliseconds between SSE keep-alive comment frames.
155+ * Defaults to `15000`; set to `0` to disable.
156+ */
157+ keepAliveMs ?: number ;
158+
151159 /**
152160 * List of protocol versions that this transport will accept.
153161 * Used to validate the `mcp-protocol-version` header in incoming requests.
@@ -247,6 +255,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
247255 private _enableDnsRebindingProtection : boolean ;
248256 private _retryInterval ?: number ;
249257 private _supportedProtocolVersions : string [ ] ;
258+ private _keepAliveMs : number ;
250259
251260 sessionId ?: string ;
252261 onclose ?: ( ) => void ;
@@ -264,6 +273,23 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
264273 this . _enableDnsRebindingProtection = options . enableDnsRebindingProtection ?? false ;
265274 this . _retryInterval = options . retryInterval ;
266275 this . _supportedProtocolVersions = options . supportedProtocolVersions ?? SUPPORTED_PROTOCOL_VERSIONS ;
276+ this . _keepAliveMs = options . keepAliveMs ?? DEFAULT_SSE_KEEP_ALIVE_MS ;
277+ }
278+
279+ private startKeepAlive (
280+ controller : ReadableStreamDefaultController < Uint8Array > ,
281+ encoder : InstanceType < typeof TextEncoder >
282+ ) : ReturnType < typeof setInterval > | undefined {
283+ if ( this . _closed ) return undefined ;
284+
285+ const timer = armSseKeepAlive ( this . _keepAliveMs , ( ) => {
286+ try {
287+ controller . enqueue ( encoder . encode ( ': keepalive\n\n' ) ) ;
288+ } catch {
289+ if ( timer !== undefined ) clearInterval ( timer ) ;
290+ }
291+ } ) ;
292+ return timer ;
267293 }
268294
269295 /**
@@ -462,13 +488,17 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
462488
463489 const encoder = new TextEncoder ( ) ;
464490 let streamController : ReadableStreamDefaultController < Uint8Array > ;
491+ // Captured by cancel/cleanup before it is assigned after stream setup.
492+ // eslint-disable-next-line prefer-const
493+ let keepAliveTimer : ReturnType < typeof setInterval > | undefined ;
465494
466495 // Create a ReadableStream with a controller we can use to push SSE events
467496 const readable = new ReadableStream < Uint8Array > ( {
468497 start : controller => {
469498 streamController = controller ;
470499 } ,
471500 cancel : ( ) => {
501+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
472502 // Stream was cancelled by client. Only drop the mapping when
473503 // it still points at THIS controller — a stale cancel must not
474504 // delete a successor stream registered by a later GET/resume.
@@ -481,7 +511,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
481511 const headers : Record < string , string > = {
482512 'Content-Type' : 'text/event-stream' ,
483513 'Cache-Control' : 'no-cache, no-transform' ,
484- Connection : 'keep-alive'
514+ Connection : 'keep-alive' ,
515+ 'X-Accel-Buffering' : 'no'
485516 } ;
486517
487518 // After initialization, always include the session ID if we have one
@@ -494,6 +525,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
494525 controller : streamController ! ,
495526 encoder,
496527 cleanup : ( ) => {
528+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
497529 this . _streamMapping . delete ( this . _standaloneSseStreamId ) ;
498530 try {
499531 streamController ! . close ( ) ;
@@ -503,6 +535,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
503535 }
504536 } ) ;
505537
538+ keepAliveTimer = this . startKeepAlive ( streamController ! , encoder ) ;
506539 return new Response ( readable , { headers } ) ;
507540 }
508541
@@ -537,7 +570,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
537570 const headers : Record < string , string > = {
538571 'Content-Type' : 'text/event-stream' ,
539572 'Cache-Control' : 'no-cache, no-transform' ,
540- Connection : 'keep-alive'
573+ Connection : 'keep-alive' ,
574+ 'X-Accel-Buffering' : 'no'
541575 } ;
542576
543577 if ( this . sessionId !== undefined ) {
@@ -547,6 +581,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
547581 // Create a ReadableStream with controller for SSE
548582 const encoder = new TextEncoder ( ) ;
549583 let streamController : ReadableStreamDefaultController < Uint8Array > ;
584+ let keepAliveTimer : ReturnType < typeof setInterval > | undefined ;
585+ let cancelled = false ;
550586 // Captured by the cancel closure below before it's assigned (after
551587 // replayEventsAfter resolves) — must be `let`.
552588 // eslint-disable-next-line prefer-const
@@ -557,6 +593,8 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
557593 streamController = controller ;
558594 } ,
559595 cancel : ( ) => {
596+ cancelled = true ;
597+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
560598 // Stream was cancelled by client — drop the mapping so a
561599 // subsequent reconnect with the same Last-Event-ID is not
562600 // refused with 409 by the conflict check above. Only delete
@@ -585,11 +623,22 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
585623 }
586624 } ) ;
587625
626+ if ( this . _closed || cancelled ) {
627+ try {
628+ streamController ! . close ( ) ;
629+ } catch {
630+ // Controller already closed/cancelled.
631+ }
632+ return this . createJsonErrorResponse ( 404 , - 32_001 , 'Session not found' ) ;
633+ }
634+
635+ this . _streamMapping . get ( replayedStreamId ) ?. cleanup ( ) ;
588636 this . _streamMapping . set ( replayedStreamId , {
589637 controller : streamController ! ,
590638 encoder,
591639 replayedEventIds,
592640 cleanup : ( ) => {
641+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
593642 this . _streamMapping . delete ( replayedStreamId ! ) ;
594643 try {
595644 streamController ! . close ( ) ;
@@ -618,6 +667,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
618667 }
619668 }
620669
670+ if ( this . _streamMapping . get ( replayedStreamId ) ?. controller === streamController ! ) {
671+ keepAliveTimer = this . startKeepAlive ( streamController ! , encoder ) ;
672+ }
621673 return new Response ( readable , { headers } ) ;
622674 } catch ( error ) {
623675 this . onerror ?.( error as Error ) ;
@@ -818,12 +870,14 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
818870 // SSE streaming mode - use ReadableStream with controller for more reliable data pushing
819871 const encoder = new TextEncoder ( ) ;
820872 let streamController : ReadableStreamDefaultController < Uint8Array > ;
873+ let keepAliveTimer : ReturnType < typeof setInterval > | undefined ;
821874
822875 const readable = new ReadableStream < Uint8Array > ( {
823876 start : controller => {
824877 streamController = controller ;
825878 } ,
826879 cancel : ( ) => {
880+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
827881 // Stream was cancelled by client. Only drop the mapping
828882 // when it still points at THIS controller — a stale cancel
829883 // (firing after a Last-Event-ID reconnect registered a
@@ -837,8 +891,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
837891
838892 const headers : Record < string , string > = {
839893 'Content-Type' : 'text/event-stream' ,
840- 'Cache-Control' : 'no-cache' ,
841- Connection : 'keep-alive'
894+ 'Cache-Control' : 'no-cache, no-transform' ,
895+ Connection : 'keep-alive' ,
896+ 'X-Accel-Buffering' : 'no'
842897 } ;
843898
844899 // After initialization, always include the session ID if we have one
@@ -854,6 +909,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
854909 controller : streamController ! ,
855910 encoder,
856911 cleanup : ( ) => {
912+ if ( keepAliveTimer !== undefined ) clearInterval ( keepAliveTimer ) ;
857913 this . _streamMapping . delete ( streamId ) ;
858914 try {
859915 streamController ! . close ( ) ;
@@ -891,6 +947,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
891947 // The server SHOULD NOT close the SSE stream before sending all JSON-RPC responses
892948 // This will be handled by the send() method when responses are ready
893949
950+ if ( this . _streamMapping . get ( streamId ) ?. controller === streamController ! ) {
951+ keepAliveTimer = this . startKeepAlive ( streamController ! , encoder ) ;
952+ }
894953 return new Response ( readable , { status : 200 , headers } ) ;
895954 } catch ( error ) {
896955 // return JSON-RPC formatted error
0 commit comments