diff --git a/src/channels/draft-stream-loop.test.ts b/src/channels/draft-stream-loop.test.ts index 2e39efb0c431..2a4f83bb6de4 100644 --- a/src/channels/draft-stream-loop.test.ts +++ b/src/channels/draft-stream-loop.test.ts @@ -260,6 +260,75 @@ describe("createDraftStreamLoop", () => { expect(sendOrEditStreamMessage).toHaveBeenNthCalledWith(2, "hello"); }); + it("preserves concurrent update text when sendOrEditStreamMessage returns false", async () => { + let deliver!: (() => void) | undefined; + const deliverPromise = new Promise((resolve) => { + deliver = resolve; + }); + const capturedArgs: string[] = []; + + const sendOrEditStreamMessage = vi + .fn<(text: string) => Promise>() + .mockImplementationOnce(async (text: string) => { + capturedArgs.push(text); + await deliverPromise; + return false; + }) + .mockImplementationOnce(async (text: string) => { + capturedArgs.push(text); + return true; + }); + + const loop = createDraftStreamLoop({ + throttleMs: 0, + isStopped: () => false, + sendOrEditStreamMessage, + }); + + loop.update("initial"); + await flushMicrotasks(); + + loop.update("concurrent"); + deliver!(); + await loop.flush(); + + expect(capturedArgs[1]).toBe("concurrent"); + }); + + it("preserves generic pending updates when consecutive sends return false", async () => { + type Update = { text: string; blocks: string[] }; + const initial = { text: "initial", blocks: ["initial-blocks"] }; + const latest = { text: "latest", blocks: ["latest-blocks"] }; + let releaseFirst: (() => void) | undefined; + const firstSend = new Promise((resolve) => { + releaseFirst = resolve; + }); + const sendOrEditStreamMessage = vi + .fn<(update: Update) => Promise>() + .mockImplementationOnce(async () => { + await firstSend; + return false; + }) + .mockResolvedValueOnce(false) + .mockResolvedValueOnce(true); + const loop = createDraftStreamLoop({ + throttleMs: 0, + isStopped: () => false, + emptyValue: { text: "", blocks: [] }, + isEmpty: (update) => !update.text, + sendOrEditStreamMessage, + }); + + loop.update(initial); + await flushMicrotasks(); + loop.update(latest); + releaseFirst?.(); + await loop.flush(); + await loop.flush(); + + expect(sendOrEditStreamMessage.mock.calls).toEqual([[initial], [latest], [latest]]); + }); + it("keeps generic payload fields atomic while newer updates queue", async () => { type Update = { text: string; blocks: string[] }; let releaseFirst: (() => void) | undefined; diff --git a/src/channels/draft-stream-loop.ts b/src/channels/draft-stream-loop.ts index debeb7812ee2..adb2cd31fcc1 100644 --- a/src/channels/draft-stream-loop.ts +++ b/src/channels/draft-stream-loop.ts @@ -84,7 +84,9 @@ export function createDraftStreamLoop(params: { throw err; } if (sent === false) { - pendingValue = value; + if (!hasPendingValue(pendingValue)) { + pendingValue = value; + } return; } lastSentAt = Date.now();