diff --git a/apps/host-cloudflare/package.json b/apps/host-cloudflare/package.json index b693834278..94ad2bca33 100644 --- a/apps/host-cloudflare/package.json +++ b/apps/host-cloudflare/package.json @@ -36,6 +36,7 @@ "@executor-js/sdk": "workspace:*", "@jitl/quickjs-wasmfile-release-sync": "catalog:", "@modelcontextprotocol/sdk": "^1.29.0", + "@sentry/cloudflare": "^10.48.0", "@tanstack/react-router": "catalog:", "drizzle-orm": "catalog:", "effect": "catalog:", diff --git a/apps/host-cloudflare/src/config.ts b/apps/host-cloudflare/src/config.ts index 68b8296399..376e9954c2 100644 --- a/apps/host-cloudflare/src/config.ts +++ b/apps/host-cloudflare/src/config.ts @@ -78,6 +78,8 @@ export interface CloudflareEnv { */ readonly AI_GATEWAY_TOKEN?: SecretsStoreBinding; readonly VITE_PUBLIC_SITE_URL?: string; + /** Sentry DSN (a `wrangler secret`). Unset leaves Sentry disabled. */ + readonly SENTRY_DSN?: string; /** * Dev/single-user escape hatch: when "true", skip Cloudflare Access entirely * and treat every request as a fixed admin. For local `wrangler dev` and @@ -124,9 +126,7 @@ type CloudflareAccessEnv = Pick< // Both ids are required: a gateway URL missing either one resolves to a 404 that // would look like "Jev found nothing" rather than "Jev was never configured". -const resolveJevGateway = ( - env: CloudflareConfigEnv, -): CloudflareConfig["jevGateway"] => { +const resolveJevGateway = (env: CloudflareConfigEnv): CloudflareConfig["jevGateway"] => { const accountId = env.CLOUDFLARE_ACCOUNT_ID?.trim() ?? ""; const gatewayId = env.AI_GATEWAY_ID?.trim() ?? ""; if (accountId.length === 0 || gatewayId.length === 0) { diff --git a/apps/host-cloudflare/src/observability.ts b/apps/host-cloudflare/src/observability.ts index 249eb52061..5edac2064d 100644 --- a/apps/host-cloudflare/src/observability.ts +++ b/apps/host-cloudflare/src/observability.ts @@ -1,7 +1,40 @@ -// Cloudflare host `ErrorCapture` — the shared console implementation with a -// `cloudflare-` trace-id prefix. Worker stdout is routed to Logpush/the -// dashboard, so the squashed cause is grep-able by the opaque 500 traceId. +// Cloudflare host `ErrorCapture`: the shared console implementation (a +// `cloudflare-` trace id, grep-able in Workers logs) plus a Sentry report +// tagged with that trace id. Sentry stays a no-op while SENTRY_DSN is unset. +import * as Sentry from "@sentry/cloudflare"; +import { Cause, Effect, Layer } from "effect"; +import { ErrorCapture } from "@executor-js/api"; import { consoleErrorCapture } from "@executor-js/api/server"; -export const ErrorCaptureLive = consoleErrorCapture("cloudflare"); +export const ErrorCaptureLive: Layer.Layer = Layer.effect( + ErrorCapture, + Effect.gen(function* () { + const consoleCapture = yield* ErrorCapture; + return ErrorCapture.of({ + captureException: (cause) => + consoleCapture.captureException(cause).pipe( + Effect.tap((traceId) => + Effect.sync(() => { + const [error] = Cause.prettyErrors(cause); + Sentry.captureException(error ?? Cause.squash(cause), { tags: { traceId } }); + }), + ), + ), + }); + }), +).pipe(Layer.provide(consoleErrorCapture("cloudflare"))); + +/** + * A call the client was waiting on that will never get a real answer: its + * session Durable Object died (memory or CPU limit, deploy) or nothing came + * back before the deadline. The dying isolate cannot report itself, so the + * front worker, which answers the client for it, reports it instead. + */ +export const reportLostCall = (detail: Readonly>): void => { + Sentry.captureMessage(`MCP call lost: ${String(detail.reason ?? "unknown")}`, { + level: "error", + tags: { reason: String(detail.reason ?? "unknown") }, + extra: detail, + }); +}; diff --git a/apps/host-cloudflare/src/sentry.ts b/apps/host-cloudflare/src/sentry.ts new file mode 100644 index 0000000000..8350ab8aaf --- /dev/null +++ b/apps/host-cloudflare/src/sentry.ts @@ -0,0 +1,7 @@ +import type { CloudflareEnv } from "./config"; + +export const sentryOptions = (env: CloudflareEnv) => ({ + dsn: env.SENTRY_DSN, + tracesSampleRate: 0, + sendDefaultPii: false, +}); diff --git a/apps/host-cloudflare/src/worker.ts b/apps/host-cloudflare/src/worker.ts index b9964fac4a..7c6d6d8802 100644 --- a/apps/host-cloudflare/src/worker.ts +++ b/apps/host-cloudflare/src/worker.ts @@ -1,14 +1,27 @@ +import * as Sentry from "@sentry/cloudflare"; + import { makeCloudflareApp } from "./app"; import { cloudflareAccessConfigErrorMessage, missingCloudflareAccessVars, type CloudflareEnv, } from "./config"; +import { McpSessionDO as McpSessionDOBase } from "./mcp"; import { mcpResourceFromPath } from "./mcp/resource"; +import { reportLostCall } from "./observability"; +import { sentryOptions } from "./sentry"; + +(globalThis as { __executorReportLostCall?: typeof reportLostCall }).__executorReportLostCall = + reportLostCall; // The MCP Durable Object classes, bound in wrangler.jsonc. They must be exported -// at the Worker entry module scope for the runtime to find them. -export { McpExecutionOwnerDirectoryDO, McpSessionDO } from "./mcp"; +// at the Worker entry module scope for the runtime to find them. Wrapping the +// session DO initialises Sentry inside its isolate. +export { McpExecutionOwnerDirectoryDO } from "./mcp"; +export const McpSessionDO = Sentry.instrumentDurableObjectWithSentry( + sentryOptions, + McpSessionDOBase, +); // --------------------------------------------------------------------------- // The Worker fetch entry. Most requests go to `ExecutorApp.make`'s Effect web @@ -41,7 +54,7 @@ const accessConfigErrorResponse = (missingVars: readonly string[]): Response => }, }); -export default { +export default Sentry.withSentry(sentryOptions, { fetch: async (request: Request, env: CloudflareEnv, ctx: ExecutionContext): Promise => { const missingAccessVars = missingCloudflareAccessVars(env); if (missingAccessVars.length > 0) { @@ -55,4 +68,4 @@ export default { } return serve.app(request); }, -}; +} satisfies ExportedHandler); diff --git a/bun.lock b/bun.lock index 78ee6bcded..dc5a396dad 100644 --- a/bun.lock +++ b/bun.lock @@ -192,6 +192,7 @@ "@executor-js/sdk": "workspace:*", "@jitl/quickjs-wasmfile-release-sync": "catalog:", "@modelcontextprotocol/sdk": "^1.29.0", + "@sentry/cloudflare": "^10.48.0", "@tanstack/react-router": "catalog:", "drizzle-orm": "catalog:", "effect": "catalog:", diff --git a/packages/hosts/cloudflare/src/mcp/agents-post-stream-loss.test.ts b/packages/hosts/cloudflare/src/mcp/agents-post-stream-loss.test.ts index 1681709ca6..503d300d49 100644 --- a/packages/hosts/cloudflare/src/mcp/agents-post-stream-loss.test.ts +++ b/packages/hosts/cloudflare/src/mcp/agents-post-stream-loss.test.ts @@ -208,6 +208,31 @@ describe("POST bridge: lost-execution visibility", () => { } }); + it("reports each lost call to the host's lost-call hook, and nothing on the happy path", async () => { + const reported: unknown[] = []; + const holder = globalThis as { __executorReportLostCall?: (detail: unknown) => void }; + holder.__executorReportLostCall = (detail) => reported.push(detail); + try { + const delivered = await postToBridge(toolCall(1)); + const deliveredBody = drainResponse(delivered.response); + emitResponse(delivered.ws, { id: 1, jsonrpc: "2.0", result: { ok: true } }); + await flushMicrotasks(); + emitAbnormalClose(delivered.ws); + await deliveredBody; + expect(reported).toEqual([]); + + const lost = await postToBridge(toolCall(2)); + const lostBody = drainResponse(lost.response); + emitAbnormalClose(lost.ws, 1006, "WebSocket disconnected without sending Close frame."); + await lostBody; + expect(reported).toEqual([ + expect.objectContaining({ outstandingCount: 1, reason: "session_reset" }), + ]); + } finally { + delete holder.__executorReportLostCall; + } + }); + it("writes nothing extra when the response was delivered before the close", async () => { const { response, ws } = await postToBridge(toolCall(1)); const drained = drainResponse(response); diff --git a/patches/agents@0.17.3.patch b/patches/agents@0.17.3.patch index f785f07484..c195a9cd83 100644 --- a/patches/agents@0.17.3.patch +++ b/patches/agents@0.17.3.patch @@ -80,7 +80,7 @@ index c8fad448e8797b89690a99d93490d1363851b225..d80f66f1532c29dda0a0477dc1ec9ef5 McpAgent, type McpAuthContext, diff --git a/dist/mcp/index.js b/dist/mcp/index.js -index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed0768462f640 100644 +index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..30dcc9bbb148db7bfe9f649de332289c89a06c12 100644 --- a/dist/mcp/index.js +++ b/dist/mcp/index.js @@ -28,13 +28,60 @@ import { WebStandardStreamableHTTPServerTransport } from "@modelcontextprotocol/ @@ -181,7 +181,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 request.headers.forEach((value, key) => { existingHeaders[key] = value; }); -@@ -206,47 +273,160 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -206,47 +273,166 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { jsonrpc: "2.0" }); return new Response(body, { status: 500 }); @@ -274,6 +274,12 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 + sessionId, + ...detail + })); ++ globalThis.__executorReportLostCall?.({ ++ outstandingCount: __outstanding.size, ++ reason, ++ sessionId, ++ ...detail ++ }); + for (const id of __outstanding.values()) __forwardSse(encoder.encode(`event: message\ndata: ${JSON.stringify(sessionResetErrorResponse(id, reason))}\n\n`)); + __outstanding.clear(); + } @@ -369,7 +375,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 }); return new Response(readable, { headers: { -@@ -279,10 +459,16 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -279,10 +465,16 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { id: null, jsonrpc: "2.0" }), { status: 400 }); @@ -390,7 +396,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 props: ctx.props, jurisdiction: options.jurisdiction }); -@@ -306,27 +492,116 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { +@@ -306,27 +498,116 @@ const createStreamingHttpHandler = (basePath, namespace, options = {}) => { if (!ws) { await writer.close(); return new Response("Failed to establish WS to DO", { status: 500 }); @@ -522,7 +528,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 return new Response(readable, { headers: { "Cache-Control": "no-cache", -@@ -389,10 +664,16 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { +@@ -389,10 +670,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(); @@ -543,7 +549,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 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 +695,94 @@ const createLegacySseHandler = (basePath, namespace, options = {}) => { +@@ -414,35 +701,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 }); @@ -656,7 +662,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 console.error("Error closing SSE connection:", error); } } -@@ -586,6 +926,8 @@ var StreamableHTTPServerTransport = class { +@@ -586,6 +932,8 @@ var StreamableHTTPServerTransport = class { constructor(options) { this._started = false; this._streamResponseIds = /* @__PURE__ */ new Map(); @@ -665,7 +671,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 const { agent } = getCurrentAgent(); if (!agent) throw new Error("McpAgent was not found in Transport constructor"); this._agent = agent; -@@ -627,23 +969,146 @@ var StreamableHTTPServerTransport = class { +@@ -627,23 +975,146 @@ var StreamableHTTPServerTransport = class { const resumedStreamId = await this._eventStore.getStreamIdForEventId?.(lastEventId); if (resumedStreamId) { const resumeState = { streamId: resumedStreamId }; @@ -818,7 +824,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 } /** * Close any connection (other than `selfId`) currently bound to -@@ -651,6 +1116,9 @@ var StreamableHTTPServerTransport = class { +@@ -651,6 +1122,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. @@ -828,7 +834,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 */ supersedePriorStreamConnections(agent, selfId, streamId) { for (const other of agent.getConnections()) { -@@ -664,12 +1132,14 @@ var StreamableHTTPServerTransport = class { +@@ -664,12 +1138,14 @@ var StreamableHTTPServerTransport = class { * Only used when resumability is enabled */ async replayEvents(lastEventId) { @@ -844,7 +850,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 this.writeSSEEvent(connection, message, eventId); } catch (error) { this.onerror?.(error); -@@ -678,6 +1148,45 @@ var StreamableHTTPServerTransport = class { +@@ -678,6 +1154,45 @@ var StreamableHTTPServerTransport = class { } catch (error) { this.onerror?.(error); } @@ -890,7 +896,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 } /** * Writes an event to the SSE stream with proper formatting -@@ -689,10 +1198,69 @@ var StreamableHTTPServerTransport = class { +@@ -689,10 +1204,69 @@ var StreamableHTTPServerTransport = class { return connection.send(JSON.stringify({ type: "cf_mcp_agent_event", event: eventData, @@ -960,7 +966,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 * Handles POST requests containing JSON-RPC messages */ async handlePostRequest(req, parsedBody) { -@@ -733,6 +1301,22 @@ var StreamableHTTPServerTransport = class { +@@ -733,6 +1307,22 @@ var StreamableHTTPServerTransport = class { }; connection.setState(postState); if (this._eventStore) await agent.setStreamRequestIds(streamId, requestIds); @@ -983,7 +989,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 for (const message of messages) { if (this.messageInterceptor) { if (await this.messageInterceptor(message, { -@@ -760,7 +1344,22 @@ var StreamableHTTPServerTransport = class { +@@ -760,7 +1350,22 @@ var StreamableHTTPServerTransport = class { * when the originating WS has dropped. */ async sendOnStream(agent, streamId, relatedIds, liveConnection, message, requestId) { @@ -1007,7 +1013,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 let shouldClose = false; if (isJSONRPCResultResponse(message) || isJSONRPCErrorResponse(message)) { let responseIds = this._streamResponseIds.get(streamId); -@@ -777,9 +1376,11 @@ var StreamableHTTPServerTransport = class { +@@ -777,9 +1382,11 @@ var StreamableHTTPServerTransport = class { } catch (error) { this.onerror?.(error); } @@ -1021,7 +1027,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 } } async send(message, options) { -@@ -798,14 +1399,19 @@ var StreamableHTTPServerTransport = class { +@@ -798,14 +1405,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 @@ -1045,7 +1051,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 if (standalone) this.writeSSEEvent(standalone, message, eventId); } /** -@@ -861,12 +1467,10 @@ var StreamableHTTPServerTransport = class { +@@ -861,12 +1473,10 @@ var StreamableHTTPServerTransport = class { * * ## Lifecycle * @@ -1062,7 +1068,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 * * Standalone GET stream events (`_GET_stream`) are *not* cleared * automatically; they accumulate for the lifetime of the DO. Bounded -@@ -893,12 +1497,34 @@ var DurableObjectEventStore = class DurableObjectEventStore { +@@ -893,12 +1503,34 @@ var DurableObjectEventStore = class DurableObjectEventStore { } async storeEvent(streamId, message) { if (streamId.includes(":")) throw new Error(`DurableObjectEventStore: streamId must not contain ':' (got ${JSON.stringify(streamId)})`); @@ -1097,7 +1103,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 return eventId; } async getStreamIdForEventId(eventId) { -@@ -915,9 +1541,59 @@ var DurableObjectEventStore = class DurableObjectEventStore { +@@ -915,9 +1547,59 @@ var DurableObjectEventStore = class DurableObjectEventStore { start: startKey, limit: DurableObjectEventStore.REPLAY_LIMIT }); @@ -1158,7 +1164,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 /** * 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 +1649,13 @@ DurableObjectEventStore.EVENT_KEY_PREFIX = "__mcp_event__:"; +@@ -973,6 +1655,13 @@ DurableObjectEventStore.EVENT_KEY_PREFIX = "__mcp_event__:"; DurableObjectEventStore.SEQ_PAD = 16; DurableObjectEventStore.DELETE_CHUNK = 128; DurableObjectEventStore.REPLAY_LIMIT = 1e3; @@ -1172,7 +1178,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 //#endregion //#region src/mcp/client-transports.ts /** -@@ -1355,6 +2038,26 @@ function experimental_createMcpHandler(server, options = {}) { +@@ -1355,6 +2044,26 @@ function experimental_createMcpHandler(server, options = {}) { } //#endregion //#region src/mcp/index.ts @@ -1199,7 +1205,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 var McpAgent = class McpAgent extends Agent { constructor(..._args) { super(..._args); -@@ -1369,18 +2072,121 @@ var McpAgent = class McpAgent extends Agent { +@@ -1369,18 +2078,121 @@ var McpAgent = class McpAgent extends Agent { async getInitializeRequest() { return this.ctx.storage.get("initializeRequest"); } @@ -1324,7 +1330,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 /** * Reverse lookup: find which POST stream a given `requestId` belongs * to, and return the stream's full `requestIds` list in the same -@@ -1407,10 +2213,14 @@ var McpAgent = class McpAgent extends Agent { +@@ -1407,10 +2219,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`); @@ -1343,7 +1349,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 } /** Read the transport type for this agent. * This relies on the naming scheme being `sse:${sessionId}`, -@@ -1498,6 +2308,12 @@ var McpAgent = class McpAgent extends Agent { +@@ -1498,6 +2314,12 @@ var McpAgent = class McpAgent extends Agent { } /** Sets up the MCP transport and server every time the Agent is started.*/ async onStart(props) { @@ -1356,7 +1362,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 if (props) await this.updateProps(props); else this.props = await this.ctx.storage.get("props"); await this.init(); -@@ -1516,23 +2332,36 @@ var McpAgent = class McpAgent extends Agent { +@@ -1516,23 +2338,36 @@ var McpAgent = class McpAgent extends Agent { return; } break; @@ -1408,7 +1414,7 @@ index 1edcf0c8c9e67aa211ae515e7672cdf79912101e..c661da89ea510827af5bf6a374eed076 } } } -@@ -1697,7 +2526,9 @@ var McpAgent = class McpAgent extends Agent { +@@ -1697,7 +2532,9 @@ var McpAgent = class McpAgent extends Agent { } }; McpAgent.STREAM_REQS_KEY_PREFIX = "__mcp_stream_reqs__:";