From d1584c44f838f4f23a22d5b841ef0f3c5f7911a3 Mon Sep 17 00:00:00 2001 From: edenbuilds <279970382+edenbuilds@users.noreply.github.com> Date: Sun, 16 Aug 2026 05:25:27 +0530 Subject: [PATCH] fix(client): preserve Streamable HTTP request provenance (#2659) Expose the originating client request ID when a server-initiated JSON-RPC request arrives on that request's SSE response stream. Fixes #2659 --- .changeset/quiet-stream-provenance.md | 5 ++++ packages/client/src/client/streamableHttp.ts | 22 +++++++++++--- .../client/test/client/streamableHttp.test.ts | 30 +++++++++++++++++++ packages/core-internal/src/types/types.ts | 7 +++++ 4 files changed, 60 insertions(+), 4 deletions(-) create mode 100644 .changeset/quiet-stream-provenance.md diff --git a/.changeset/quiet-stream-provenance.md b/.changeset/quiet-stream-provenance.md new file mode 100644 index 0000000000..b04cb9c2f6 --- /dev/null +++ b/.changeset/quiet-stream-provenance.md @@ -0,0 +1,5 @@ +--- +"@modelcontextprotocol/client": patch +--- + +Expose the originating client request ID for server-initiated requests received on a Streamable HTTP response stream. diff --git a/packages/client/src/client/streamableHttp.ts b/packages/client/src/client/streamableHttp.ts index ace0663158..9fcb02fa17 100644 --- a/packages/client/src/client/streamableHttp.ts +++ b/packages/client/src/client/streamableHttp.ts @@ -1,6 +1,6 @@ import type { ReadableWritablePair } from 'node:stream/web'; -import type { FetchLike, JSONRPCMessage, Transport } from '@modelcontextprotocol/core-internal'; +import type { FetchLike, JSONRPCMessage, MessageExtraInfo, RequestId, Transport } from '@modelcontextprotocol/core-internal'; import { createFetchWithInit, encodeMcpParamValue, @@ -54,6 +54,12 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp * Options for starting or authenticating an SSE connection */ export interface StartSSEOptions { + /** + * The client request whose POST response stream carries this SSE stream. + * Standalone GET streams have no related request. + */ + relatedRequestId?: RequestId; + /** * The resumption token used to continue long-running requests that were interrupted. * @@ -330,7 +336,7 @@ export class StreamableHTTPClientTransport implements Transport { onclose?: () => void; onerror?: (error: Error) => void; - onmessage?: (message: JSONRPCMessage) => void; + onmessage?: (message: JSONRPCMessage, extra?: MessageExtraInfo) => void; /** * Streamable HTTP opens one POST (and SSE response stream) per outbound @@ -713,7 +719,7 @@ export class StreamableHTTPClientTransport implements Transport { options.onRequestStreamEnd?.(); return; } - const { onresumptiontoken, replayMessageId, requestSignal, onRequestStreamEnd } = options; + const { onresumptiontoken, replayMessageId, relatedRequestId, requestSignal, onRequestStreamEnd } = options; // An intentional abort — transport-wide close OR a per-request abort // (McpSubscription.close() aborting its `requestSignal`) — must read as // a clean shutdown: no misleading "SSE stream disconnected" onerror, @@ -775,7 +781,11 @@ export class StreamableHTTPClientTransport implements Transport { message.id = replayMessageId; } } - this.onmessage?.(message); + if (relatedRequestId === undefined || !isJSONRPCRequest(message)) { + this.onmessage?.(message); + } else { + this.onmessage?.(message, { relatedRequestId }); + } } catch (error) { this.onerror?.(error as Error); } @@ -794,6 +804,7 @@ export class StreamableHTTPClientTransport implements Transport { resumptionToken: lastEventId, onresumptiontoken, replayMessageId, + relatedRequestId, requestSignal, onRequestStreamEnd }, @@ -827,6 +838,7 @@ export class StreamableHTTPClientTransport implements Transport { resumptionToken: lastEventId, onresumptiontoken, replayMessageId, + relatedRequestId, requestSignal, onRequestStreamEnd }, @@ -1121,6 +1133,7 @@ export class StreamableHTTPClientTransport implements Transport { const messages = Array.isArray(message) ? message : [message]; const hasRequests = messages.some(msg => 'method' in msg && 'id' in msg && msg.id !== undefined); + const relatedRequestId = messages.length === 1 && isJSONRPCRequest(messages[0]) ? messages[0].id : undefined; // Check the response type (parsed media type — see mediaTypeEssence) const contentType = response.headers.get('content-type'); @@ -1135,6 +1148,7 @@ export class StreamableHTTPClientTransport implements Transport { response.body, { onresumptiontoken, + relatedRequestId, requestSignal: options?.requestSignal, onRequestStreamEnd: options?.onRequestStreamEnd }, diff --git a/packages/client/test/client/streamableHttp.test.ts b/packages/client/test/client/streamableHttp.test.ts index a36bbc0ad3..3af65c8204 100644 --- a/packages/client/test/client/streamableHttp.test.ts +++ b/packages/client/test/client/streamableHttp.test.ts @@ -524,6 +524,36 @@ describe('StreamableHTTPClientTransport', () => { ).toBe(true); }); + it('attributes server requests received on a POST SSE stream to the originating request', async () => { + const encoder = new TextEncoder(); + const stream = new ReadableStream({ + start(controller) { + controller.enqueue( + encoder.encode( + 'event: message\ndata: {"jsonrpc":"2.0","id":"elicitation-1","method":"elicitation/create","params":{}}\n\n' + ) + ); + } + }); + + (globalThis.fetch as Mock).mockResolvedValueOnce({ + ok: true, + status: 200, + headers: new Headers({ 'content-type': 'text/event-stream' }), + body: stream + }); + + const messageSpy = vi.fn(); + transport.onmessage = messageSpy; + + await transport.send({ jsonrpc: '2.0', id: 'tool-call-1', method: 'tools/call', params: { name: 'route' } }); + await new Promise(resolve => setTimeout(resolve, 50)); + + expect(messageSpy).toHaveBeenCalledWith(expect.objectContaining({ id: 'elicitation-1', method: 'elicitation/create' }), { + relatedRequestId: 'tool-call-1' + }); + }); + it('declares hasPerRequestStream so the protocol layer routes 2026-era cancellation to stream-close', () => { // Spec basic/patterns/cancellation §Transport-Specific (2026-07-28): // closing the per-request SSE stream IS the cancel signal on diff --git a/packages/core-internal/src/types/types.ts b/packages/core-internal/src/types/types.ts index f2bc9d67fc..6ac0ecc32b 100644 --- a/packages/core-internal/src/types/types.ts +++ b/packages/core-internal/src/types/types.ts @@ -888,6 +888,13 @@ export interface MessageClassification { * Extra information about a message. */ export interface MessageExtraInfo { + /** + * The client request whose Streamable HTTP response stream carried this message. + * Set by the client transport for server-initiated requests received on a + * per-request SSE stream; absent for the standalone GET stream. + */ + relatedRequestId?: RequestId; + /** * The original HTTP request. */