Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/quiet-stream-provenance.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@modelcontextprotocol/client": patch
---

Expose the originating client request ID for server-initiated requests received on a Streamable HTTP response stream.
22 changes: 18 additions & 4 deletions packages/client/src/client/streamableHttp.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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.
*
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}
Expand All @@ -794,6 +804,7 @@ export class StreamableHTTPClientTransport implements Transport {
resumptionToken: lastEventId,
onresumptiontoken,
replayMessageId,
relatedRequestId,
requestSignal,
onRequestStreamEnd
},
Expand Down Expand Up @@ -827,6 +838,7 @@ export class StreamableHTTPClientTransport implements Transport {
resumptionToken: lastEventId,
onresumptiontoken,
replayMessageId,
relatedRequestId,
requestSignal,
onRequestStreamEnd
},
Expand Down Expand Up @@ -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');
Expand All @@ -1135,6 +1148,7 @@ export class StreamableHTTPClientTransport implements Transport {
response.body,
{
onresumptiontoken,
relatedRequestId,
requestSignal: options?.requestSignal,
onRequestStreamEnd: options?.onRequestStreamEnd
},
Expand Down
30 changes: 30 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions packages/core-internal/src/types/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
Loading