diff --git a/client/src/hooks/Chat/__tests__/queue.spec.ts b/client/src/hooks/Chat/__tests__/queue.spec.ts new file mode 100644 index 00000000000..ea127ac5e8e --- /dev/null +++ b/client/src/hooks/Chat/__tests__/queue.spec.ts @@ -0,0 +1,118 @@ +import type { TMessage } from 'librechat-data-provider'; +import { resolveDetachedRunEnd } from '../queue'; + +const CONVO_ID = 'convo-detached'; +const USER_ID = 'user-1'; + +const response = (overrides: Partial = {}): TMessage => + ({ + messageId: 'response-1', + parentMessageId: USER_ID, + conversationId: CONVO_ID, + isCreatedByUser: false, + text: 'done', + ...overrides, + }) as TMessage; + +const userMessage = { messageId: USER_ID, isCreatedByUser: true } as TMessage; + +describe('resolveDetachedRunEnd', () => { + it('reports a completed run with its persisted response', () => { + const end = resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, [ + userMessage, + response(), + ]); + expect(end).toEqual( + expect.objectContaining({ + conversationId: CONVO_ID, + outcome: 'completed', + responseMessageId: 'response-1', + }), + ); + }); + + it('reports a failed run as an error, so the queue waits for a manual send', () => { + const end = resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, [ + response({ error: true }), + ]); + expect(end?.outcome).toBe('error'); + expect(end?.responseMessageId).toBeUndefined(); + }); + + it('reports an unfinished response as aborted', () => { + const end = resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, [ + response({ unfinished: true }), + ]); + expect(end?.outcome).toBe('aborted'); + }); + + it('returns nothing while no response to the run is persisted', () => { + expect(resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, [userMessage])).toBeNull(); + expect(resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, undefined)).toBeNull(); + }); + + it('picks the response the run created among regenerated siblings', () => { + const end = resolveDetachedRunEnd( + CONVO_ID, + { userMessageId: USER_ID, responseMessageId: 'response-2_' }, + [response({ messageId: 'response-1', error: true }), response({ messageId: 'response-2' })], + ); + expect(end).toEqual( + expect.objectContaining({ outcome: 'completed', responseMessageId: 'response-2' }), + ); + }); + + it('matches a persisted response id that itself ends in an underscore', () => { + const end = resolveDetachedRunEnd( + CONVO_ID, + { userMessageId: USER_ID, responseMessageId: 'response-2_' }, + [response({ messageId: 'response-1' }), response({ messageId: 'response-2_' })], + ); + expect(end?.responseMessageId).toBe('response-2_'); + }); + + it('leaves a run unresolved when the response it named is not loaded', () => { + expect( + resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID, responseMessageId: 'response-2' }, [ + response({ messageId: 'response-1' }), + ]), + ).toBeNull(); + }); + + it('treats the padded user id placeholder as naming no response', () => { + const end = resolveDetachedRunEnd( + CONVO_ID, + { userMessageId: USER_ID, responseMessageId: `${USER_ID}_` }, + [response({ messageId: 'server-response' })], + ); + expect(end?.responseMessageId).toBe('server-response'); + }); + + it('leaves a regeneration unresolved, since history still holds the reply it replaces', () => { + expect( + resolveDetachedRunEnd( + CONVO_ID, + { userMessageId: USER_ID, responseMessageId: 'response-1_', isRegenerate: true }, + [response({ messageId: 'response-1' })], + ), + ).toBeNull(); + }); + + it('carries the run epoch so the drain can match server admission receipts', () => { + const end = resolveDetachedRunEnd( + CONVO_ID, + { userMessageId: USER_ID, generationCreatedAt: 4200 }, + [response()], + ); + expect(end?.generationCreatedAt).toBe(4200); + }); + + it('does not guess between siblings when the run named no response', () => { + expect( + resolveDetachedRunEnd(CONVO_ID, { userMessageId: USER_ID }, [ + response({ messageId: 'response-1' }), + response({ messageId: 'response-2' }), + ]), + ).toBeNull(); + }); +}); diff --git a/client/src/hooks/Chat/queue.ts b/client/src/hooks/Chat/queue.ts index e4f5308d8ab..2f879f4fe09 100644 --- a/client/src/hooks/Chat/queue.ts +++ b/client/src/hooks/Chat/queue.ts @@ -224,6 +224,79 @@ export const drainAfterAbortByIndex = atomFamily((_index: string | number) => atom(false), ); +/** A run whose stream this pane closed while it was still generating, because the user left its + * chat. The server keeps going and deletes the job once it finishes, so the run's end is learned + * from the persisted response when the user comes back instead of from the stream. */ +export type DetachedRun = { + /** The user message the run answers. */ + userMessageId: string; + /** The response placeholder's id, when the submission carried one. */ + responseMessageId?: string; + /** The run's generation epoch, when the start response installed one. The queue drain matches + * server admission receipts against it. */ + generationCreatedAt?: number; + /** A regeneration rewrites a response that already exists in history, so history cannot tell + * whether the row it finds is the old reply or the new one. */ + isRegenerate?: boolean; +}; + +export const detachedRunByConvoId = atomFamily((_conversationId: string) => + atom(null), +); + +/** The user pressed Stop on the conversation's current run. A stopped run's response persists + * like a completed one, so a run stopped and then left before its abort event arrived must not + * be resolved as detached and drain the queue. Reset when the next run starts. */ +export const stopRequestedByConvoId = atomFamily((_conversationId: string) => atom(false)); + +/** + * The run end a detached run implies, read from the persisted response: `completed` lets the + * queue drain, while a stopped (`unfinished`) or failed response leaves the queue for a manual + * send, as it would had the stream been attached. Returns `null` while no response to the run is + * persisted, so nothing is drained on a guess. + */ +export function resolveDetachedRunEnd( + conversationId: string, + run: DetachedRun, + messages: TMessage[] | undefined, +): RunEnd | null { + if (run.isRegenerate === true) { + return null; + } + const responses = (messages ?? []).filter( + (message) => message.isCreatedByUser === false && message.parentMessageId === run.userMessageId, + ); + /** A fresh turn's placeholder is the user message id padded with `_`, which names no response; + * a real id is matched exactly first, since persisted ids may themselves end in `_`. */ + const exact = run.responseMessageId; + const unpadded = exact?.replace(/_+$/, ''); + const namesResponse = exact != null && unpadded !== run.userMessageId; + let response: TMessage | undefined; + if (namesResponse) { + response = + responses.find((message) => message.messageId === exact) ?? + responses.find((message) => message.messageId === unpadded); + } else if (responses.length === 1) { + response = responses[0]; + } + if (response == null) { + return null; + } + let outcome: RunEnd['outcome'] = 'completed'; + if (response.error === true) { + outcome = 'error'; + } else if (response.unfinished === true) { + outcome = 'aborted'; + } + return { + conversationId, + outcome, + endedAt: Date.now(), + ...(run.generationCreatedAt != null && { generationCreatedAt: run.generationCreatedAt }), + ...(outcome === 'completed' && { responseMessageId: response.messageId }), + }; +} + const clearFamily = (family: { getParams(): Iterable; remove(param: Param): void; @@ -244,4 +317,6 @@ export function resetQueueFamilies(): void { clearFamily(drainAfterAbortByIndex); clearFamily(runEndsByIndex); clearFamily(runEndByIndex); + clearFamily(detachedRunByConvoId); + clearFamily(stopRequestedByConvoId); } diff --git a/client/src/hooks/Chat/useChatHelpers.ts b/client/src/hooks/Chat/useChatHelpers.ts index 6d2e30c6d3f..7819821d5bb 100644 --- a/client/src/hooks/Chat/useChatHelpers.ts +++ b/client/src/hooks/Chat/useChatHelpers.ts @@ -10,9 +10,9 @@ import { useAbortStreamMutation, supportsGenerationProtocolV2, } from '~/data-provider'; +import { stopRequestedByConvoId, drainAfterAbortByIndex, runEndByIndex } from '~/hooks/Chat/queue'; import { useLatestMessage, useLatestMessageId } from '~/hooks/Messages/useLatestMessage'; import { siblingIdxFamily, siblingKey } from '~/components/Chat/Messages/Thread/state'; -import { drainAfterAbortByIndex, runEndByIndex } from '~/hooks/Chat/queue'; import useChatFunctions from '~/hooks/Chat/useChatFunctions'; import useSteerConvert from '~/hooks/Chat/useSteerConvert'; import { resolveAbortSteerTarget } from '~/utils'; @@ -241,6 +241,7 @@ export default function useChatHelpers(index = 0, paramId?: string): ChatContrac // start the NEXT submission) while the abort response is in flight; // the fallback clear below must not tear down that new run. const submissionAtAbort = captureSubmission(); + queueStore.set(stopRequestedByConvoId(conversationId), true); try { console.log('[useChatHelpers] Calling abort mutation for:', conversationId); const response = await abortStream({ @@ -352,6 +353,7 @@ export default function useChatHelpers(index = 0, paramId?: string): ChatContrac conversationId, endpoint, endpointType, + queueStore, activeGenerationCreatedAt, activeGenerationProtocolVersion, abortStream, diff --git a/client/src/hooks/SSE/__tests__/useResumableSSE.spec.ts b/client/src/hooks/SSE/__tests__/useResumableSSE.spec.ts index 7aea121390c..e186e79549d 100644 --- a/client/src/hooks/SSE/__tests__/useResumableSSE.spec.ts +++ b/client/src/hooks/SSE/__tests__/useResumableSSE.spec.ts @@ -328,7 +328,12 @@ import useResumableSSE, { ABORT_SWEEP_STATUSES, } from '~/hooks/SSE/useResumableSSE'; import useSSE from '~/hooks/SSE/useSSE'; -import { queuedMessagesByConvoId, resetQueueFamilies } from '~/hooks/Chat/queue'; +import { + queuedMessagesByConvoId, + stopRequestedByConvoId, + detachedRunByConvoId, + resetQueueFamilies, +} from '~/hooks/Chat/queue'; const CONV_ID = 'conv-abc-123'; @@ -1015,6 +1020,96 @@ describe('useResumableSSE', () => { unmount(); }); + describe('a run the user leaves mid-stream', () => { + const renderLeavable = async () => { + const chatHelpers = buildChatHelpers(); + const rendered = renderHook( + ({ current }: { current: TSubmission | null }) => useResumableSSE(current, chatHelpers), + { initialProps: { current: buildSubmission() as TSubmission | null } }, + ); + await flushMicrotasks(); + expect(mockSSEInstances.length).toBeGreaterThan(0); + return rendered; + }; + const detachedRun = () => getDefaultStore().get(detachedRunByConvoId(CONV_ID)); + + it('remembers the run when navigation clears the submission, so its end is read on return', async () => { + const { rerender, unmount } = await renderLeavable(); + rerender({ current: {} as TSubmission }); + expect(detachedRun()).toEqual({ userMessageId: 'msg-1', responseMessageId: 'resp-1' }); + unmount(); + }); + + it('remembers the run when switching to a saved chat clears the submission to null', async () => { + const { rerender, unmount } = await renderLeavable(); + rerender({ current: null }); + expect(detachedRun()).toEqual({ userMessageId: 'msg-1', responseMessageId: 'resp-1' }); + unmount(); + }); + + it('marks a regeneration and carries a resumed epoch when the user leaves', async () => { + const chatHelpers = buildChatHelpers(); + const resumed = { + ...buildSubmission(), + isRegenerate: true, + resumeStreamId: CONV_ID, + resumeGenerationCreatedAt: 4200, + } as TSubmission; + const { unmount } = renderHook(() => useResumableSSE(resumed, chatHelpers)); + await flushMicrotasks(); + unmount(); + expect(detachedRun()).toEqual( + expect.objectContaining({ isRegenerate: true, generationCreatedAt: 4200 }), + ); + }); + + it('remembers the run when the chat unmounts mid-stream', async () => { + const { unmount } = await renderLeavable(); + unmount(); + expect(detachedRun()).toEqual({ userMessageId: 'msg-1', responseMessageId: 'resp-1' }); + }); + + it('forgets a left run once a new run starts in that conversation', async () => { + const { rerender, unmount } = await renderLeavable(); + rerender({ current: null }); + expect(detachedRun()).not.toBeNull(); + rerender({ + current: buildSubmission({ + userMessage: { + messageId: 'msg-2', + conversationId: CONV_ID, + text: 'Next', + isCreatedByUser: true, + sender: 'User', + parentMessageId: 'resp-1', + }, + }), + }); + expect(detachedRun()).toBeNull(); + unmount(); + }); + + it('records no detached run once the run end already reached the drain', async () => { + mockFetchStreamStatus.mockResolvedValue({ active: false }); + const { rerender, unmount } = await renderLeavable(); + await act(async () => { + getLastSSE()._emit('error', { responseCode: 404 }); + }); + await waitFor(() => expect(mockSetRunEnd).toHaveBeenCalled()); + rerender({ current: {} as TSubmission }); + expect(detachedRun()).toBeNull(); + unmount(); + }); + + it('records no detached run for a run the user stopped before leaving', async () => { + const { rerender, unmount } = await renderLeavable(); + getDefaultStore().set(stopRequestedByConvoId(CONV_ID), true); + rerender({ current: {} as TSubmission }); + expect(detachedRun()).toBeNull(); + unmount(); + }); + }); + it('authorizes only the exact failed recovery source past the conversion tombstone', async () => { const parked = [ { steerId: 'failed-source', text: 'retry these words', createdAt: 1 }, diff --git a/client/src/hooks/SSE/__tests__/useResumeOnLoad.spec.tsx b/client/src/hooks/SSE/__tests__/useResumeOnLoad.spec.tsx index de535a194e3..226396b3984 100644 --- a/client/src/hooks/SSE/__tests__/useResumeOnLoad.spec.tsx +++ b/client/src/hooks/SSE/__tests__/useResumeOnLoad.spec.tsx @@ -1,4 +1,4 @@ -import { renderHook, act } from '@testing-library/react'; +import { renderHook, act, waitFor } from '@testing-library/react'; import { RecoilRoot, useRecoilValue, useSetRecoilState } from 'recoil'; import { QueryClient, QueryClientProvider } from '@tanstack/react-query'; import { Provider as JotaiProvider, createStore, useAtomValue } from 'jotai'; @@ -6,10 +6,15 @@ import { Constants, ContentTypes, QueryKeys } from 'librechat-data-provider'; import type { Agents, TMessage, TConversation, TSubmission } from 'librechat-data-provider'; import type { MutableSnapshot } from 'recoil'; import type { ReactNode } from 'react'; -import type { QueuedMessage } from '~/hooks/Chat/queue'; -import type { PendingSteer } from '~/hooks/Chat/queue'; +import type { QueuedMessage, PendingSteer, DetachedRun } from '~/hooks/Chat/queue'; +import { + settledQueuedTurnReceiptsByConvoId, + queuedMessagesByConvoId, + pendingRunEndByConvoId, + detachedRunByConvoId, + resetQueueFamilies, +} from '~/hooks/Chat/queue'; import { siblingIdxFamily, siblingKey } from '~/components/Chat/Messages/Thread/state'; -import { queuedMessagesByConvoId, resetQueueFamilies } from '~/hooks/Chat/queue'; import { pendingApprovalActionFamily } from '~/components/Chat/approval/state'; import { agentQueuedTurnsQueryKey } from '~/data-provider/SSE/queuedTurns'; import { revealedQueuedTurnFamily } from '~/store/steer'; @@ -107,6 +112,9 @@ function renderUseResumeOnLoad({ onQueuedMessages, submissionStart, onSubmissionStart, + detachedRun, + seedJotai, + seedQueryClient, }: { messages?: TMessage[]; getMessages?: () => TMessage[] | undefined; @@ -123,12 +131,20 @@ function renderUseResumeOnLoad({ onQueuedMessages?: (queued: QueuedMessage[]) => void; submissionStart?: number; onSubmissionStart?: (submissionStart: number | null) => void; + detachedRun?: DetachedRun; + seedJotai?: (store: ReturnType) => void; + seedQueryClient?: (client: QueryClient) => void; }) { const getMessages = jest.fn(getMessagesOverride ?? (() => messages)); const jotaiStore = createStore(); + if (detachedRun != null) { + jotaiStore.set(detachedRunByConvoId(conversationId), detachedRun); + } + seedJotai?.(jotaiStore); const queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } }, }); + seedQueryClient?.(queryClient); let setSubmissionState: ((submission: TSubmission | null) => void) | undefined; let setIsSubmittingState: ((value: boolean) => void) | undefined; let setAttachedEpochState: ((value: number | null) => void) | undefined; @@ -2731,6 +2747,235 @@ describe('useResumeOnLoad', () => { expect(observedSteers[observedSteers.length - 1]).toEqual([failedChip]); }); + it('parks the end of a run that finished after the user left, so its queue drains', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const completedResponse = { + messageId: 'response-detached', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'finished while away', + } as TMessage; + + const { jotaiStore } = renderUseResumeOnLoad({ + messages: [buildUserMessage(CONVERSATION_ID), completedResponse], + detachedRun: { userMessageId: USER_MESSAGE_ID }, + }); + await act(async () => { + await Promise.resolve(); + }); + + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toEqual( + expect.objectContaining({ + conversationId: CONVERSATION_ID, + outcome: 'completed', + responseMessageId: 'response-detached', + }), + ); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).toBeNull(); + }); + + it('keeps the detached run until a history refetch shows its response', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + let history: TMessage[] = [buildUserMessage(CONVERSATION_ID)]; + const { jotaiStore } = renderUseResumeOnLoad({ + getMessages: () => history, + detachedRun: { userMessageId: USER_MESSAGE_ID }, + }); + /** The first read found no response; it is saved before the history refetch settles. */ + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toBeNull(); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).not.toBeNull(); + history = [ + ...history, + { + messageId: 'response-late', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'saved after the first read', + } as TMessage, + ]; + + await waitFor(() => + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toEqual( + expect.objectContaining({ outcome: 'completed', responseMessageId: 'response-late' }), + ), + ); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).toBeNull(); + }); + + it('resolves a detached run from refetched history, not the response it left cached', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const historyKey = [QueryKeys.messages, CONVERSATION_ID]; + const leftResponse = { + messageId: 'response-detached', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'streamed before leaving', + } as TMessage; + let client: QueryClient | undefined; + const { jotaiStore } = renderUseResumeOnLoad({ + getMessages: () => client?.getQueryData(historyKey), + detachedRun: { userMessageId: USER_MESSAGE_ID, responseMessageId: 'response-detached' }, + seedQueryClient: (queryClient) => { + client = queryClient; + queryClient.setQueryDefaults(historyKey, { + queryFn: async () => [ + buildUserMessage(CONVERSATION_ID), + { ...leftResponse, unfinished: true }, + ], + }); + queryClient.setQueryData(historyKey, [buildUserMessage(CONVERSATION_ID), leftResponse]); + }, + }); + + await waitFor(() => + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toEqual( + expect.objectContaining({ outcome: 'aborted' }), + ), + ); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).toBeNull(); + }); + + it('keeps the detached run for the next visit when the history refetch fails', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const historyKey = [QueryKeys.messages, CONVERSATION_ID]; + const queryFn = jest.fn(async (): Promise => { + throw new Error('history unavailable'); + }); + let client: QueryClient | undefined; + const { jotaiStore } = renderUseResumeOnLoad({ + getMessages: () => client?.getQueryData(historyKey), + detachedRun: { userMessageId: USER_MESSAGE_ID, responseMessageId: 'response-detached' }, + seedQueryClient: (queryClient) => { + client = queryClient; + queryClient.setQueryDefaults(historyKey, { queryFn }); + queryClient.setQueryData(historyKey, [ + buildUserMessage(CONVERSATION_ID), + { + messageId: 'response-detached', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'streamed before leaving', + } as TMessage, + ]); + }, + }); + + await waitFor(() => expect(queryFn).toHaveBeenCalled()); + await act(async () => { + await Promise.resolve(); + }); + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toBeNull(); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).not.toBeNull(); + }); + + it('leaves a run without an epoch for a manual send when the server shares the queue', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const { jotaiStore } = renderUseResumeOnLoad({ + messages: [ + buildUserMessage(CONVERSATION_ID), + { + messageId: 'response-detached', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'finished while away', + } as TMessage, + ], + detachedRun: { userMessageId: USER_MESSAGE_ID }, + seedJotai: (store) => + store.set(settledQueuedTurnReceiptsByConvoId(CONVERSATION_ID), [ + { clientRequestId: 'queued-1', status: 'admitted', effectivePredecessorCreatedAt: 41 }, + ]), + }); + await act(async () => { + await Promise.resolve(); + }); + + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toBeNull(); + expect(jotaiStore.get(detachedRunByConvoId(CONVERSATION_ID))).toBeNull(); + }); + + it('parks a shared-queue run end with its epoch, so its admission receipt is matched', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const { jotaiStore } = renderUseResumeOnLoad({ + messages: [ + buildUserMessage(CONVERSATION_ID), + { + messageId: 'response-detached', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'finished while away', + } as TMessage, + ], + detachedRun: { userMessageId: USER_MESSAGE_ID, generationCreatedAt: 41 }, + seedJotai: (store) => + store.set(settledQueuedTurnReceiptsByConvoId(CONVERSATION_ID), [ + { clientRequestId: 'queued-1', status: 'admitted', effectivePredecessorCreatedAt: 41 }, + ]), + }); + await act(async () => { + await Promise.resolve(); + }); + + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toEqual( + expect.objectContaining({ outcome: 'completed', generationCreatedAt: 41 }), + ); + }); + + it('parks nothing for a conversation this pane never left mid-run', async () => { + mockUseStreamStatus.mockReturnValue({ + isSuccess: true, + isFetching: false, + data: { active: false }, + }); + const { jotaiStore } = renderUseResumeOnLoad({ + messages: [ + buildUserMessage(CONVERSATION_ID), + { + messageId: 'response-old', + parentMessageId: USER_MESSAGE_ID, + conversationId: CONVERSATION_ID, + isCreatedByUser: false, + text: 'history', + } as TMessage, + ], + }); + await act(async () => { + await Promise.resolve(); + }); + + expect(jotaiStore.get(pendingRunEndByConvoId(CONVERSATION_ID))).toBeNull(); + }); + it('converts resumeState.pendingSteers to queued when inactive (expired action, unparked queue)', async () => { const observedSteers: PendingSteer[][] = []; const observedQueues: QueuedMessage[][] = []; diff --git a/client/src/hooks/SSE/useResumableSSE.ts b/client/src/hooks/SSE/useResumableSSE.ts index 6b68f37b7e4..86a16d35b9c 100644 --- a/client/src/hooks/SSE/useResumableSSE.ts +++ b/client/src/hooks/SSE/useResumableSSE.ts @@ -41,8 +41,8 @@ import type { TContextUsageEvent, ChatStreamConnection, } from 'librechat-data-provider'; +import type { QueuedMessageOrigin, DrainAfterAbort, RunEnd } from '~/hooks/Chat/queue'; import type { ActiveJobsResponse, StreamStatusResponse } from '~/data-provider'; -import type { DrainAfterAbort, QueuedMessageOrigin } from '~/hooks/Chat/queue'; import type { GenerationProtocolVersion } from '~/data-provider'; import type { EventHandlerParams } from './useEventHandlers'; import type { TResData, TFinalResData } from '~/common'; @@ -86,6 +86,13 @@ import { supportsGenerationProtocolV2, GENERATION_PROTOCOL_VERSION, } from '~/data-provider'; +import { + stopRequestedByConvoId, + detachedRunByConvoId, + drainAfterAbortByIndex, + queuedMessagesByConvoId, + runEndByIndex, +} from '~/hooks/Chat/queue'; import { recoveryDispositionsFamily, canRestoreRecovery, @@ -95,7 +102,6 @@ import useEventHandlers, { buildCreatedInitialResponse, keepLocalCodeApprovalMode, } from './useEventHandlers'; -import { drainAfterAbortByIndex, queuedMessagesByConvoId, runEndByIndex } from '~/hooks/Chat/queue'; import { pendingApprovalActionFamily } from '~/components/Chat/approval/state'; import { withSubmittedCodeDecision } from '~/hooks/Agents/codeDecision'; import { useChatTransport } from '~/Providers/ChatTransportContext'; @@ -1333,6 +1339,8 @@ export default function useResumableSSE( [convertSteersToQueued], ); + /** The epoch this pane last installed, kept so a run left mid-stream can carry it. */ + const liveEpochRef = useRef<{ conversationId: string; createdAt: number } | null>(null); /** Conversation ids are reused by successive agent turns. Keep the exact * live epoch beside the UI controls, and only clear it if the terminal event * still belongs to the epoch that installed it. */ @@ -1349,6 +1357,7 @@ export default function useResumableSSE( return false; } set(state, next); + liveEpochRef.current = next == null ? null : { conversationId, createdAt: next }; set( store.activeGenerationProtocolVersionByConvoId(conversationId), next == null ? 1 : generationProtocolVersion, @@ -1358,7 +1367,17 @@ export default function useResumableSSE( [], ); - const setRunEnd = useSetAtom(runEndByIndex(runIndex)); + const publishRunEnd = useSetAtom(runEndByIndex(runIndex)); + /** Whether the current submission's run end already reached the queue drain. A run that + * ended attached must not also be resolved as detached when the user later leaves. */ + const runEndPublishedRef = useRef(false); + const setRunEnd = useCallback( + (end: RunEnd) => { + runEndPublishedRef.current = true; + publishRunEnd(end); + }, + [publishRunEnd], + ); const setDrainAfterAbort = useSetAtom(drainAfterAbortByIndex(runIndex)); const clearDrainAfterAbort = useCallback( (conversationId: string, generationCreatedAt?: number) => { @@ -4289,6 +4308,13 @@ export default function useResumableSSE( }); submissionRef.current = submission; + runEndPublishedRef.current = false; + if (submission.conversation?.conversationId != null) { + /** A run starting here owns the conversation's end from now on: an earlier run this pane + * left is superseded, and a Stop belonged to that earlier run. */ + jotaiStore.set(stopRequestedByConvoId(submission.conversation.conversationId), false); + jotaiStore.set(detachedRunByConvoId(submission.conversation.conversationId), null); + } const startController = new AbortController(); const { signal } = startController; const isCurrentEffect = () => !signal.aborted && submissionRef.current === submission; @@ -4814,6 +4840,35 @@ export default function useResumableSSE( stopForegroundReattachRef.current = null; // Reset reconnect counter before closing (so abort handler doesn't think we're reconnecting) reconnectAttemptRef.current = 0; + /** The pane is moving off a run that is still generating (another chat, a new chat, or + * an unmount): the server keeps going with no subscriber and deletes the job when done, + * so no terminal event will reach the queue drain. Remember the run; on return its end is + * read from history. A run that already published its end, or that the user stopped, is + * not remembered. */ + const closing = submissionRef.current; + const closingConvoId = closing?.conversation?.conversationId; + const closingUserMessageId = closing?.userMessage?.messageId; + if ( + streamRef.current != null && + !runEndPublishedRef.current && + closingConvoId != null && + closingConvoId !== Constants.NEW_CONVO && + closingUserMessageId != null && + !jotaiStore.get(stopRequestedByConvoId(closingConvoId)) + ) { + const liveEpoch = liveEpochRef.current; + const generationCreatedAt = + liveEpoch?.conversationId === closingConvoId + ? liveEpoch.createdAt + : (closing as TSubmission & { resumeGenerationCreatedAt?: number }) + .resumeGenerationCreatedAt; + jotaiStore.set(detachedRunByConvoId(closingConvoId), { + userMessageId: closingUserMessageId, + responseMessageId: closing?.initialResponse?.messageId, + ...(generationCreatedAt != null && { generationCreatedAt }), + ...(closing?.isRegenerate === true && { isRegenerate: true }), + }); + } streamRef.current?.abort(); streamRef.current = null; // Clear handler maps to prevent memory leaks and stale state diff --git a/client/src/hooks/SSE/useResumeOnLoad.ts b/client/src/hooks/SSE/useResumeOnLoad.ts index b6ee8e6be37..fd6f36aa2bf 100644 --- a/client/src/hooks/SSE/useResumeOnLoad.ts +++ b/client/src/hooks/SSE/useResumeOnLoad.ts @@ -34,6 +34,13 @@ import { extendActiveJobsGrace, ACTIVE_JOBS_SUCCESSOR_GRACE_MS, } from '~/data-provider'; +import { + settledQueuedTurnReceiptsByConvoId, + queuedMessagesByConvoId, + resolveDetachedRunEnd, + pendingRunEndByConvoId, + detachedRunByConvoId, +} from '~/hooks/Chat/queue'; import { getGenerationProtocolVersion, supportsGenerationProtocolV2, @@ -1001,10 +1008,50 @@ export default function useResumeOnLoad( // `unrecoveredSteers` above — same empty-list reconcile as the resume path. settleAppliedSteerParts(conversationId, getMessages()); restoreSteerChips(conversationId, undefined); + /** A run this pane stopped watching when the user left has ended on the server, which + * deleted the job. Its persisted response tells how it ended; parking that end lets the + * queue drain send a follow-up queued during the run, as an attached run would have. + * The cached history still holds the optimistic response the pane left, which lacks the + * persisted stop and error flags, so only a successful refetch of history may resolve it; + * when that fails or predates the saved response, the next visit tries again. */ + const detachedFamily = detachedRunByConvoId(conversationId); + const detachedRun = jotaiStore.get(detachedFamily); + const parkDetachedEnd = () => { + if (detachedRun == null || jotaiStore.get(detachedFamily) !== detachedRun) { + return; + } + /** Without the run's epoch the drain cannot match a server admission receipt, so a run + * whose queue the server shares is left for a manual send rather than guessed at. */ + const serverSharesQueue = + jotaiStore.get(settledQueuedTurnReceiptsByConvoId(conversationId)).length > 0 || + jotaiStore + .get(queuedMessagesByConvoId(conversationId)) + .some((item) => item.server != null); + if (detachedRun.generationCreatedAt == null && serverSharesQueue) { + jotaiStore.set(detachedFamily, null); + return; + } + const end = resolveDetachedRunEnd(conversationId, detachedRun, getMessages()); + if (end == null) { + return; + } + jotaiStore.set(detachedFamily, null); + jotaiStore.set(pendingRunEndByConvoId(conversationId), end); + }; + if (detachedRun != null) { + queryClient + .refetchQueries( + { queryKey: [QueryKeys.messages, conversationId], exact: true }, + { throwOnError: true }, + ) + .then(parkDetachedEnd, () => undefined); + } processedConvoRef.current = conversationId; return; } + /** Still generating: the stream this resume attaches will deliver the run's own end. */ + jotaiStore.set(detachedRunByConvoId(conversationId), null); processedConvoRef.current = conversationId; if (handoffGenerationKey != null) { consumedHandoffGenerationRef.current = handoffGenerationKey; diff --git a/e2e/specs/mock/scenarios/queue-owners.spec.ts b/e2e/specs/mock/scenarios/queue-owners.spec.ts index 05fc9de584b..c39440ba098 100644 --- a/e2e/specs/mock/scenarios/queue-owners.spec.ts +++ b/e2e/specs/mock/scenarios/queue-owners.spec.ts @@ -1,8 +1,12 @@ import { expect, test } from '@playwright/test'; import type { Page } from '@playwright/test'; +import type { TMessage } from 'librechat-data-provider'; import { selectMockEndpoint, + getAccessToken, messagesView, + requestJson, + fetchJson, sendMessage, replyPrompt, replyText, @@ -13,12 +17,13 @@ import { /** * The follow-up queue, its run-end signals and the interrupt-drain flag are chat-owned Jotai - * state. Queueing and draining are covered by the composer-queue scenarios; this one drives the - * interrupt-drain flag through a real stopped run. + * state. Queueing and draining are covered by the composer-queue scenarios; these drive the + * interrupt-drain flag through a real stopped run, and a run that ends while its chat is left. */ const messageInput = (page: Page) => page.getByRole('textbox', { name: 'Message input' }); const duringRunSendButton = (page: Page) => page.getByTestId('during-run-send-button'); +const queuedRows = (page: Page) => page.getByTestId('queued-message-row'); const messageTurns = (page: Page) => messagesView(page).locator('.message-render'); const uniqueLabel = (prefix: string) => `${prefix}-${Date.now()}-${Math.floor(Math.random() * 1e6)}`; @@ -38,6 +43,29 @@ async function typeDuringRun(page: Page, text: string) { await expect(duringRunSendButton(page)).toBeVisible({ timeout: 5000 }); } +/** Waits until the server has saved the slow run's complete reply. */ +async function waitForServerToFinish(page: Page, conversationId: string) { + const token = await getAccessToken(page); + await expect + .poll( + async () => { + const messages = await fetchJson( + page, + `/api/messages/${encodeURIComponent(conversationId)}`, + token, + ); + return messages.some( + (message) => + !message.isCreatedByUser && + message.unfinished !== true && + JSON.stringify(message.content ?? message.text ?? '').includes('chunk-159'), + ); + }, + { timeout: 60000 }, + ) + .toBe(true); +} + test.describe('chat-owned queue state', () => { test.beforeEach(async ({ page }) => { /** Steer is the plain-Enter default here, so Cmd/Ctrl+Enter is the queue path. */ @@ -69,4 +97,120 @@ test.describe('chat-owned queue state', () => { await expect(messageTurns(page).nth(5)).toContainText(MOCK_REPLY_TEXT, { timeout: 30000 }); await expect(messagesView(page).getByText('chunk-159')).toHaveCount(0); }); + + test('a follow-up queued in a chat the user left sends on return @scenario:parked-run-end-drains-on-return', async ({ + page, + }) => { + test.setTimeout(150000); + const label = uniqueLabel('parked'); + const followUp = `Parked follow-up ${label}`; + + await page.goto(NEW_CHAT_PATH, { timeout: 10000 }); + await selectMockEndpoint(page, MOCK_ENDPOINTS[0]); + const conversationId = await establishConversation(page, `parked-setup-${label}`); + + const run = await sendMessage(page, `E2E_SLOW_REPLY:${label}`); + expect(run.ok()).toBeTruthy(); + await typeDuringRun(page, followUp); + await messageInput(page).press('ControlOrMeta+Enter'); + await expect(queuedRows(page).filter({ hasText: followUp })).toBeVisible({ timeout: 10000 }); + + /** Leave through the router, not a reload, so the in-memory queue survives the visit. */ + await page.evaluate((path) => { + window.history.pushState({}, '', path); + window.dispatchEvent(new PopStateEvent('popstate')); + }, NEW_CHAT_PATH); + await expect(page).toHaveURL(/\/c\/new$/); + + /** The run finishes while its chat is not on screen. */ + await waitForServerToFinish(page, conversationId); + await expect(messagesView(page).getByText(followUp)).toHaveCount(0); + + await page.goBack(); + await expect(page).toHaveURL(new RegExp(`/c/${conversationId}(\\?.*)?$`)); + await expect( + messagesView(page).locator('.user-turn').filter({ hasText: followUp }), + ).toBeVisible({ timeout: 30000 }); + await expect(queuedRows(page).filter({ hasText: followUp })).toHaveCount(0); + }); + + test('a follow-up queued before the user stopped and left the run stays queued on return @scenario:stopped-run-left-keeps-follow-up-queued', async ({ + page, + }) => { + test.setTimeout(150000); + const label = uniqueLabel('stopped'); + const followUp = `Stopped follow-up ${label}`; + + await page.goto(NEW_CHAT_PATH, { timeout: 10000 }); + await selectMockEndpoint(page, MOCK_ENDPOINTS[0]); + const conversationId = await establishConversation(page, `stopped-setup-${label}`); + + const run = await sendMessage(page, `E2E_SLOW_REPLY:${label}`); + expect(run.ok()).toBeTruthy(); + await expect(messagesView(page).getByText('chunk-010')).toBeVisible({ timeout: 15000 }); + await typeDuringRun(page, followUp); + await messageInput(page).press('ControlOrMeta+Enter'); + await expect(queuedRows(page).filter({ hasText: followUp })).toBeVisible({ timeout: 10000 }); + + const stop = page.getByRole('button', { name: 'Stop generating' }); + await stop.click(); + await expect(stop).toBeHidden({ timeout: 15000 }); + + await page.evaluate((path) => { + window.history.pushState({}, '', path); + window.dispatchEvent(new PopStateEvent('popstate')); + }, NEW_CHAT_PATH); + await expect(page).toHaveURL(/\/c\/new$/); + + await page.goBack(); + await expect(page).toHaveURL(new RegExp(`/c/${conversationId}(\\?.*)?$`)); + await expect(queuedRows(page).filter({ hasText: followUp })).toBeVisible({ timeout: 15000 }); + await expect( + messagesView(page).locator('.user-turn').filter({ hasText: followUp }), + ).toHaveCount(0); + }); + + test('a follow-up queued in a chat the user left for a saved chat sends on return @scenario:parked-run-end-drains-after-switching-chats', async ({ + page, + }) => { + const width = page.viewportSize()?.width ?? 0; + test.skip(width < 768, 'the conversation list is in the drawer below md'); + test.setTimeout(150000); + const label = uniqueLabel('switched'); + const followUp = `Switched follow-up ${label}`; + const otherTitle = `Other chat ${label}`; + + /** A saved chat to switch to, titled so its list row can be found. */ + await page.goto(NEW_CHAT_PATH, { timeout: 10000 }); + await selectMockEndpoint(page, MOCK_ENDPOINTS[0]); + const otherId = await establishConversation(page, `switched-other-${label}`); + await requestJson(page, { + path: '/api/convos/update', + token: await getAccessToken(page), + method: 'POST', + body: { arg: { conversationId: otherId, title: otherTitle } }, + }); + + await page.goto(NEW_CHAT_PATH, { timeout: 10000 }); + await selectMockEndpoint(page, MOCK_ENDPOINTS[0]); + const conversationId = await establishConversation(page, `switched-setup-${label}`); + const run = await sendMessage(page, `E2E_SLOW_REPLY:${label}`); + expect(run.ok()).toBeTruthy(); + await typeDuringRun(page, followUp); + await messageInput(page).press('ControlOrMeta+Enter'); + await expect(queuedRows(page).filter({ hasText: followUp })).toBeVisible({ timeout: 10000 }); + + /** Leave the way a user does: pick the other saved chat in the conversation list. */ + await page.getByTestId('convo-item').filter({ hasText: otherTitle }).click(); + await expect(page).toHaveURL(new RegExp(`/c/${otherId}(\\?.*)?$`)); + + await waitForServerToFinish(page, conversationId); + + await page.goBack(); + await expect(page).toHaveURL(new RegExp(`/c/${conversationId}(\\?.*)?$`)); + await expect( + messagesView(page).locator('.user-turn').filter({ hasText: followUp }), + ).toBeVisible({ timeout: 30000 }); + await expect(queuedRows(page).filter({ hasText: followUp })).toHaveCount(0); + }); });