Skip to content

Commit b678f44

Browse files
refactor(server): start keepalive inside the stream start callback
Avoids new non-null assertions on the stream controller and makes startKeepAlive always return a stop function (no-op when disabled), so cleanup sites call it unconditionally. Also drops a redundant comment in buildContext.
1 parent c2f32f0 commit b678f44

2 files changed

Lines changed: 14 additions & 19 deletions

File tree

packages/server/src/server/server.ts

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -160,9 +160,6 @@ export class Server extends Protocol<ServerContext> {
160160
mcpReq: {
161161
...ctx.mcpReq,
162162
log: (level, data, logger) => this.sendLoggingMessage({ level, data, logger }),
163-
// Associate the nested server→client request with the request being handled so
164-
// transports can route it onto the originating request's response stream
165-
// (servers must only send these requests in association with a client request).
166163
elicitInput: (params, options) => this.elicitInput(params, { relatedRequestId: ctx.mcpReq.id, ...options }),
167164
requestSampling: (params, options) => this.createMessage(params, { relatedRequestId: ctx.mcpReq.id, ...options })
168165
},

packages/server/src/server/streamableHttp.ts

Lines changed: 14 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,9 @@ import {
1818
SUPPORTED_PROTOCOL_VERSIONS
1919
} from '@modelcontextprotocol/core';
2020

21+
/** Placeholder stop-function used until a stream's keepalive is started (or when keepalive is disabled). */
22+
const noKeepAlive = (): void => {};
23+
2124
export type StreamId = string;
2225
export type EventId = string;
2326

@@ -453,11 +456,13 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
453456

454457
const encoder = new TextEncoder();
455458
let streamController: ReadableStreamDefaultController<Uint8Array>;
459+
let stopKeepAlive: () => void = noKeepAlive;
456460

457461
// Create a ReadableStream with a controller we can use to push SSE events
458462
const readable = new ReadableStream<Uint8Array>({
459463
start: controller => {
460464
streamController = controller;
465+
stopKeepAlive = this.startKeepAlive(controller, encoder);
461466
},
462467
cancel: () => {
463468
// Stream was cancelled by client
@@ -476,15 +481,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
476481
headers['mcp-session-id'] = this.sessionId;
477482
}
478483

479-
// Keepalive stops via cleanup(); after a client cancel it self-clears on the next write attempt.
480-
const stopKeepAlive = this.startKeepAlive(streamController!, encoder);
481-
482484
// Store the stream mapping with the controller for pushing data
483485
this._streamMapping.set(this._standaloneSseStreamId, {
484486
controller: streamController!,
485487
encoder,
486488
cleanup: () => {
487-
stopKeepAlive?.();
489+
stopKeepAlive();
488490
this._streamMapping.delete(this._standaloneSseStreamId);
489491
try {
490492
streamController!.close();
@@ -538,10 +540,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
538540
// Create a ReadableStream with controller for SSE
539541
const encoder = new TextEncoder();
540542
let streamController: ReadableStreamDefaultController<Uint8Array>;
543+
let stopKeepAlive: () => void = noKeepAlive;
541544

542545
const readable = new ReadableStream<Uint8Array>({
543546
start: controller => {
544547
streamController = controller;
548+
stopKeepAlive = this.startKeepAlive(controller, encoder);
545549
},
546550
cancel: () => {
547551
// Stream was cancelled by client
@@ -563,13 +567,11 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
563567
}
564568
});
565569

566-
const stopKeepAlive = this.startKeepAlive(streamController!, encoder);
567-
568570
this._streamMapping.set(replayedStreamId, {
569571
controller: streamController!,
570572
encoder,
571573
cleanup: () => {
572-
stopKeepAlive?.();
574+
stopKeepAlive();
573575
this._streamMapping.delete(replayedStreamId);
574576
try {
575577
streamController!.close();
@@ -591,12 +593,9 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
591593
* Returns a stop function, or `undefined` when keepalive is not configured or disabled
592594
* (non-positive interval).
593595
*/
594-
private startKeepAlive(
595-
controller: ReadableStreamDefaultController<Uint8Array>,
596-
encoder: InstanceType<typeof TextEncoder>
597-
): (() => void) | undefined {
596+
private startKeepAlive(controller: ReadableStreamDefaultController<Uint8Array>, encoder: InstanceType<typeof TextEncoder>): () => void {
598597
if (this._keepAliveInterval === undefined || this._keepAliveInterval <= 0) {
599-
return undefined;
598+
return () => {};
600599
}
601600
const timer = setInterval(() => {
602601
try {
@@ -792,10 +791,12 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
792791
// SSE streaming mode - use ReadableStream with controller for more reliable data pushing
793792
const encoder = new TextEncoder();
794793
let streamController: ReadableStreamDefaultController<Uint8Array>;
794+
let stopKeepAlive: () => void = noKeepAlive;
795795

796796
const readable = new ReadableStream<Uint8Array>({
797797
start: controller => {
798798
streamController = controller;
799+
stopKeepAlive = this.startKeepAlive(controller, encoder);
799800
},
800801
cancel: () => {
801802
// Stream was cancelled by client
@@ -814,9 +815,6 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
814815
headers['mcp-session-id'] = this.sessionId;
815816
}
816817

817-
// Keepalive stops via cleanup(); after a client cancel it self-clears on the next write attempt.
818-
const stopKeepAlive = this.startKeepAlive(streamController!, encoder);
819-
820818
// Store the response for this request to send messages back through this connection
821819
// We need to track by request ID to maintain the connection
822820
for (const message of messages) {
@@ -825,7 +823,7 @@ export class WebStandardStreamableHTTPServerTransport implements Transport {
825823
controller: streamController!,
826824
encoder,
827825
cleanup: () => {
828-
stopKeepAlive?.();
826+
stopKeepAlive();
829827
this._streamMapping.delete(streamId);
830828
try {
831829
streamController!.close();

0 commit comments

Comments
 (0)