Ignore onTriggeredRun callback failures during streaming
Co-authored-by: Eric Allam <eric@trigger.dev>
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user