From c8df2ac07876b764c1d51468b13e39d48c0239ce Mon Sep 17 00:00:00 2001 From: fforres Date: Sat, 26 Sep 2026 00:55:16 -0700 Subject: [PATCH] chore(observability): record why an MCP request bridge ended abnormally Requests fail with "session was reset" whenever the front worker's bridge WebSocket to the session Durable Object closes without its own "SSE response delivered" code. A real object reset and a resumed stream superseding the bridge both land there, and the log could not tell them apart. Log the close code, reason and wasClean (or the error) with mcp_post_stream_lost. Also drops empty .bun-tag marker entries that bun patch had recorded in the patch. Claude-Session: https://claude.ai/code/session_01VMqJkxznzTQaFVHHcttxpJ --- patches/agents@0.17.3.patch | 76 +++++++++++++++---------------------- 1 file changed, 31 insertions(+), 45 deletions(-) diff --git a/patches/agents@0.17.3.patch b/patches/agents@0.17.3.patch index 7eec44acf0..f785f07484 100644 --- a/patches/agents@0.17.3.patch +++ b/patches/agents@0.17.3.patch @@ -1,18 +1,3 @@ -diff --git a/node_modules/agents/.bun-tag-37e67b70862be2b2 b/.bun-tag-37e67b70862be2b2 -new file mode 100644 -index 0000000000000000000000000000000000000000..e69de29bb2d1d6434b8b29ae775ad8c2e48c5391 -diff --git a/node_modules/agents/.bun-tag-61b2f1517ab5ced4 b/.bun-tag-61b2f1517ab5ced4 -new file mode 100644 -index 0000000000000000000000000000000000000000..e69de29bb2d1d6434b8b29ae775ad8c2e48c5391 -diff --git a/node_modules/agents/.bun-tag-61d29b56a6079f1d b/.bun-tag-61d29b56a6079f1d -new file mode 100644 -index 0000000000000000000000000000000000000000..e69de29bb2d1d6434b8b29ae775ad8c2e48c5391 -diff --git a/node_modules/agents/.bun-tag-a6a47855632a2623 b/.bun-tag-a6a47855632a2623 -new file mode 100644 -index 0000000000000000000000000000000000000000..e69de29bb2d1d6434b8b29ae775ad8c2e48c5391 -diff --git a/node_modules/agents/.bun-tag-c0c639aa2299e502 b/.bun-tag-c0c639aa2299e502 -new file mode 100644 -index 0000000000000000000000000000000000000000..e69de29bb2d1d6434b8b29ae775ad8c2e48c5391 diff --git a/dist/agent-tool-types-CNyE1iz_.d.ts b/dist/agent-tool-types-CNyE1iz_.d.ts index 571eececebd5a1eaf7f2fbf5278801c4d34728ba..010811ae74ef2d27f12f7fb2506848f1939d91d6 100644 --- a/dist/agent-tool-types-CNyE1iz_.d.ts @@ -95,7 +80,7 @@ index c8fad448e8797b89690a99d93490d1363851b225..d80f66f1532c29dda0a0477dc1ec9ef5 McpAgent, type McpAuthContext, diff --git a/dist/mcp/index.js b/dist/mcp/index.js -index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c28786cca392 100644 +index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed0768462f640 100644 --- a/dist/mcp/index.js +++ b/dist/mcp/index.js @@ -28,13 +28,60 @@ import { WebStandardStreamableHTTPServerTransport } from "@modelcontextprotocol/ @@ -196,7 +181,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 request.headers.forEach((value, key) => { existingHeaders[key] = value; }); -@@ -206,47 +273,159 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -206,47 +273,160 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { jsonrpc: "2.0" }); return new Response(body, { status: 500 }); @@ -276,7 +261,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 + * The happy path is byte-identical: every id has been reported + * responded by then, so the set is empty and nothing is written. + */ -+ const __finishAbnormally = async (reason) => { ++ const __finishAbnormally = async (reason, detail = {}) => { + if (__finished) return; + __finished = true; + clearTimeout(__responseDeadline); @@ -286,7 +271,8 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 + event: "mcp_post_stream_lost", + outstandingCount: __outstanding.size, + reason, -+ sessionId ++ sessionId, ++ ...detail + })); + for (const id of __outstanding.values()) __forwardSse(encoder.encode(`event: message\ndata: ${JSON.stringify(sessionResetErrorResponse(id, reason))}\n\n`)); + __outstanding.clear(); @@ -351,7 +337,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 + async function onError(_error) { + // The DO end of the bridge failed. Nothing still + // outstanding will ever be answered. -+ await __finishAbnormally("session_reset"); ++ await __finishAbnormally("session_reset", { bridge: "error", error: String(_error?.message ?? _error ?? "") }); } onError(error).catch(console.error); }); @@ -376,14 +362,14 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 + await writer.close().catch(() => {}); + return; + } -+ await __finishAbnormally("session_reset"); ++ await __finishAbnormally("session_reset", { bridge: "close", closeCode: closeEvent?.code ?? null, closeReason: closeEvent?.reason ?? null, wasClean: closeEvent?.wasClean ?? null }); } - onClose().catch(console.error); + onClose(event).catch(console.error); }); return new Response(readable, { headers: { -@@ -279,10 +458,16 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -279,10 +459,16 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { id: null, jsonrpc: "2.0" }), { status: 400 }); @@ -404,7 +390,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 props: ctx.props, jurisdiction: options.jurisdiction }); -@@ -306,27 +491,116 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -306,27 +492,116 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { if (!ws) { await writer.close(); return new Response("Failed to establish WS to DO", { status: 500 }); @@ -536,7 +522,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 return new Response(readable, { headers: { "Cache-Control": "no-cache", -@@ -389,10 +663,16 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { +@@ -389,10 +664,16 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { const url = new URL(request.url); if (request.method === "GET" && basePattern.test(url)) { const sessionId = url.searchParams.get("sessionId") || namespace.newUniqueId().toString(); @@ -557,7 +543,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 endpointUrl.pathname = encodeURI(`${basePath}/message`); endpointUrl.searchParams.set("sessionId", sessionId); const endpointMessage = `event: endpoint\ndata: ${endpointUrl.pathname + endpointUrl.search + endpointUrl.hash}\n\n`; -@@ -414,35 +694,94 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { +@@ -414,35 +695,94 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { console.error("Failed to establish WebSocket connection"); await writer.close(); return new Response("Failed to establish WebSocket connection", { status: 500 }); @@ -670,7 +656,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 console.error("Error closing SSE connection:", error); } } -@@ -586,6 +925,8 @@ var StreamableHTTPServerTransport = class { +@@ -586,6 +926,8 @@ var StreamableHTTPServerTransport = class { constructor(options) { this._started = false; this._streamResponseIds = /* @__PURE__ */ new Map(); @@ -679,7 +665,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 const { agent } = getCurrentAgent(); if (!agent) throw new Error("McpAgent was not found in Transport constructor"); this._agent = agent; -@@ -627,23 +968,146 @@ var StreamableHTTPServerTransport = class { +@@ -627,23 +969,146 @@ var StreamableHTTPServerTransport = class { const resumedStreamId = await this._eventStore.getStreamIdForEventId?.(lastEventId); if (resumedStreamId) { const resumeState = { streamId: resumedStreamId }; @@ -832,7 +818,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 } /** * Close any connection (other than `selfId`) currently bound to -@@ -651,6 +1115,9 @@ var StreamableHTTPServerTransport = class { +@@ -651,6 +1116,9 @@ var StreamableHTTPServerTransport = class { * Closing rather than mutating sibling state mirrors how the SDK's * single `_streamMapping` entry gives last-writer-wins for free, and * keeps `send()` from routing to a stale bridge. @@ -842,7 +828,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 */ supersedePriorStreamConnections(agent, selfId, streamId) { for (const other of agent.getConnections()) { -@@ -664,12 +1131,14 @@ var StreamableHTTPServerTransport = class { +@@ -664,12 +1132,14 @@ var StreamableHTTPServerTransport = class { * Only used when resumability is enabled */ async replayEvents(lastEventId) { @@ -858,7 +844,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 this.writeSSEEvent(connection, message, eventId); } catch (error) { this.onerror?.(error); -@@ -678,6 +1147,45 @@ var StreamableHTTPServerTransport = class { +@@ -678,6 +1148,45 @@ var StreamableHTTPServerTransport = class { } catch (error) { this.onerror?.(error); } @@ -904,7 +890,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 } /** * Writes an event to the SSE stream with proper formatting -@@ -689,10 +1197,69 @@ var StreamableHTTPServerTransport = class { +@@ -689,10 +1198,69 @@ var StreamableHTTPServerTransport = class { return connection.send(JSON.stringify({ type: "cf_mcp_agent_event", event: eventData, @@ -974,7 +960,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 * Handles POST requests containing JSON-RPC messages */ async handlePostRequest(req, parsedBody) { -@@ -733,6 +1300,22 @@ var StreamableHTTPServerTransport = class { +@@ -733,6 +1301,22 @@ var StreamableHTTPServerTransport = class { }; connection.setState(postState); if (this._eventStore) await agent.setStreamRequestIds(streamId, requestIds); @@ -997,7 +983,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 for (const message of messages) { if (this.messageInterceptor) { if (await this.messageInterceptor(message, { -@@ -760,7 +1343,22 @@ var StreamableHTTPServerTransport = class { +@@ -760,7 +1344,22 @@ var StreamableHTTPServerTransport = class { * when the originating WS has dropped. */ async sendOnStream(agent, streamId, relatedIds, liveConnection, message, requestId) { @@ -1021,7 +1007,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 let shouldClose = false; if (isJSONRPCResultResponse(message) || isJSONRPCErrorResponse(message)) { let responseIds = this._streamResponseIds.get(streamId); -@@ -777,9 +1375,11 @@ var StreamableHTTPServerTransport = class { +@@ -777,9 +1376,11 @@ var StreamableHTTPServerTransport = class { } catch (error) { this.onerror?.(error); } @@ -1035,7 +1021,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 } } async send(message, options) { -@@ -798,14 +1398,19 @@ var StreamableHTTPServerTransport = class { +@@ -798,14 +1399,19 @@ var StreamableHTTPServerTransport = class { * * Sent on exactly one stream, per MCP: "the server MUST send each of * its JSON-RPC messages on only one of the connected streams; it MUST @@ -1059,7 +1045,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 if (standalone) this.writeSSEEvent(standalone, message, eventId); } /** -@@ -861,12 +1466,10 @@ var StreamableHTTPServerTransport = class { +@@ -861,12 +1467,10 @@ var StreamableHTTPServerTransport = class { * * ## Lifecycle * @@ -1076,7 +1062,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 * * Standalone GET stream events (`_GET_stream`) are *not* cleared * automatically; they accumulate for the lifetime of the DO. Bounded -@@ -893,12 +1496,34 @@ var DurableObjectEventStore = class DurableObjectEventStore { +@@ -893,12 +1497,34 @@ var DurableObjectEventStore = class DurableObjectEventStore { } async storeEvent(streamId, message) { if (streamId.includes(":")) throw new Error(`DurableObjectEventStore: streamId must not contain ':' (got ${JSON.stringify(streamId)})`); @@ -1111,7 +1097,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 return eventId; } async getStreamIdForEventId(eventId) { -@@ -915,9 +1540,59 @@ var DurableObjectEventStore = class DurableObjectEventStore { +@@ -915,9 +1541,59 @@ var DurableObjectEventStore = class DurableObjectEventStore { start: startKey, limit: DurableObjectEventStore.REPLAY_LIMIT }); @@ -1172,7 +1158,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 /** * Drop the event log for a single stream. Called by the transport * immediately after a POST's final response has been written to the -@@ -973,6 +1648,13 @@ DurableObjectEventStore.EVENT_KEY_PREFIX = "__mcp_event__:"; +@@ -973,6 +1649,13 @@ DurableObjectEventStore.EVENT_KEY_PREFIX = "__mcp_event__:"; DurableObjectEventStore.SEQ_PAD = 16; DurableObjectEventStore.DELETE_CHUNK = 128; DurableObjectEventStore.REPLAY_LIMIT = 1e3; @@ -1186,7 +1172,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 //#endregion //#region src/mcp/client-transports.ts /** -@@ -1355,6 +2037,26 @@ function experimental_createMcpHandler(server, options = {}) { +@@ -1355,6 +2038,26 @@ function experimental_createMcpHandler(server, options = {}) { } //#endregion //#region src/mcp/index.ts @@ -1213,7 +1199,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 var McpAgent = class McpAgent extends Agent { constructor(..._args) { super(..._args); -@@ -1369,18 +2071,121 @@ var McpAgent = class McpAgent extends Agent { +@@ -1369,18 +2072,121 @@ var McpAgent = class McpAgent extends Agent { async getInitializeRequest() { return this.ctx.storage.get("initializeRequest"); } @@ -1338,7 +1324,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 /** * Reverse lookup: find which POST stream a given `requestId` belongs * to, and return the stream's full `requestIds` list in the same -@@ -1407,10 +2212,14 @@ var McpAgent = class McpAgent extends Agent { +@@ -1407,10 +2213,14 @@ var McpAgent = class McpAgent extends Agent { limit: STREAM_REQS_SCAN_LIMIT }); if (rows.size === STREAM_REQS_SCAN_LIMIT) console.warn(`McpAgent: getStreamForRequestId hit the ${STREAM_REQS_SCAN_LIMIT}-key scan cap; stale __mcp_stream_reqs__ entries may be accumulating from abandoned POSTs`); @@ -1357,7 +1343,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 } /** Read the transport type for this agent. * This relies on the naming scheme being `sse:${sessionId}`, -@@ -1498,6 +2307,12 @@ var McpAgent = class McpAgent extends Agent { +@@ -1498,6 +2308,12 @@ var McpAgent = class McpAgent extends Agent { } /** Sets up the MCP transport and server every time the Agent is started.*/ async onStart(props) { @@ -1370,7 +1356,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 if (props) await this.updateProps(props); else this.props = await this.ctx.storage.get("props"); await this.init(); -@@ -1516,23 +2331,36 @@ var McpAgent = class McpAgent extends Agent { +@@ -1516,23 +2332,36 @@ var McpAgent = class McpAgent extends Agent { return; } break; @@ -1422,7 +1408,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..231220c26b995a48e40dba993773c287 } } } -@@ -1697,7 +2525,9 @@ var McpAgent = class McpAgent extends Agent { +@@ -1697,7 +2526,9 @@ var McpAgent = class McpAgent extends Agent { } }; McpAgent.STREAM_REQS_KEY_PREFIX = "__mcp_stream_reqs__:";