diff --git a/packages/trigger-sdk/src/v3/ai.ts b/packages/trigger-sdk/src/v3/ai.ts index e10306681..51dd4ec3b 100644 --- a/packages/trigger-sdk/src/v3/ai.ts +++ b/packages/trigger-sdk/src/v3/ai.ts @@ -5062,7 +5062,7 @@ function makeChannelStreamEditor( let latest = ""; let lastSent = ""; let timer: ReturnType | undefined; - let inFlight = false; + let inFlightPromise: Promise | undefined; let stopped = false; const arm = () => { @@ -5074,7 +5074,7 @@ function makeChannelStreamEditor( }; const edit = async () => { - if (stopped || inFlight) return; + if (stopped || inFlightPromise) return; const text = latest; if (text === lastSent) return; const message = outbound({ @@ -5084,22 +5084,24 @@ function makeChannelStreamEditor( stopped: false, }); if (!message) return; - inFlight = true; - try { - await connector.send!(message, { - event: channelEvent.event as TEvent, - deliveryId: channelEvent.deliveryId, - previousRef: ackRef, - mode: "stream", - final: false, - }); - lastSent = text; - } catch (error) { - logger.warn("chat.agent: channel stream edit failed", { error }); - } finally { - inFlight = false; - if (!stopped && latest !== lastSent) arm(); - } + inFlightPromise = (async () => { + try { + await connector.send!(message, { + event: channelEvent.event as TEvent, + deliveryId: channelEvent.deliveryId, + previousRef: ackRef, + mode: "stream", + final: false, + }); + lastSent = text; + } catch (error) { + logger.warn("chat.agent: channel stream edit failed", { error }); + } finally { + inFlightPromise = undefined; + if (!stopped && latest !== lastSent) arm(); + } + })(); + await inFlightPromise; }; return { @@ -5111,17 +5113,22 @@ function makeChannelStreamEditor( arm(); } }, - stop() { + async stop() { stopped = true; if (timer) { clearTimeout(timer); timer = undefined; } + if (inFlightPromise) { + try { + await inFlightPromise; + } catch {} + } }, }; } -type ChannelStreamEditor = { observe(chunk: unknown): void; stop(): void }; +type ChannelStreamEditor = { observe(chunk: unknown): void; stop(): void | Promise }; /** * Wrap the reply stream so it debounce-edits the channel message as it streams. @@ -5132,18 +5139,18 @@ type ChannelStreamEditor = { observe(chunk: unknown): void; stop(): void }; function makeChannelStreamTap(editor: ChannelStreamEditor): TransformStream { const transformer: { transform(chunk: unknown, controller: TransformStreamDefaultController): void; - flush(): void; - cancel?: (reason?: unknown) => void; + flush(): Promise; + cancel?: (reason?: unknown) => Promise; } = { transform(chunk, controller) { editor.observe(chunk); controller.enqueue(chunk); }, - flush() { - editor.stop(); + async flush() { + await editor.stop(); }, - cancel() { - editor.stop(); + async cancel() { + await editor.stop(); }, }; return new TransformStream(transformer); @@ -7938,7 +7945,7 @@ function chatAgent< } } finally { msgSub.off(); - channelStreamEditor?.stop(); + await channelStreamEditor?.stop(); } // Wait for onFinish to fire — on abort this may resolve slightly diff --git a/packages/trigger-sdk/test/chatChannels.test.ts b/packages/trigger-sdk/test/chatChannels.test.ts index f8c53b4e3..79398689a 100644 --- a/packages/trigger-sdk/test/chatChannels.test.ts +++ b/packages/trigger-sdk/test/chatChannels.test.ts @@ -350,7 +350,12 @@ describe("makeChannelStreamEditor", () => { describe("makeChannelStreamTap", () => { it("stops the editor when the stream completes (flush)", async () => { let stops = 0; - const tap = __makeChannelStreamTapForTests({ observe: () => {}, stop: () => (stops += 1) }); + const tap = __makeChannelStreamTapForTests({ + observe: () => {}, + stop: () => { + stops += 1; + }, + }); const source = new ReadableStream({ start(controller) { controller.enqueue({ type: "text-delta", delta: "x" }); @@ -367,7 +372,12 @@ describe("makeChannelStreamTap", () => { it("stops the editor when the stream is cancelled mid-flight (abort)", async () => { let stops = 0; - const tap = __makeChannelStreamTapForTests({ observe: () => {}, stop: () => (stops += 1) }); + const tap = __makeChannelStreamTapForTests({ + observe: () => {}, + stop: () => { + stops += 1; + }, + }); const source = new ReadableStream({ start(controller) { controller.enqueue({ type: "text-delta", delta: "x" });