diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 506649fb9..4a249aadb 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 diff --git a/packages/coding-agent/src/mcp/transports/http.ts b/packages/coding-agent/src/mcp/transports/http.ts index 86b22909a..ccec75634 100644 --- a/packages/coding-agent/src/mcp/transports/http.ts +++ b/packages/coding-agent/src/mcp/transports/http.ts @@ -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 { + if (signal.aborted) throw signal.reason; + const { promise, resolve, reject } = Promise.withResolvers(); + 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(); + 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 => { + let current = response; try { - for await (const raw of readSseJson(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 = { + 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; diff --git a/packages/coding-agent/test/mcp-http-transport.test.ts b/packages/coding-agent/test/mcp-http-transport.test.ts index e1f2244df..09915e993 100644 --- a/packages/coding-agent/test/mcp-http-transport.test.ts +++ b/packages/coding-agent/test/mcp-http-transport.test.ts @@ -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("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); + }); +}); diff --git a/packages/utils/CHANGELOG.md b/packages/utils/CHANGELOG.md index be043a33a..a60e0e239 100644 --- a/packages/utils/CHANGELOG.md +++ b/packages/utils/CHANGELOG.md @@ -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 diff --git a/packages/utils/src/stream.ts b/packages/utils/src/stream.ts index 9ae9a7d0c..36ff5380e 100644 --- a/packages/utils/src/stream.ts +++ b/packages/utils/src/stream.ts @@ -281,11 +281,16 @@ export async function* readSseJson( * - `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; } diff --git a/packages/utils/test/stream.test.ts b/packages/utils/test/stream.test.ts index 3b3348782..9306424c7 100644 --- a/packages/utils/test/stream.test.ts +++ b/packages/utils/test/stream.test.ts @@ -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));