fix(mcp): resumed interrupted HTTP response streams
Implemented 2025-11-25 Streamable HTTP polling semantics for POST SSE responses: retain event IDs and retry intervals, wait as instructed, and reconnect with GET plus Last-Event-ID until the originating JSON-RPC response arrives. Extended the shared SSE parser to expose valid id/retry fields and control-only events so reconnecting consumers do not need to reparse raw lines. Fixes #8264
This commit is contained in:
@@ -4,7 +4,7 @@
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed the MCP Streamable HTTP transport never sending the `MCP-Protocol-Version` header and negotiating the stale `2025-03-26` revision, which made spec-current servers (e.g. AWS Bedrock AgentCore Gateway with an outbound per-user OAuth target) reject every `tools/call` with a generic internal error. The client now negotiates `2025-11-25` and echoes the negotiated version in the `MCP-Protocol-Version` header on every request after `initialize` ([#8264](https://github.com/can1357/oh-my-pi/issues/8264)).
|
||||
- Fixed the MCP Streamable HTTP transport never sending the `MCP-Protocol-Version` header and negotiating the stale `2025-03-26` revision, which made spec-current servers (e.g. AWS Bedrock AgentCore Gateway with an outbound per-user OAuth target) reject every `tools/call` with a generic internal error. The client now negotiates `2025-11-25`, echoes the negotiated version on every request after `initialize`, and resumes server-closed POST response streams with `Last-Event-ID` after the requested SSE retry interval ([#8264](https://github.com/can1357/oh-my-pi/issues/8264)).
|
||||
|
||||
## [17.2.13] - 2026-08-11
|
||||
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
* header on every request (see `MCP_PROTOCOL_VERSION`).
|
||||
*/
|
||||
import * as AIError from "@oh-my-pi/pi-ai/error";
|
||||
import { logger, readSseJson } from "@oh-my-pi/pi-utils";
|
||||
import { logger, readSseEvents, readSseJson } from "@oh-my-pi/pi-utils";
|
||||
import type {
|
||||
JsonRpcError,
|
||||
JsonRpcMessage,
|
||||
@@ -23,6 +23,27 @@ import { createMCPTimeout, getNeverAbortSignal, isMCPTimeoutEnabled, resolveMCPT
|
||||
import { type MCPFetchInit, mcpFetch, withoutHeader } from "./header-policy";
|
||||
|
||||
const HTTP_SSE_CONNECT_TIMEOUT_MS = 1_000;
|
||||
const DEFAULT_SSE_RETRY_MS = 3_000;
|
||||
|
||||
interface SSEResumeState {
|
||||
lastEventId: string | null;
|
||||
retryMs: number;
|
||||
}
|
||||
|
||||
/** Wait for the server-provided SSE retry interval while remaining abortable. */
|
||||
async function waitForSSERetry(ms: number, signal: AbortSignal): Promise<void> {
|
||||
if (signal.aborted) throw signal.reason;
|
||||
const { promise, resolve, reject } = Promise.withResolvers<void>();
|
||||
const timer = setTimeout(resolve, ms);
|
||||
const onAbort = (): void => reject(signal.reason);
|
||||
signal.addEventListener("abort", onAbort, { once: true });
|
||||
try {
|
||||
await promise;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
signal.removeEventListener("abort", onAbort);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Best-effort startup deadline for the optional Streamable HTTP GET SSE listener.
|
||||
*
|
||||
@@ -323,40 +344,64 @@ export class HttpTransport implements MCPTransport {
|
||||
const signal = operation.signal ?? getNeverAbortSignal();
|
||||
|
||||
const { promise, resolve, reject } = Promise.withResolvers<T>();
|
||||
const resume: SSEResumeState = { lastEventId: null, retryMs: DEFAULT_SSE_RETRY_MS };
|
||||
let captured = false;
|
||||
|
||||
// Drain the SSE stream from a single iterator. We resolve the deferred
|
||||
// promise as soon as the matching response arrives, then keep iterating
|
||||
// in the background to pick up piggybacked notifications/requests.
|
||||
// Re-reading `response.body` after `for await` breaks would lock the
|
||||
// stream a second time and surface as "ReadableStream already has a
|
||||
// controller", so we must not exit the loop early.
|
||||
// Drain each physical SSE connection without leaving its iterator early.
|
||||
// A server may close a connection without terminating the logical stream;
|
||||
// when it supplied an event ID, resume that stream via GET + Last-Event-ID.
|
||||
const drain = async (): Promise<void> => {
|
||||
let current = response;
|
||||
try {
|
||||
for await (const raw of readSseJson<JsonRpcMessage | JsonRpcMessage[]>(response.body!, signal)) {
|
||||
const messages = Array.isArray(raw) ? raw : [raw];
|
||||
for (const message of messages) {
|
||||
if (
|
||||
!captured &&
|
||||
"id" in message &&
|
||||
message.id === expectedId &&
|
||||
("result" in message || "error" in message)
|
||||
) {
|
||||
captured = true;
|
||||
operation.clear();
|
||||
if (message.error) {
|
||||
reject(new Error(`MCP error ${message.error.code}: ${message.error.message}`));
|
||||
} else {
|
||||
resolve(message.result as T);
|
||||
for (;;) {
|
||||
if (!current.body) throw new Error("SSE response did not include a body");
|
||||
for await (const event of readSseEvents(current.body, signal)) {
|
||||
if (event.id !== undefined) resume.lastEventId = event.id || null;
|
||||
if (event.retry !== undefined) resume.retryMs = event.retry;
|
||||
if (event.data === "") continue;
|
||||
const raw = JSON.parse(event.data) as JsonRpcMessage | JsonRpcMessage[];
|
||||
const messages = Array.isArray(raw) ? raw : [raw];
|
||||
for (const message of messages) {
|
||||
if (
|
||||
!captured &&
|
||||
"id" in message &&
|
||||
message.id === expectedId &&
|
||||
("result" in message || "error" in message)
|
||||
) {
|
||||
captured = true;
|
||||
operation.clear();
|
||||
if (message.error) {
|
||||
reject(new Error(`MCP error ${message.error.code}: ${message.error.message}`));
|
||||
} else {
|
||||
resolve(message.result as T);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
continue;
|
||||
if (!this.#connected) continue;
|
||||
this.#dispatchSSEMessage(message);
|
||||
}
|
||||
if (!this.#connected) continue;
|
||||
this.#dispatchSSEMessage(message);
|
||||
}
|
||||
}
|
||||
if (!captured) {
|
||||
reject(new Error(`No response received for request ID ${expectedId}`));
|
||||
if (captured) return;
|
||||
if (resume.lastEventId === null) {
|
||||
throw new Error(`No response received for request ID ${expectedId}`);
|
||||
}
|
||||
|
||||
await waitForSSERetry(resume.retryMs, signal);
|
||||
const generated: Record<string, string> = {
|
||||
Accept: "text/event-stream",
|
||||
"Last-Event-ID": resume.lastEventId,
|
||||
};
|
||||
if (this.#sessionId) generated["Mcp-Session-Id"] = this.#sessionId;
|
||||
current = await this.#fetch({ method: "GET", signal }, generated);
|
||||
if (!current.ok) {
|
||||
const text = await current.text();
|
||||
throw new Error(`HTTP ${current.status} resuming MCP SSE stream: ${text}`);
|
||||
}
|
||||
const contentType = current.headers.get("Content-Type") ?? "";
|
||||
if (!contentType.includes("text/event-stream")) {
|
||||
await current.body?.cancel();
|
||||
throw new Error(`MCP SSE resume returned unsupported Content-Type: ${contentType || "(missing)"}`);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
if (captured) return;
|
||||
|
||||
@@ -170,3 +170,47 @@ describe("MCP Streamable HTTP protocol version header", () => {
|
||||
expect(seen.post).toBe("2025-11-25");
|
||||
});
|
||||
});
|
||||
|
||||
describe("MCP Streamable HTTP POST response resumption", () => {
|
||||
it("resumes a closed response stream with Last-Event-ID after the requested retry delay", async () => {
|
||||
const observed: {
|
||||
lastEventId: string | null;
|
||||
protocolVersion: string | null;
|
||||
postClosedAt: number;
|
||||
resumedAt: number;
|
||||
} = { lastEventId: null, protocolVersion: null, postClosedAt: 0, resumedAt: 0 };
|
||||
server = Bun.serve({
|
||||
port: 0,
|
||||
fetch(req) {
|
||||
if (req.method === "POST") {
|
||||
observed.postClosedAt = performance.now();
|
||||
return new Response("id: stream-1\nretry: 20\ndata:\n\n", {
|
||||
headers: { "Content-Type": "text/event-stream" },
|
||||
});
|
||||
}
|
||||
observed.resumedAt = performance.now();
|
||||
observed.lastEventId = req.headers.get("Last-Event-ID");
|
||||
observed.protocolVersion = req.headers.get("MCP-Protocol-Version");
|
||||
return new Response(
|
||||
'id: stream-2\ndata: {"jsonrpc":"2.0","id":1,"result":{"tools":[{"name":"resumed","inputSchema":{"type":"object"}}]}}\n\n',
|
||||
{ headers: { "Content-Type": "text/event-stream" } },
|
||||
);
|
||||
},
|
||||
});
|
||||
if (!server) throw new Error("Test server was not started");
|
||||
const transport = new HttpTransport({
|
||||
type: "http",
|
||||
url: `http://127.0.0.1:${server.port}/mcp`,
|
||||
timeout: GUARD_TIMEOUT_MS,
|
||||
});
|
||||
await transport.connect();
|
||||
transport.setProtocolVersion("2025-11-25");
|
||||
|
||||
await expect(withPendingGuard(transport.request<ToolList>("tools/list"), "request")).resolves.toEqual({
|
||||
tools: [{ name: "resumed", inputSchema: { type: "object" } }],
|
||||
});
|
||||
expect(observed.lastEventId).toBe("stream-1");
|
||||
expect(observed.protocolVersion).toBe("2025-11-25");
|
||||
expect(observed.resumedAt - observed.postClosedAt).toBeGreaterThanOrEqual(15);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Changed
|
||||
|
||||
- Extended parsed Server-Sent Events with optional `id` and `retry` fields, including control-only events, so reconnecting transports can retain stream cursors and server-requested retry intervals.
|
||||
|
||||
## [17.2.13] - 2026-08-11
|
||||
|
||||
### Changed
|
||||
|
||||
@@ -281,11 +281,16 @@ export async function* readSseJson<T>(
|
||||
* - `raw` is the list of decoded non-empty lines that made up the event,
|
||||
* preserved for diagnostic context (error reporting, debugging). The
|
||||
* dispatching blank line is not included.
|
||||
* - `id` and `retry` are present only when the event carried valid fields with
|
||||
* those names. Control-only events are yielded so reconnecting transports can
|
||||
* retain the cursor and server-requested retry interval.
|
||||
*/
|
||||
export interface ServerSentEvent {
|
||||
event: string | null;
|
||||
data: string;
|
||||
raw: string[];
|
||||
id?: string;
|
||||
retry?: number;
|
||||
}
|
||||
|
||||
interface SseEventState {
|
||||
@@ -296,6 +301,8 @@ interface SseEventState {
|
||||
// seen yet" (distinct from a `data:` field with an empty value).
|
||||
data: string | null;
|
||||
raw: string[];
|
||||
id?: string;
|
||||
retry?: number;
|
||||
}
|
||||
|
||||
// Complete lines are decoded in one batch per source chunk. Each batch ends on
|
||||
@@ -303,7 +310,7 @@ interface SseEventState {
|
||||
const SSE_DECODER = new TextDecoder("utf-8");
|
||||
|
||||
function flushSseEvent(state: SseEventState): ServerSentEvent | null {
|
||||
if (state.event === null && state.data === null) {
|
||||
if (state.event === null && state.data === null && state.id === undefined && state.retry === undefined) {
|
||||
state.raw = [];
|
||||
return null;
|
||||
}
|
||||
@@ -312,9 +319,13 @@ function flushSseEvent(state: SseEventState): ServerSentEvent | null {
|
||||
data: state.data ?? "",
|
||||
raw: state.raw,
|
||||
};
|
||||
if (state.id !== undefined) event.id = state.id;
|
||||
if (state.retry !== undefined) event.retry = state.retry;
|
||||
state.event = null;
|
||||
state.data = null;
|
||||
state.raw = [];
|
||||
state.id = undefined;
|
||||
state.retry = undefined;
|
||||
return event;
|
||||
}
|
||||
|
||||
@@ -348,9 +359,22 @@ function pushSseLine(line: string, state: SseEventState): ServerSentEvent | null
|
||||
state.data += "\n";
|
||||
state.data += value;
|
||||
}
|
||||
} else if (fieldName === "id") {
|
||||
if (!value.includes("\0")) state.id = value;
|
||||
} else if (fieldName === "retry" && value.length > 0) {
|
||||
let valid = true;
|
||||
for (let index = 0; index < value.length; index++) {
|
||||
const code = value.charCodeAt(index);
|
||||
if (code < 0x30 || code > 0x39) {
|
||||
valid = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (valid) {
|
||||
const retry = Number(value);
|
||||
if (Number.isSafeInteger(retry)) state.retry = retry;
|
||||
}
|
||||
}
|
||||
// `id` and `retry` are intentionally ignored — the providers we consume
|
||||
// don't use them, and the underlying transport handles reconnects itself.
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
@@ -396,6 +396,21 @@ describe("readSseEvents", () => {
|
||||
expect(evt.raw).toEqual(["event: ping", "data: ok"]);
|
||||
});
|
||||
|
||||
it("yields control-only id/retry events for reconnecting transports", async () => {
|
||||
const stream = bytesStreamFromChunks([encoder.encode("id: stream-1\nretry: 25\n\n")]);
|
||||
const events = await collectAsync(readSseEvents(stream));
|
||||
|
||||
expect(events).toEqual([
|
||||
{
|
||||
event: null,
|
||||
data: "",
|
||||
raw: ["id: stream-1", "retry: 25"],
|
||||
id: "stream-1",
|
||||
retry: 25,
|
||||
},
|
||||
] satisfies ServerSentEvent[]);
|
||||
});
|
||||
|
||||
it("strips a single optional space after the field colon (and only one)", async () => {
|
||||
const stream = bytesStreamFromChunks([encoder.encode("event: spaced\ndata: body\n\n")]);
|
||||
const [evt] = await collectAsync(readSseEvents(stream));
|
||||
|
||||
Reference in New Issue
Block a user