From bc1a10c28e4adbfe32cd65e86741da567bd22c21 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 15 Feb 2026 00:06:51 +0000 Subject: [PATCH] Ignore onTriggeredRun callback failures during streaming Co-authored-by: Eric Allam --- packages/ai/README.md | 2 +- packages/ai/src/chatTransport.test.ts | 62 +++++++++++++++++++++++++++ packages/ai/src/chatTransport.ts | 6 ++- 3 files changed, 68 insertions(+), 2 deletions(-) diff --git a/packages/ai/README.md b/packages/ai/README.md index 3f3fc1833..d642510a5 100644 --- a/packages/ai/README.md +++ b/packages/ai/README.md @@ -117,7 +117,7 @@ class MemoryStore implements TriggerChatRunStore { ``` `onTriggeredRun` can also be async, which is useful for persisting run IDs before -the chat stream is consumed. +the chat stream is consumed. Callback failures are ignored so chat streaming can continue. ## `ai.tool(...)` example diff --git a/packages/ai/src/chatTransport.test.ts b/packages/ai/src/chatTransport.test.ts index 391cfd1f0..ede31f1cd 100644 --- a/packages/ai/src/chatTransport.test.ts +++ b/packages/ai/src/chatTransport.test.ts @@ -688,6 +688,68 @@ describe("TriggerChatTransport", function () { expect(callbackCompleted).toBe(true); }); + it("continues streaming when onTriggeredRun callback throws", async function () { + let callbackCalled = false; + + const server = await startServer(function (req, res) { + if (req.method === "POST" && req.url === "/api/v1/tasks/chat-task/trigger") { + res.writeHead(200, { + "content-type": "application/json", + "x-trigger-jwt": "pk_run_callback_error", + }); + res.end(JSON.stringify({ id: "run_callback_error" })); + return; + } + + if ( + req.method === "GET" && + req.url === "/realtime/v1/streams/run_callback_error/chat-stream" + ) { + res.writeHead(200, { + "content-type": "text/event-stream", + }); + writeSSE( + res, + "1-0", + JSON.stringify({ type: "text-start", id: "callback_error_1" }) + ); + writeSSE( + res, + "2-0", + JSON.stringify({ type: "text-end", id: "callback_error_1" }) + ); + res.end(); + return; + } + + res.writeHead(404); + res.end(); + }); + + const transport = new TriggerChatTransport({ + task: "chat-task", + stream: "chat-stream", + accessToken: "pk_trigger", + baseURL: server.url, + onTriggeredRun: async function onTriggeredRun() { + callbackCalled = true; + throw new Error("callback failed"); + }, + }); + + const stream = await transport.sendMessages({ + trigger: "submit-message", + chatId: "chat-callback-error", + messageId: undefined, + messages: [], + abortSignal: undefined, + }); + + const chunks = await readChunks(stream); + expect(callbackCalled).toBe(true); + expect(chunks).toHaveLength(2); + }); + it("cleans run store state when stream completes", async function () { const trackedRunStore = new TrackedRunStore(); diff --git a/packages/ai/src/chatTransport.ts b/packages/ai/src/chatTransport.ts index 7186c9d39..d520e39ac 100644 --- a/packages/ai/src/chatTransport.ts +++ b/packages/ai/src/chatTransport.ts @@ -180,7 +180,11 @@ export class TriggerChatTransport< await this.runStore.set(runState); if (this.onTriggeredRun) { - await this.onTriggeredRun(runState); + try { + await this.onTriggeredRun(runState); + } catch { + // Ignore callback errors so chat streaming can continue. + } } const stream = await this.fetchRunStream(runState, options.abortSignal);