From 8ab26e59e37563a64f60d35a8dbcff0c0d8f97e6 Mon Sep 17 00:00:00 2001 From: Cyril Date: Mon, 7 Sep 2026 14:18:08 +0900 Subject: [PATCH 1/3] fix(client): reset message queue processing flag when sendFn throws processMessageQueue set isProcessingQueueRef.current back to false only after the while loop completed normally. When sendFn rejected (e.g. a failed file upload), the flag stayed true for the component's lifetime, so every subsequent sendMessage call was silently swallowed by the early-return guard: no error, no queued message, nothing rendered. Wrap the loop body in try/finally so the flag always resets, whether the loop finishes or throws. The error still propagates to the sendMessage caller unchanged. The message that was sending when the error occurred was already shifted off pendingMessagesRef before the send attempt, so it is not retried; any messages queued after it stay queued and are processed on the next sendMessage call, same as today. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01X61A3hcJvq7HfxcrYKEfNZ --- .../client/src/hooks/useMessageQueue.test.ts | 38 ++++++++ packages/client/src/hooks/useMessageQueue.ts | 96 ++++++++++--------- 2 files changed, 87 insertions(+), 47 deletions(-) create mode 100644 packages/client/src/hooks/useMessageQueue.test.ts diff --git a/packages/client/src/hooks/useMessageQueue.test.ts b/packages/client/src/hooks/useMessageQueue.test.ts new file mode 100644 index 00000000..218d2108 --- /dev/null +++ b/packages/client/src/hooks/useMessageQueue.test.ts @@ -0,0 +1,38 @@ +import { describe, test, expect, mock } from 'bun:test'; +import { renderHook, act } from '@testing-library/react'; +import { useMessageQueue } from './useMessageQueue'; + +function createOptions(sendFn: (message: string) => Promise) { + return { + sendFn, + createNewChat: mock(async () => 'chat-id'), + connected: true, + loading: false, + hasPendingApproval: false, + }; +} + +describe('useMessageQueue', () => { + test('processes a later message after an earlier send rejects', async () => { + const sendFn = mock((message: string) => { + if (message === 'fails') { + return Promise.reject(new Error('upload failed')); + } + return Promise.resolve(); + }); + + const { result } = renderHook(() => useMessageQueue(createOptions(sendFn))); + + await act(async () => { + await expect(result.current.sendMessage('fails')).rejects.toThrow('upload failed'); + }); + + await act(async () => { + await result.current.sendMessage('second'); + }); + + expect(sendFn).toHaveBeenCalledTimes(2); + expect(sendFn).toHaveBeenNthCalledWith(1, 'fails', undefined, undefined); + expect(sendFn).toHaveBeenNthCalledWith(2, 'second', undefined, undefined); + }); +}); diff --git a/packages/client/src/hooks/useMessageQueue.ts b/packages/client/src/hooks/useMessageQueue.ts index 087ab19a..f468c6f5 100644 --- a/packages/client/src/hooks/useMessageQueue.ts +++ b/packages/client/src/hooks/useMessageQueue.ts @@ -97,56 +97,58 @@ export function useMessageQueue({ isProcessingQueueRef.current = true; - while (pendingMessagesRef.current.length > 0) { - const { message, options } = pendingMessagesRef.current.shift()!; - const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {}; - - if (newChat) { - await createNewChatRef.current({ metadata }); - } - - // Convert File[] to FileAttachment[] - const fileAttachments: FileAttachment[] = await Promise.all( - attachments.map(async (file) => { - let preview: string | undefined; - if (file.type.startsWith('image/')) { - preview = await new Promise((resolve) => { - const reader = new FileReader(); - reader.onload = () => resolve(typeof reader.result === 'string' ? reader.result : undefined); - reader.onerror = () => resolve(undefined); - reader.readAsDataURL(file); - }); - } - return { - id: crypto.randomUUID(), - file, - preview, + try { + while (pendingMessagesRef.current.length > 0) { + const { message, options } = pendingMessagesRef.current.shift()!; + const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {}; + + if (newChat) { + await createNewChatRef.current({ metadata }); + } + + // Convert File[] to FileAttachment[] + const fileAttachments: FileAttachment[] = await Promise.all( + attachments.map(async (file) => { + let preview: string | undefined; + if (file.type.startsWith('image/')) { + preview = await new Promise((resolve) => { + const reader = new FileReader(); + reader.onload = () => resolve(typeof reader.result === 'string' ? reader.result : undefined); + reader.onerror = () => resolve(undefined); + reader.readAsDataURL(file); + }); + } + return { + id: crypto.randomUUID(), + file, + preview, + }; + }) + ); + + await sendFnRef.current(message, fileAttachments.length > 0 ? fileAttachments : undefined, forwardedProps); + + if (openChat && setOpenRef.current) { + setOpenRef.current(true); + } + + // Wait for loading and pending approval to complete before processing next message + await new Promise((resolve) => { + const checkReady = () => { + setTimeout(() => { + if (!loadingRef.current && !hasPendingApprovalRef.current) { + resolve(); + } else { + checkReady(); + } + }, 100); }; - }) - ); - - await sendFnRef.current(message, fileAttachments.length > 0 ? fileAttachments : undefined, forwardedProps); - - if (openChat && setOpenRef.current) { - setOpenRef.current(true); + checkReady(); + }); } - - // Wait for loading and pending approval to complete before processing next message - await new Promise((resolve) => { - const checkReady = () => { - setTimeout(() => { - if (!loadingRef.current && !hasPendingApprovalRef.current) { - resolve(); - } else { - checkReady(); - } - }, 100); - }; - checkReady(); - }); + } finally { + isProcessingQueueRef.current = false; } - - isProcessingQueueRef.current = false; }, []); const sendMessage = useCallback(async (message: string, options?: SendMessageOptions): Promise => { From a348bab2b9542e1575f3c21073f1ea7baaec99aa Mon Sep 17 00:00:00 2001 From: Cyril Date: Mon, 7 Sep 2026 18:24:33 +0900 Subject: [PATCH 2/3] fix(client): clear pending message queue when a queued send rejects A rejecting sendFn left later queued messages in pendingMessagesRef even though the finally block lowered isProcessingQueueRef. A later unrelated sendMessage would then drain and send the abandoned message ahead of its own. Clear the queue in a catch before rethrowing. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_014qPDCJH1oDExv1xYg2DkV4 --- .../client/src/hooks/useMessageQueue.test.ts | 39 +++++++++++++++++++ packages/client/src/hooks/useMessageQueue.ts | 3 ++ 2 files changed, 42 insertions(+) diff --git a/packages/client/src/hooks/useMessageQueue.test.ts b/packages/client/src/hooks/useMessageQueue.test.ts index 218d2108..3d87abc8 100644 --- a/packages/client/src/hooks/useMessageQueue.test.ts +++ b/packages/client/src/hooks/useMessageQueue.test.ts @@ -35,4 +35,43 @@ describe('useMessageQueue', () => { expect(sendFn).toHaveBeenNthCalledWith(1, 'fails', undefined, undefined); expect(sendFn).toHaveBeenNthCalledWith(2, 'second', undefined, undefined); }); + + test('clears a message queued behind a failing send instead of delivering it to a later caller', async () => { + let rejectA: (error: Error) => void = () => {}; + const sendFn = mock((message: string) => { + if (message === 'A') { + return new Promise((_resolve, reject) => { + rejectA = reject; + }); + } + return Promise.resolve(); + }); + + const { result } = renderHook(() => useMessageQueue(createOptions(sendFn))); + + let sendAPromise: Promise = Promise.resolve(); + await act(async () => { + sendAPromise = result.current.sendMessage('A'); + // Let the queue processor reach the in-flight sendFn('A') call. + await new Promise((resolve) => setTimeout(resolve, 0)); + }); + + await act(async () => { + // Queued behind the in-flight 'A' drain; today this is unreachable until 'A' settles. + await result.current.sendMessage('B'); + }); + + await act(async () => { + rejectA(new Error('upload failed')); + await expect(sendAPromise).rejects.toThrow('upload failed'); + }); + + await act(async () => { + await result.current.sendMessage('C'); + }); + + const sentMessages = sendFn.mock.calls.map((call) => call[0]); + expect(sentMessages).not.toContain('B'); + expect(sentMessages).toEqual(['A', 'C']); + }); }); diff --git a/packages/client/src/hooks/useMessageQueue.ts b/packages/client/src/hooks/useMessageQueue.ts index f468c6f5..98424777 100644 --- a/packages/client/src/hooks/useMessageQueue.ts +++ b/packages/client/src/hooks/useMessageQueue.ts @@ -146,6 +146,9 @@ export function useMessageQueue({ checkReady(); }); } + } catch (error) { + pendingMessagesRef.current = []; + throw error; } finally { isProcessingQueueRef.current = false; } From 4427ed59a740a27173c58b12e425b1c8ecf1a3fd Mon Sep 17 00:00:00 2001 From: Zachary Davison Date: Tue, 8 Sep 2026 13:04:41 +0200 Subject: [PATCH 3/3] fix(client): settle each queued message's caller on its own outcome A send that rejected used to clear every message queued behind it, even though those callers had already been told their send succeeded. Each queued entry now carries its own promise: a failing send rejects only its caller and the loop continues to the next message. Co-Authored-By: Claude Fable 5.1 --- .../client/src/hooks/useMessageQueue.test.ts | 38 ++++-- packages/client/src/hooks/useMessageQueue.ts | 129 ++++++++++-------- 2 files changed, 97 insertions(+), 70 deletions(-) diff --git a/packages/client/src/hooks/useMessageQueue.test.ts b/packages/client/src/hooks/useMessageQueue.test.ts index 3d87abc8..2ce21945 100644 --- a/packages/client/src/hooks/useMessageQueue.test.ts +++ b/packages/client/src/hooks/useMessageQueue.test.ts @@ -36,7 +36,7 @@ describe('useMessageQueue', () => { expect(sendFn).toHaveBeenNthCalledWith(2, 'second', undefined, undefined); }); - test('clears a message queued behind a failing send instead of delivering it to a later caller', async () => { + function queueBehindFailingSend() { let rejectA: (error: Error) => void = () => {}; const sendFn = mock((message: string) => { if (message === 'A') { @@ -46,32 +46,42 @@ describe('useMessageQueue', () => { } return Promise.resolve(); }); - const { result } = renderHook(() => useMessageQueue(createOptions(sendFn))); + return { sendFn, result, rejectA: (error: Error) => rejectA(error) }; + } - let sendAPromise: Promise = Promise.resolve(); + test('resolves a queued caller only once its message has been sent', async () => { + const { result } = queueBehindFailingSend(); + + let bSettled = false; await act(async () => { - sendAPromise = result.current.sendMessage('A'); - // Let the queue processor reach the in-flight sendFn('A') call. + result.current.sendMessage('A').catch(() => {}); + await new Promise((resolve) => setTimeout(resolve, 0)); + result.current.sendMessage('B').then(() => { bSettled = true; }, () => { bSettled = true; }); await new Promise((resolve) => setTimeout(resolve, 0)); }); + expect(bSettled).toBe(false); + }); + + test('delivers a message queued behind a failing send', async () => { + const { sendFn, result, rejectA } = queueBehindFailingSend(); + + let sendAPromise: Promise = Promise.resolve(); + let sendBPromise: Promise = Promise.resolve(); await act(async () => { - // Queued behind the in-flight 'A' drain; today this is unreachable until 'A' settles. - await result.current.sendMessage('B'); + sendAPromise = result.current.sendMessage('A'); + await new Promise((resolve) => setTimeout(resolve, 0)); + sendBPromise = result.current.sendMessage('B'); + await new Promise((resolve) => setTimeout(resolve, 0)); }); await act(async () => { rejectA(new Error('upload failed')); await expect(sendAPromise).rejects.toThrow('upload failed'); + await sendBPromise; }); - await act(async () => { - await result.current.sendMessage('C'); - }); - - const sentMessages = sendFn.mock.calls.map((call) => call[0]); - expect(sentMessages).not.toContain('B'); - expect(sentMessages).toEqual(['A', 'C']); + expect(sendFn.mock.calls.map((call) => call[0])).toEqual(['A', 'B']); }); }); diff --git a/packages/client/src/hooks/useMessageQueue.ts b/packages/client/src/hooks/useMessageQueue.ts index 98424777..a2aaec33 100644 --- a/packages/client/src/hooks/useMessageQueue.ts +++ b/packages/client/src/hooks/useMessageQueue.ts @@ -46,6 +46,13 @@ export interface UseMessageQueueReturn { sendMessage: (message: string, options?: SendMessageOptions) => Promise; } +interface QueuedMessage { + message: string; + options?: SendMessageOptions; + resolve: () => void; + reject: (error: unknown) => void; +} + /** * Hook for queuing and sending programmatic messages. * @@ -64,7 +71,7 @@ export function useMessageQueue({ loading, hasPendingApproval, }: UseMessageQueueOptions): UseMessageQueueReturn { - const pendingMessagesRef = useRef>([]); + const pendingMessagesRef = useRef([]); const isProcessingQueueRef = useRef(false); // Use refs for callbacks that may change between renders. @@ -90,77 +97,87 @@ export function useMessageQueue({ hasPendingApprovalRef.current = hasPendingApproval; }, [hasPendingApproval]); + const sendQueuedMessage = useCallback(async (message: string, options?: SendMessageOptions) => { + const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {}; + + if (newChat) { + await createNewChatRef.current({ metadata }); + } + + // Convert File[] to FileAttachment[] + const fileAttachments: FileAttachment[] = await Promise.all( + attachments.map(async (file) => { + let preview: string | undefined; + if (file.type.startsWith('image/')) { + preview = await new Promise((resolve) => { + const reader = new FileReader(); + reader.onload = () => resolve(typeof reader.result === 'string' ? reader.result : undefined); + reader.onerror = () => resolve(undefined); + reader.readAsDataURL(file); + }); + } + return { + id: crypto.randomUUID(), + file, + preview, + }; + }) + ); + + await sendFnRef.current(message, fileAttachments.length > 0 ? fileAttachments : undefined, forwardedProps); + + if (openChat && setOpenRef.current) { + setOpenRef.current(true); + } + }, []); + + const waitUntilIdle = useCallback(() => new Promise((resolve) => { + const checkReady = () => { + setTimeout(() => { + if (!loadingRef.current && !hasPendingApprovalRef.current) { + resolve(); + } else { + checkReady(); + } + }, 100); + }; + checkReady(); + }), []); + + // Each queued message settles its own caller's promise, so one failing send + // neither blocks nor discards the messages queued behind it. const processMessageQueue = useCallback(async () => { - if (isProcessingQueueRef.current || pendingMessagesRef.current.length === 0) { + if (isProcessingQueueRef.current) { return; } isProcessingQueueRef.current = true; - try { while (pendingMessagesRef.current.length > 0) { - const { message, options } = pendingMessagesRef.current.shift()!; - const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {}; - - if (newChat) { - await createNewChatRef.current({ metadata }); + const entry = pendingMessagesRef.current.shift()!; + try { + await sendQueuedMessage(entry.message, entry.options); + } catch (error) { + entry.reject(error); + continue; } - - // Convert File[] to FileAttachment[] - const fileAttachments: FileAttachment[] = await Promise.all( - attachments.map(async (file) => { - let preview: string | undefined; - if (file.type.startsWith('image/')) { - preview = await new Promise((resolve) => { - const reader = new FileReader(); - reader.onload = () => resolve(typeof reader.result === 'string' ? reader.result : undefined); - reader.onerror = () => resolve(undefined); - reader.readAsDataURL(file); - }); - } - return { - id: crypto.randomUUID(), - file, - preview, - }; - }) - ); - - await sendFnRef.current(message, fileAttachments.length > 0 ? fileAttachments : undefined, forwardedProps); - - if (openChat && setOpenRef.current) { - setOpenRef.current(true); - } - - // Wait for loading and pending approval to complete before processing next message - await new Promise((resolve) => { - const checkReady = () => { - setTimeout(() => { - if (!loadingRef.current && !hasPendingApprovalRef.current) { - resolve(); - } else { - checkReady(); - } - }, 100); - }; - checkReady(); - }); + entry.resolve(); + await waitUntilIdle(); } - } catch (error) { - pendingMessagesRef.current = []; - throw error; } finally { isProcessingQueueRef.current = false; } - }, []); + }, [sendQueuedMessage, waitUntilIdle]); - const sendMessage = useCallback(async (message: string, options?: SendMessageOptions): Promise => { + const sendMessage = useCallback((message: string, options?: SendMessageOptions): Promise => { if (!connected) { - throw new Error('Not connected to UseAI server'); + return Promise.reject(new Error('Not connected to UseAI server')); } - pendingMessagesRef.current.push({ message, options }); - await processMessageQueue(); + return new Promise((resolve, reject) => { + pendingMessagesRef.current.push({ message, options, resolve, reject }); + void processMessageQueue(); + }); }, [connected, processMessageQueue]); return { sendMessage };