diff --git a/packages/ai/src/chatTransport.test.ts b/packages/ai/src/chatTransport.test.ts index 89f67451d..1896ec1f1 100644 --- a/packages/ai/src/chatTransport.test.ts +++ b/packages/ai/src/chatTransport.test.ts @@ -476,6 +476,33 @@ describe("TriggerChatTransport", function () { expect(stream).toBeNull(); }); + it("removes inactive run entries during reconnect attempts", async function () { + const runStore = new TrackedRunStore(); + runStore.set({ + chatId: "chat-inactive", + runId: "run_inactive", + publicAccessToken: "pk_inactive", + streamKey: "chat-stream", + lastEventId: "10-0", + isActive: false, + }); + + const transport = new TriggerChatTransport({ + task: "chat-task", + stream: "chat-stream", + accessToken: "pk_trigger", + runStore, + }); + + const stream = await transport.reconnectToStream({ + chatId: "chat-inactive", + }); + + expect(stream).toBeNull(); + expect(runStore.deleteCalls).toContain("chat-inactive"); + expect(runStore.get("chat-inactive")).toBeUndefined(); + }); + it("supports custom payload mapping and trigger options resolver", async function () { let receivedTriggerBody: Record | undefined; let receivedResolverChatId: string | undefined; diff --git a/packages/ai/src/chatTransport.ts b/packages/ai/src/chatTransport.ts index d80b2fdeb..c65f5b532 100644 --- a/packages/ai/src/chatTransport.ts +++ b/packages/ai/src/chatTransport.ts @@ -195,7 +195,12 @@ export class TriggerChatTransport< ): Promise | null> { const runState = await this.runStore.get(options.chatId); - if (!runState || !runState.isActive) { + if (!runState) { + return null; + } + + if (!runState.isActive) { + await this.runStore.delete(options.chatId); return null; }