-
-
Notifications
You must be signed in to change notification settings - Fork 9.3k
📬 fix: Drain a Follow-Up Queued in a Chat the User Left #16632
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: dev
Are you sure you want to change the base?
Changes from all commits
fa338dc
23c416b
bf95c79
576cecd
6c596f1
4062008
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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(); | ||
| }); | ||
| }); |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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); | ||
| } 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 }), | ||
|
berry-13 marked this conversation as resolved.
|
||
| ...(outcome === 'completed' && { responseMessageId: response.messageId }), | ||
|
Comment on lines
+291
to
+296
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
When an Agent server-queued turn is admitted while the predecessor chat is detached, its receipt records the consumed predecessor in Useful? React with 👍 / 👎.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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; | ||
|
|
@@ -244,4 +317,6 @@ export function resetQueueFamilies(): void { | |
| clearFamily(drainAfterAbortByIndex); | ||
| clearFamily(runEndsByIndex); | ||
| clearFamily(runEndByIndex); | ||
| clearFamily(detachedRunByConvoId); | ||
| clearFamily(stopRequestedByConvoId); | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
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 useexistingId_while the pre-runexistingIdrow already exists, meaningparkDetachedEnddoes 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 👍 / 👎.
There was a problem hiding this comment.
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.