Skip to content
Open
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
118 changes: 118 additions & 0 deletions client/src/hooks/Chat/__tests__/queue.spec.ts
Original file line number Diff line number Diff line change
@@ -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> = {}): 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();
});
});
75 changes: 75 additions & 0 deletions client/src/hooks/Chat/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -224,6 +224,79 @@ export const drainAfterAbortByIndex = atomFamily((_index: string | number) =>
atom<DrainAfterAbort | false>(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<DetachedRun | null>(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);
Comment on lines +276 to +278

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Reject the pre-run regeneration row as terminal

When a regeneration is left and its history request completes before terminal persistence, the messages still contain the previous response under the unpadded ID. This fallback immediately accepts that old row as the detached run's result and usually classifies it as completed, so a regeneration that subsequently aborts or fails can auto-send the queued follow-up instead of leaving it for manual send. The fresh evidence beyond the earlier missing-response race is that regenerations deliberately use existingId_ while the pre-run existingId row already exists, meaning parkDetachedEnd does not take the new refetch/retry path at all. Require evidence that this row was updated by the detached generation before publishing its outcome.

AGENTS.md reference: AGENTS.md:L49-L52

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in d8269cc: a detached regeneration is no longer resolved from history (resolveDetachedRunEnd returns null for isRegenerate), so the follow-up stays queued for a manual send; covered in queue.spec and useResumableSSE.spec.

} 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 }),
Comment thread
berry-13 marked this conversation as resolved.
...(outcome === 'completed' && { responseMessageId: response.messageId }),
Comment on lines +291 to +296

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve the detached generation epoch

When an Agent server-queued turn is admitted while the predecessor chat is detached, its receipt records the consumed predecessor in effectivePredecessorCreatedAt, but this synthesized end omits generationCreatedAt. useQueueDrain therefore cannot match and consume that admission receipt; once receipt reconciliation removes the admitted server row, any remaining local row is treated as eligible on this predecessor end and can be submitted while the server-started successor already owns the conversation. Capture the detached run's active generation epoch and include it in the reconstructed RunEnd so the existing admission-boundary guard remains effective.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in d8269cc: the pane captures the run's live generation epoch when it leaves and the parked end carries it as generationCreatedAt, so the drain still matches admission receipts; a detached run with no epoch whose queue the server shares is left for a manual send (useResumeOnLoad.spec covers both).

};
}

const clearFamily = <Param>(family: {
getParams(): Iterable<Param>;
remove(param: Param): void;
Expand All @@ -244,4 +317,6 @@ export function resetQueueFamilies(): void {
clearFamily(drainAfterAbortByIndex);
clearFamily(runEndsByIndex);
clearFamily(runEndByIndex);
clearFamily(detachedRunByConvoId);
clearFamily(stopRequestedByConvoId);
}
4 changes: 3 additions & 1 deletion client/src/hooks/Chat/useChatHelpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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);
Comment thread
berry-13 marked this conversation as resolved.
try {
console.log('[useChatHelpers] Calling abort mutation for:', conversationId);
const response = await abortStream({
Expand Down Expand Up @@ -352,6 +353,7 @@ export default function useChatHelpers(index = 0, paramId?: string): ChatContrac
conversationId,
endpoint,
endpointType,
queueStore,
activeGenerationCreatedAt,
activeGenerationProtocolVersion,
abortStream,
Expand Down
97 changes: 96 additions & 1 deletion client/src/hooks/SSE/__tests__/useResumableSSE.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -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 },
Expand Down
Loading
Loading