Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 87 additions & 0 deletions packages/client/src/hooks/useMessageQueue.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
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<void>) {
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);
});

function queueBehindFailingSend() {
let rejectA: (error: Error) => void = () => {};
const sendFn = mock((message: string) => {
if (message === 'A') {
return new Promise<void>((_resolve, reject) => {
rejectA = reject;
});
}
return Promise.resolve();
});
const { result } = renderHook(() => useMessageQueue(createOptions(sendFn)));
return { sendFn, result, rejectA: (error: Error) => rejectA(error) };
}

test('resolves a queued caller only once its message has been sent', async () => {
const { result } = queueBehindFailingSend();

let bSettled = false;
await act(async () => {
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<void> = Promise.resolve();
let sendBPromise: Promise<void> = Promise.resolve();
await act(async () => {
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;
});

expect(sendFn.mock.calls.map((call) => call[0])).toEqual(['A', 'B']);
});
});
134 changes: 78 additions & 56 deletions packages/client/src/hooks/useMessageQueue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,13 @@ export interface UseMessageQueueReturn {
sendMessage: (message: string, options?: SendMessageOptions) => Promise<void>;
}

interface QueuedMessage {
message: string;
options?: SendMessageOptions;
resolve: () => void;
reject: (error: unknown) => void;
}

/**
* Hook for queuing and sending programmatic messages.
*
Expand All @@ -64,7 +71,7 @@ export function useMessageQueue({
loading,
hasPendingApproval,
}: UseMessageQueueOptions): UseMessageQueueReturn {
const pendingMessagesRef = useRef<Array<{ message: string; options?: SendMessageOptions }>>([]);
const pendingMessagesRef = useRef<QueuedMessage[]>([]);
const isProcessingQueueRef = useRef(false);

// Use refs for callbacks that may change between renders.
Expand All @@ -90,72 +97,87 @@ export function useMessageQueue({
hasPendingApprovalRef.current = hasPendingApproval;
}, [hasPendingApproval]);

const processMessageQueue = useCallback(async () => {
if (isProcessingQueueRef.current || pendingMessagesRef.current.length === 0) {
return;
}
const sendQueuedMessage = useCallback(async (message: string, options?: SendMessageOptions) => {
const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {};

isProcessingQueueRef.current = true;
if (newChat) {
await createNewChatRef.current({ metadata });
}

while (pendingMessagesRef.current.length > 0) {
const { message, options } = pendingMessagesRef.current.shift()!;
const { newChat = false, attachments = [], openChat = true, metadata, forwardedProps } = options ?? {};
// 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<string | undefined>((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,
};
})
);

if (newChat) {
await createNewChatRef.current({ metadata });
}
await sendFnRef.current(message, fileAttachments.length > 0 ? fileAttachments : undefined, forwardedProps);

// 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<string | undefined>((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);
}
if (openChat && setOpenRef.current) {
setOpenRef.current(true);
}
}, []);

// Wait for loading and pending approval to complete before processing next message
await new Promise<void>((resolve) => {
const checkReady = () => {
setTimeout(() => {
if (!loadingRef.current && !hasPendingApprovalRef.current) {
resolve();
} else {
checkReady();
}
}, 100);
};
checkReady();
});
const waitUntilIdle = useCallback(() => new Promise<void>((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) {
return;
}

isProcessingQueueRef.current = false;
}, []);
isProcessingQueueRef.current = true;
try {
while (pendingMessagesRef.current.length > 0) {
const entry = pendingMessagesRef.current.shift()!;
try {
await sendQueuedMessage(entry.message, entry.options);
} catch (error) {
entry.reject(error);
continue;
}
entry.resolve();
await waitUntilIdle();
}
} finally {
isProcessingQueueRef.current = false;
}
}, [sendQueuedMessage, waitUntilIdle]);

const sendMessage = useCallback(async (message: string, options?: SendMessageOptions): Promise<void> => {
const sendMessage = useCallback((message: string, options?: SendMessageOptions): Promise<void> => {
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<void>((resolve, reject) => {
pendingMessagesRef.current.push({ message, options, resolve, reject });
void processMessageQueue();
});
}, [connected, processMessageQueue]);

return { sendMessage };
Expand Down
Loading