Skip to content
Merged
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
76 changes: 31 additions & 45 deletions patches/agents@0.17.3.patch
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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/
Expand Down Expand Up @@ -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 });
Expand Down Expand Up @@ -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);
Expand All @@ -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();
Expand Down Expand Up @@ -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);
});
Expand All @@ -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 });
Expand All @@ -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 });
Expand Down Expand Up @@ -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();
Expand All @@ -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 });
Expand Down Expand Up @@ -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();
Expand All @@ -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 };
Expand Down Expand Up @@ -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.
Expand All @@ -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) {
Expand All @@ -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);
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
Expand All @@ -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) {
Expand All @@ -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);
}
Expand All @@ -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
Expand All @@ -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
*
Expand All @@ -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)})`);
Expand Down Expand Up @@ -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
});
Expand Down Expand Up @@ -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;
Expand All @@ -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
Expand All @@ -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");
}
Expand Down Expand Up @@ -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`);
Expand All @@ -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) {
Expand All @@ -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;
Expand Down Expand Up @@ -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__:";
Expand Down
Loading