From 91f6d82d668bdcfbd3f25e0ec7af52c819edb5bf Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:07:57 +0200 Subject: [PATCH 01/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Finalize=20Compacti?= =?UTF-8?q?on=20Snapshots=20and=20Move=20Abort=20Policy=20Into=20packages/?= =?UTF-8?q?api?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/controllers/agents/request.js | 39 ++++++++++ api/server/routes/agents/index.js | 6 +- packages/api/src/agents/compaction.spec.ts | 83 ++++++++++++++++++++++ packages/api/src/agents/compaction.ts | 46 ++++++++++++ 4 files changed, 170 insertions(+), 4 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index b75182e19ba..90f13089e64 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -63,6 +63,7 @@ const { stampPreliminaryPrivateTextMessage, announceReply, announceErrorTurn, + resolveFinalizedCompactionTurn, markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); @@ -463,6 +464,11 @@ async function saveErrorTurn( '_id', ); if (existing.length > 0) { + await finalizeFailedCompactionTurn(req, { + userId, + conversationId, + messageId: errorMessageId, + }); return; } if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { @@ -471,6 +477,11 @@ async function saveErrorTurn( '_id', ); if (partial.length > 0) { + await finalizeFailedCompactionTurn(req, { + userId, + conversationId, + messageId: liveResponseMessageId, + }); return; } } @@ -596,6 +607,34 @@ async function saveErrorTurn( } } +/** + * The disconnect save is marker-only while the run is still live; a failed + * turn is what settles it, so a compaction's partial row is finalized here + * with the terminal outcome instead of keeping the snapshot's live-run + * marking. The decision lives in @librechat/api; this is the wiring. + */ +async function finalizeFailedCompactionTurn(req, { userId, conversationId, messageId }) { + const [partialRow] = await getMessages({ user: userId, messageId, conversationId }); + const finalized = resolveFinalizedCompactionTurn(partialRow, req.body); + if (finalized == null) { + return; + } + await saveMessage( + { + userId, + isTemporary: + req?._agentEventBindingRetention?.isTemporary ?? + req?.resolvedConversation?.isTemporary ?? + req?.body?.isTemporary, + expiredAt: + req?._agentEventBindingRetention?.expiredAt ?? req?.resolvedConversation?.expiredAt, + interfaceConfig: req?.config?.interfaceConfig, + }, + { messageId, conversationId, ...finalized }, + { context: 'api/server/controllers/agents/request.js - finalize failed compaction turn' }, + ); +} + function classifyScheduledFailure(error, aborted = false) { if (aborted || error?.code === 'SCHEDULE_NO_LONGER_ACTIVE') { return { status: 'interrupted', error: error?.message }; diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index b2579f3dc84..ef9f90d4caf 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -6,6 +6,7 @@ const { TERMINAL_PUBLICATION_RECONNECT_ERROR, hasPersistableAbortContent, announceStoppedReply, + shouldPersistAbortAnchor, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -789,10 +790,7 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * its parent and the preliminary-parent fence correctly rejects it. */ const shouldPersistAbortedTurn = hasPersistableAbortContent(content) || jobData?.createdEventEmitted === true; - /** A compaction's `userMessage` is the already-persisted leaf - * projected for identity only; upserting it would erase a user - * leaf's text or turn an assistant leaf into an empty user row. */ - const shouldPersistAnchor = jobData?.compact !== true; + const shouldPersistAnchor = shouldPersistAbortAnchor(jobData); if ( jobData?.userMessage?.messageId && diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 5e6c2b924a4..e4f88f9efce 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -19,6 +19,8 @@ import { markCompactionOutcome, resolveFailedTurnContent, resolveCheckpointMessage, + resolveFinalizedCompactionTurn, + shouldPersistAbortAnchor, restoreCompactionSemanticIndex, restoreCompactionSemanticIndexSnapshot, stripUnusableSummaryParts, @@ -432,6 +434,87 @@ describe('markAbortedCompactionContent', () => { }); }); +describe('shouldPersistAbortAnchor', () => { + it('keeps the prerequisite user write for an ordinary turn', () => { + expect(shouldPersistAbortAnchor({})).toBe(true); + expect(shouldPersistAbortAnchor(null)).toBe(true); + }); + + /** The compaction's `userMessage` is the persisted leaf projected for + * identity only; upserting it would erase the leaf. */ + it('skips the prerequisite write for a compaction', () => { + expect(shouldPersistAbortAnchor({ compact: true })).toBe(false); + }); +}); + +describe('resolveFinalizedCompactionTurn', () => { + const compactionFailed = JSON.stringify({ type: ErrorTypes.COMPACTION_FAILED }); + + it('leaves a partial row of a turn that was not a compaction alone', () => { + const row = { content: [{ type: ContentTypes.TEXT, text: 'Partial answer' }] }; + + expect(resolveFinalizedCompactionTurn(row, {})).toBeNull(); + }); + + /** The disconnect snapshot is marker-only, so a run that fails afterwards + * leaves a row with no summary or error part and no marker at all. */ + it('finalizes a snapshot without a summary or error part with the typed failure', () => { + const row = { content: [{ type: ContentTypes.THINK, think: 'Picking what to summarize' }] }; + + expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + content: [ + { type: ContentTypes.THINK, think: 'Picking what to summarize' }, + { type: ContentTypes.ERROR, error: compactionFailed, initiatedBy: 'user' }, + ], + }); + }); + + it('marks a partial summary failed beside its text', () => { + const row = { + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + initiatedBy: 'user', + }, + ], + }; + + expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + initiatedBy: 'user', + failed: true, + }, + ], + }); + }); + + it('leaves a row that already carries a terminal outcome alone', () => { + const failedSummary = { + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + failed: true, + initiatedBy: 'user', + }, + ], + }; + const recordedFailure = { + content: [{ type: ContentTypes.ERROR, error: 'Summarization failed', initiatedBy: 'user' }], + }; + + expect(resolveFinalizedCompactionTurn(failedSummary, { compact: true })).toBeNull(); + expect(resolveFinalizedCompactionTurn(recordedFailure, { compact: true })).toBeNull(); + }); +}); + describe('findCheckpointSummaryPart', () => { const legacySummary = { type: ContentTypes.SUMMARY, text: 'Summary of conversation' }; diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 063870965ac..c1070d7323b 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -263,6 +263,52 @@ export function markAbortedCompactionContent( return contentParts; } +/** + * Whether the abort route's prerequisite user write applies to a job: a + * compaction's `userMessage` is the already-persisted branch leaf projected + * for identity only, so upserting it would erase a user leaf's text or turn + * an assistant leaf into an empty user row. The anchor exists by definition; + * only the aborted response needs writing. + */ +export function shouldPersistAbortAnchor( + jobData: { compact?: boolean } | null | undefined, +): boolean { + return jobData?.compact !== true; +} + +/** + * The content a failed compaction finalizes its already-persisted partial row + * with. The disconnect save is marker-only because the run is still live when + * it fires, so when the run then fails that snapshot is the row that stays, + * and it must carry the terminal outcome a compaction that failed without a + * snapshot records: a partial summary is marked failed beside its text, and a + * snapshot with no summary or error part gets the typed failure. Null when + * the row is not a compaction's or already carries a terminal outcome. + */ +export function resolveFinalizedCompactionTurn( + partialRow: { content?: unknown } | null | undefined, + requestBody: { compact?: boolean } | null | undefined, +): { content: TMessageContentParts[] } | null { + if (requestBody?.compact !== true) { + return null; + } + const content = Array.isArray(partialRow?.content) + ? (partialRow.content as TMessageContentParts[]) + : []; + for (const part of content) { + if (part?.type === ContentTypes.ERROR) { + return null; + } + if ( + part?.type === ContentTypes.SUMMARY && + (part.failed === true || isUsableSummaryPart(part)) + ) { + return null; + } + } + return { content: markAbortedCompactionContent(content, true) }; +} + /** * Stamps `initiatedBy: 'user'` on the part that carries a manual compaction's * outcome, which is the turn's only record of having been one: the run emits no From e75e70887ba6b905fb0ab80db4cf2dc7833b7a49 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:29:44 +0200 Subject: [PATCH 02/19] =?UTF-8?q?=F0=9F=94=95=20fix:=20Inspect=20Every=20P?= =?UTF-8?q?art=20When=20Finalizing=20and=20Reuse=20the=20Loaded=20Row?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/controllers/agents/request.js | 31 +++++++++++------- packages/api/src/agents/compaction.spec.ts | 38 ++++++++++++++++++++++ packages/api/src/agents/compaction.ts | 23 ++++++++----- 3 files changed, 73 insertions(+), 19 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 90f13089e64..45dec8da25b 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -459,28 +459,34 @@ async function saveErrorTurn( } const userId = req.user.id; - const existing = await getMessages( - { user: userId, messageId: errorMessageId, conversationId }, - '_id', - ); + /** Full documents: the compaction finalization below reuses the loaded + * row rather than reading it a second time. */ + const existing = await getMessages({ + user: userId, + messageId: errorMessageId, + conversationId, + }); if (existing.length > 0) { await finalizeFailedCompactionTurn(req, { userId, conversationId, messageId: errorMessageId, + partialRow: existing[0], }); return; } if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { - const partial = await getMessages( - { user: userId, messageId: liveResponseMessageId, conversationId }, - '_id', - ); + const partial = await getMessages({ + user: userId, + messageId: liveResponseMessageId, + conversationId, + }); if (partial.length > 0) { await finalizeFailedCompactionTurn(req, { userId, conversationId, messageId: liveResponseMessageId, + partialRow: partial[0], }); return; } @@ -611,10 +617,13 @@ async function saveErrorTurn( * The disconnect save is marker-only while the run is still live; a failed * turn is what settles it, so a compaction's partial row is finalized here * with the terminal outcome instead of keeping the snapshot's live-run - * marking. The decision lives in @librechat/api; this is the wiring. + * marking. The decision lives in @librechat/api; this is the wiring, reusing + * the row the caller already loaded. */ -async function finalizeFailedCompactionTurn(req, { userId, conversationId, messageId }) { - const [partialRow] = await getMessages({ user: userId, messageId, conversationId }); +async function finalizeFailedCompactionTurn( + req, + { userId, conversationId, messageId, partialRow }, +) { const finalized = resolveFinalizedCompactionTurn(partialRow, req.body); if (finalized == null) { return; diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index e4f88f9efce..4f28e8dea96 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -513,6 +513,44 @@ describe('resolveFinalizedCompactionTurn', () => { expect(resolveFinalizedCompactionTurn(failedSummary, { compact: true })).toBeNull(); expect(resolveFinalizedCompactionTurn(recordedFailure, { compact: true })).toBeNull(); }); + + /** A row can hold an earlier round's terminal outcome beside a later + * unfinished summary: only the inspection of every part catches it, and + * the failure lands on the summary that never finished. */ + it('finalizes a later unfinished summary beside an earlier terminal outcome', () => { + const row = { + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'An earlier checkpoint.' }], + boundary: completedBoundary, + }, + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'A later partial round.' }], + summarizing: true, + }, + ], + }; + + expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'An earlier checkpoint.' }], + boundary: completedBoundary, + initiatedBy: 'user', + }, + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'A later partial round.' }], + summarizing: true, + initiatedBy: 'user', + failed: true, + }, + ], + }); + }); }); describe('findCheckpointSummaryPart', () => { diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index c1070d7323b..941ab973c63 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -295,17 +295,24 @@ export function resolveFinalizedCompactionTurn( const content = Array.isArray(partialRow?.content) ? (partialRow.content as TMessageContentParts[]) : []; + /** Every part is inspected: a row can hold an earlier round's terminal + * outcome beside a later unfinished summary, and that summary still needs + * its failure marked. */ + let sawSummaryOrError = false; + let unfinishedSummary = false; for (const part of content) { - if (part?.type === ContentTypes.ERROR) { - return null; - } - if ( - part?.type === ContentTypes.SUMMARY && - (part.failed === true || isUsableSummaryPart(part)) - ) { - return null; + if (part?.type === ContentTypes.SUMMARY) { + sawSummaryOrError = true; + if (part.failed !== true && !isUsableSummaryPart(part)) { + unfinishedSummary = true; + } + } else if (part?.type === ContentTypes.ERROR) { + sawSummaryOrError = true; } } + if (sawSummaryOrError && !unfinishedSummary) { + return null; + } return { content: markAbortedCompactionContent(content, true) }; } From 62257fead5cb9c5dd1c45daddbbb14fbce7a5f28 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 14:27:12 +0200 Subject: [PATCH 03/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Settle=20Finalized?= =?UTF-8?q?=20Compaction=20Rows=20and=20Verify=20Their=20Anchor?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/controllers/agents/request.js | 46 +++++++----- .../routes/agents/__tests__/abort.spec.js | 53 ++++++++++++++ api/server/routes/agents/index.js | 23 +++++- packages/api/src/agents/compaction.spec.ts | 73 +++++++++++++++++-- packages/api/src/agents/compaction.ts | 61 ++++++++++++++-- 5 files changed, 219 insertions(+), 37 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 45dec8da25b..e37db4f4f04 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -63,7 +63,7 @@ const { stampPreliminaryPrivateTextMessage, announceReply, announceErrorTurn, - resolveFinalizedCompactionTurn, + persistFinalizedCompactionTurn, markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); @@ -620,28 +620,36 @@ async function saveErrorTurn( * marking. The decision lives in @librechat/api; this is the wiring, reusing * the row the caller already loaded. */ +/** + * The disconnect save is marker-only while the run is still live; a failed + * turn is what settles it, so a compaction's partial row is finalized with the + * terminal outcome instead of keeping the snapshot's live-run marking. The + * operation lives in @librechat/api; this wiring supplies the caller's + * persistence and the row saveErrorTurn already loaded. + */ async function finalizeFailedCompactionTurn( req, { userId, conversationId, messageId, partialRow }, ) { - const finalized = resolveFinalizedCompactionTurn(partialRow, req.body); - if (finalized == null) { - return; - } - await saveMessage( - { - userId, - isTemporary: - req?._agentEventBindingRetention?.isTemporary ?? - req?.resolvedConversation?.isTemporary ?? - req?.body?.isTemporary, - expiredAt: - req?._agentEventBindingRetention?.expiredAt ?? req?.resolvedConversation?.expiredAt, - interfaceConfig: req?.config?.interfaceConfig, - }, - { messageId, conversationId, ...finalized }, - { context: 'api/server/controllers/agents/request.js - finalize failed compaction turn' }, - ); + await persistFinalizedCompactionTurn(partialRow, req.body, { + messageId, + conversationId, + saveMessage: (message) => + saveMessage( + { + userId, + isTemporary: + req?._agentEventBindingRetention?.isTemporary ?? + req?.resolvedConversation?.isTemporary ?? + req?.body?.isTemporary, + expiredAt: + req?._agentEventBindingRetention?.expiredAt ?? req?.resolvedConversation?.expiredAt, + interfaceConfig: req?.config?.interfaceConfig, + }, + message, + { context: 'api/server/controllers/agents/request.js - finalize failed compaction turn' }, + ), + }); } function classifyScheduledFailure(error, aborted = false) { diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index 1656bb0a447..15655e3a748 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -27,6 +27,7 @@ const mockSaveMessage = jest.fn(); const mockHasPersistedPrivateText = jest.fn(); const mockGetPrivateMessageTexts = jest.fn(); const mockSaveConvo = jest.fn(); +const mockGetMessages = jest.fn(async () => [{ _id: 'existing-anchor' }]); const mockRecordScheduleOutcome = jest.fn(); const mockBeginScheduledStop = jest.fn(); @@ -55,6 +56,7 @@ jest.mock('~/models', () => ({ (await mockHasPersistedPrivateText(...args)) ? 'protected-row-id' : null, getPrivateMessageTexts: (...args) => mockGetPrivateMessageTexts(...args), saveConvo: (...args) => mockSaveConvo(...args), + getMessages: (...args) => mockGetMessages(...args), })); jest.mock('~/server/services/Schedules', () => ({ @@ -443,6 +445,57 @@ describe('Agent Abort Endpoint', () => { expect.objectContaining({ context: expect.stringContaining('abort endpoint') }), ); }); + + /** Stop can win the race before the branch loaded, so the projected + * anchor names a row that was never written: a response persisted + * there would be orphaned on reload. */ + it('persists nothing when the compaction anchor was never written', async () => { + mockGetMessages.mockResolvedValueOnce([]); + const jobStreamId = 'test-stream-compact-unanchored'; + + mockGenerationJobManager.getJob.mockResolvedValue({ + metadata: { userId: 'test-user-123', generationProtocolVersion: 2 }, + }); + + const abortResult = { + success: true, + jobData: { + compact: true, + createdEventEmitted: true, + userMessage: { + messageId: 'never-persisted-leaf', + parentMessageId: 'older-response', + conversationId: jobStreamId, + text: '', + }, + responseMessageId: 'compaction-response-2', + conversationId: jobStreamId, + endpoint: 'agents', + sender: 'TestAgent', + model: 'agent-1', + }, + content: [ + { + type: 'error', + error: JSON.stringify({ type: 'compaction_failed' }), + initiatedBy: 'user', + }, + ], + text: '', + }; + mockGenerationJobManager.abortJob.mockImplementation(async (_streamId, options) => { + await options.beforePublish(abortResult); + return abortResult; + }); + + const response = await request(app) + .post('/api/agents/chat/abort') + .set('X-LibreChat-Generation-Protocol', '2') + .send({ conversationId: jobStreamId, generationProtocolVersion: 2 }); + + expect(response.status).toBe(200); + expect(mockSaveMessage).not.toHaveBeenCalled(); + }); }); describe('Partial Response Saving', () => { diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index ef9f90d4caf..b3f0a9f4cba 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -6,7 +6,7 @@ const { TERMINAL_PUBLICATION_RECONNECT_ERROR, hasPersistableAbortContent, announceStoppedReply, - shouldPersistAbortAnchor, + resolveAbortAnchorDecision, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -57,6 +57,7 @@ const { } = require('~/server/controllers/agents/protocol'); const { getFiles, + getMessages, saveMessage, saveConvo, getPersistedPrivateTextId, @@ -790,12 +791,26 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * its parent and the preliminary-parent fence correctly rejects it. */ const shouldPersistAbortedTurn = hasPersistableAbortContent(content) || jobData?.createdEventEmitted === true; - const shouldPersistAnchor = shouldPersistAbortAnchor(jobData); + /** A compaction anchors on the persisted branch leaf, so its + * existence decides the prerequisite write: present, the projected + * anchor is never upserted over it; absent (Stop won the race + * before the branch loaded), there is nothing to parent the aborted + * response onto and nothing is written. */ + let anchorDecision = resolveAbortAnchorDecision(jobData, true); + if (jobData?.compact === true && jobData?.userMessage?.messageId && req?.user?.id) { + const anchorRows = await getMessages({ + user: req.user.id, + messageId: jobData.userMessage.messageId, + conversationId: jobData.conversationId, + }); + anchorDecision = resolveAbortAnchorDecision(jobData, anchorRows.length > 0); + } if ( jobData?.userMessage?.messageId && jobData?.responseMessageId && - shouldPersistAbortedTurn + shouldPersistAbortedTurn && + anchorDecision !== 'skip-turn' ) { const messageContext = { userId: req?.user?.id, @@ -857,7 +872,7 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * operation gets a chance to succeed. A compaction skips the * prerequisite: its anchor is the persisted leaf itself. */ let persistedRequestId; - if (shouldPersistAnchor) { + if (anchorDecision === 'persist') { try { const persistedRequest = await saveAbortedUserMessage( { saveMessage, getPersistedPrivateTextId, getPrivateMessageTexts }, diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 4f28e8dea96..517e1308556 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -17,10 +17,11 @@ import { getSummaryPartText, markAbortedCompactionContent, markCompactionOutcome, + persistFinalizedCompactionTurn, + resolveAbortAnchorDecision, resolveFailedTurnContent, resolveCheckpointMessage, resolveFinalizedCompactionTurn, - shouldPersistAbortAnchor, restoreCompactionSemanticIndex, restoreCompactionSemanticIndexSnapshot, stripUnusableSummaryParts, @@ -434,16 +435,76 @@ describe('markAbortedCompactionContent', () => { }); }); -describe('shouldPersistAbortAnchor', () => { +describe('resolveAbortAnchorDecision', () => { it('keeps the prerequisite user write for an ordinary turn', () => { - expect(shouldPersistAbortAnchor({})).toBe(true); - expect(shouldPersistAbortAnchor(null)).toBe(true); + expect(resolveAbortAnchorDecision({}, true)).toBe('persist'); + expect(resolveAbortAnchorDecision(null, false)).toBe('persist'); }); /** The compaction's `userMessage` is the persisted leaf projected for * identity only; upserting it would erase the leaf. */ - it('skips the prerequisite write for a compaction', () => { - expect(shouldPersistAbortAnchor({ compact: true })).toBe(false); + it('skips the prerequisite write for a compaction anchored on a persisted leaf', () => { + expect(resolveAbortAnchorDecision({ compact: true }, true)).toBe('skip-anchor'); + }); + + /** Stop can win the race before the branch loaded, leaving the projection + * with no row behind it: a response written there would be orphaned. */ + it('skips the whole turn when the compaction anchor was never persisted', () => { + expect(resolveAbortAnchorDecision({ compact: true }, false)).toBe('skip-turn'); + }); +}); + +describe('persistFinalizedCompactionTurn', () => { + it('writes the finalized content with the terminal envelope', async () => { + const saved: Record[] = []; + const partialRow = { + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + }, + ], + }; + + await persistFinalizedCompactionTurn( + partialRow, + { compact: true }, + { + messageId: 'response-1', + conversationId: 'conversation-1', + saveMessage: async (message) => { + saved.push(message); + return message; + }, + }, + ); + + /** The snapshot was saved `unfinished` with no error while the run was + * live; the settled row must not keep reading as an incomplete + * response. */ + expect(saved).toHaveLength(1); + expect(saved[0]).toMatchObject({ + messageId: 'response-1', + conversationId: 'conversation-1', + unfinished: false, + error: true, + }); + expect(saved[0].content).toEqual([ + expect.objectContaining({ type: ContentTypes.SUMMARY, failed: true, initiatedBy: 'user' }), + ]); + }); + + it('writes nothing when the row needs no finalization', async () => { + const saveMessage = jest.fn(); + + await persistFinalizedCompactionTurn( + { content: [{ type: ContentTypes.TEXT, text: 'An ordinary partial' }] }, + {}, + { messageId: 'response-1', conversationId: 'conversation-1', saveMessage }, + ); + + expect(saveMessage).not.toHaveBeenCalled(); }); }); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 941ab973c63..46eba54a8d0 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -263,17 +263,62 @@ export function markAbortedCompactionContent( return contentParts; } +/** How the abort route persists a stopped turn's prerequisite rows. */ +export type AbortAnchorDecision = 'persist' | 'skip-anchor' | 'skip-turn'; + /** - * Whether the abort route's prerequisite user write applies to a job: a - * compaction's `userMessage` is the already-persisted branch leaf projected - * for identity only, so upserting it would erase a user leaf's text or turn - * an assistant leaf into an empty user row. The anchor exists by definition; - * only the aborted response needs writing. + * A compaction's `userMessage` is the branch leaf projected for identity + * only: when the leaf is already persisted, the projection must never be + * upserted over it (an ordinary prerequisite write would erase a user leaf's + * text or turn an assistant leaf into an empty user row), so only the aborted + * response is written. When the leaf is NOT persisted, Stop won the race + * before the branch loaded and there is nothing to anchor the response onto, + * so nothing is written at all. Ordinary turns keep the prerequisite write. */ -export function shouldPersistAbortAnchor( +export function resolveAbortAnchorDecision( jobData: { compact?: boolean } | null | undefined, -): boolean { - return jobData?.compact !== true; + anchorExists: boolean, +): AbortAnchorDecision { + if (jobData?.compact !== true) { + return 'persist'; + } + return anchorExists ? 'skip-anchor' : 'skip-turn'; +} + +/** + * Finalizes a failed compaction's already-persisted partial row, with the + * write injected so the operation runs against whatever persistence the + * caller owns. The row settles with the terminal envelope the error path + * writes (an errored, finished turn): with the snapshot's `unfinished` flag + * left in place, restored sessions and downstream readers would keep + * classifying the failed turn as an incomplete response. Returns whether a + * write happened. + */ +export async function persistFinalizedCompactionTurn( + partialRow: { content?: unknown } | null | undefined, + requestBody: { compact?: boolean } | null | undefined, + { + messageId, + conversationId, + saveMessage, + }: { + messageId: string; + conversationId: string; + saveMessage: (message: Record) => Promise; + }, +): Promise { + const finalized = resolveFinalizedCompactionTurn(partialRow, requestBody); + if (finalized == null) { + return false; + } + await saveMessage({ + messageId, + conversationId, + unfinished: false, + error: true, + ...finalized, + }); + return true; } /** From 3d33a756235f2a196e965233c5d51673e2a3c9e1 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 14:49:10 +0200 Subject: [PATCH 04/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Withhold=20the=20Un?= =?UTF-8?q?anchored=20Abort=20Final=20and=20Settle=20Marked=20Snapshots?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../routes/agents/__tests__/abort.spec.js | 11 +++- api/server/routes/agents/index.js | 9 +++ packages/api/src/agents/compaction.spec.ts | 61 +++++++++++++++++-- packages/api/src/agents/compaction.ts | 57 +++++++++++------ 4 files changed, 115 insertions(+), 23 deletions(-) diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index 15655e3a748..d823a76305b 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -483,8 +483,15 @@ describe('Agent Abort Endpoint', () => { ], text: '', }; + let beforePublishError; mockGenerationJobManager.abortJob.mockImplementation(async (_streamId, options) => { - await options.beforePublish(abortResult); + try { + await options.beforePublish(abortResult); + } catch (error) { + /** The manager catches this failure and publishes a + * reconciliation frame instead of the normal FINAL. */ + beforePublishError = error; + } return abortResult; }); @@ -495,6 +502,8 @@ describe('Agent Abort Endpoint', () => { expect(response.status).toBe(200); expect(mockSaveMessage).not.toHaveBeenCalled(); + expect(beforePublishError).toBeInstanceOf(Error); + expect(beforePublishError.message).toContain('anchor was never persisted'); }); }); diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index b3f0a9f4cba..53e065e9c08 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -805,6 +805,15 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { }); anchorDecision = resolveAbortAnchorDecision(jobData, anchorRows.length > 0); } + if (anchorDecision === 'skip-turn') { + /** Throwing here is the contract for "do not publish the normal + * FINAL": the manager emits a reconciliation frame instead of + * one whose response points at a row deliberately never + * persisted. */ + persistenceErrors.push( + new Error('Compaction anchor was never persisted; abort turn withheld'), + ); + } if ( jobData?.userMessage?.messageId && diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 517e1308556..32bf57f2a5a 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -495,6 +495,34 @@ describe('persistFinalizedCompactionTurn', () => { ]); }); + it('writes only the envelope when the parts already carry the failure', async () => { + const saved: Record[] = []; + const partialRow = { + content: [{ type: ContentTypes.ERROR, error: 'Summarization failed', initiatedBy: 'user' }], + }; + + await persistFinalizedCompactionTurn( + partialRow, + { compact: true }, + { + messageId: 'response-1', + conversationId: 'conversation-1', + saveMessage: async (message) => { + saved.push(message); + return message; + }, + }, + ); + + expect(saved).toHaveLength(1); + expect(saved[0]).toEqual({ + messageId: 'response-1', + conversationId: 'conversation-1', + unfinished: false, + error: true, + }); + }); + it('writes nothing when the row needs no finalization', async () => { const saveMessage = jest.fn(); @@ -514,7 +542,7 @@ describe('resolveFinalizedCompactionTurn', () => { it('leaves a partial row of a turn that was not a compaction alone', () => { const row = { content: [{ type: ContentTypes.TEXT, text: 'Partial answer' }] }; - expect(resolveFinalizedCompactionTurn(row, {})).toBeNull(); + expect(resolveFinalizedCompactionTurn(row, {})).toEqual({ write: false }); }); /** The disconnect snapshot is marker-only, so a run that fails afterwards @@ -523,6 +551,7 @@ describe('resolveFinalizedCompactionTurn', () => { const row = { content: [{ type: ContentTypes.THINK, think: 'Picking what to summarize' }] }; expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + write: true, content: [ { type: ContentTypes.THINK, think: 'Picking what to summarize' }, { type: ContentTypes.ERROR, error: compactionFailed, initiatedBy: 'user' }, @@ -543,6 +572,7 @@ describe('resolveFinalizedCompactionTurn', () => { }; expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + write: true, content: [ { type: ContentTypes.SUMMARY, @@ -555,7 +585,10 @@ describe('resolveFinalizedCompactionTurn', () => { }); }); - it('leaves a row that already carries a terminal outcome alone', () => { + /** The parts already carry the failure, but the snapshot's live-run flags + * are still unsettled: the write settles the envelope without touching + * content. */ + it('settles only the envelope of a row whose parts already carry the failure', () => { const failedSummary = { content: [ { @@ -571,8 +604,27 @@ describe('resolveFinalizedCompactionTurn', () => { content: [{ type: ContentTypes.ERROR, error: 'Summarization failed', initiatedBy: 'user' }], }; - expect(resolveFinalizedCompactionTurn(failedSummary, { compact: true })).toBeNull(); - expect(resolveFinalizedCompactionTurn(recordedFailure, { compact: true })).toBeNull(); + expect(resolveFinalizedCompactionTurn(failedSummary, { compact: true })).toEqual({ + write: true, + }); + expect(resolveFinalizedCompactionTurn(recordedFailure, { compact: true })).toEqual({ + write: true, + }); + }); + + /** A checkpoint the run completed before failing stands exactly as it is. */ + it('leaves a completed checkpoint untouched', () => { + const row = { + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'A finished checkpoint.' }], + boundary: completedBoundary, + }, + ], + }; + + expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ write: false }); }); /** A row can hold an earlier round's terminal outcome beside a later @@ -595,6 +647,7 @@ describe('resolveFinalizedCompactionTurn', () => { }; expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ + write: true, content: [ { type: ContentTypes.SUMMARY, diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 46eba54a8d0..07708c63e91 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -308,7 +308,7 @@ export async function persistFinalizedCompactionTurn( }, ): Promise { const finalized = resolveFinalizedCompactionTurn(partialRow, requestBody); - if (finalized == null) { + if (!finalized.write) { return false; } await saveMessage({ @@ -316,26 +316,37 @@ export async function persistFinalizedCompactionTurn( conversationId, unfinished: false, error: true, - ...finalized, + ...(finalized.content != null && { content: finalized.content }), }); return true; } +/** What a failed compaction does with its already-persisted partial row. */ +export type FinalizedCompactionTurn = + /** Not the failed run's row, or one holding nothing but a completed + * checkpoint worth keeping exactly as it stands. */ + | { write: false } + /** The parts already carry the failure (an error part, a failed summary); + * only the snapshot's live-run flags remain to settle. */ + | { write: true; content?: undefined } + /** The parts need the terminal marking applied. */ + | { write: true; content: TMessageContentParts[] }; + /** - * The content a failed compaction finalizes its already-persisted partial row - * with. The disconnect save is marker-only because the run is still live when - * it fires, so when the run then fails that snapshot is the row that stays, - * and it must carry the terminal outcome a compaction that failed without a - * snapshot records: a partial summary is marked failed beside its text, and a - * snapshot with no summary or error part gets the typed failure. Null when - * the row is not a compaction's or already carries a terminal outcome. + * The disconnect save is marker-only because the run is still live when it + * fires, so when the run then fails that snapshot is the row that stays: a + * partial summary is marked failed beside its text, a snapshot with no + * summary or error part gets the typed failure, and a snapshot whose parts + * already carry the failure still settles its live-run flags. Only a + * completed checkpoint (the run produced its summary before failing) is left + * untouched, and rows of turns that were not compactions are never written. */ export function resolveFinalizedCompactionTurn( partialRow: { content?: unknown } | null | undefined, requestBody: { compact?: boolean } | null | undefined, -): { content: TMessageContentParts[] } | null { +): FinalizedCompactionTurn { if (requestBody?.compact !== true) { - return null; + return { write: false }; } const content = Array.isArray(partialRow?.content) ? (partialRow.content as TMessageContentParts[]) @@ -343,22 +354,32 @@ export function resolveFinalizedCompactionTurn( /** Every part is inspected: a row can hold an earlier round's terminal * outcome beside a later unfinished summary, and that summary still needs * its failure marked. */ - let sawSummaryOrError = false; + let sawFailure = false; + let sawCheckpoint = false; let unfinishedSummary = false; for (const part of content) { if (part?.type === ContentTypes.SUMMARY) { - sawSummaryOrError = true; - if (part.failed !== true && !isUsableSummaryPart(part)) { + if (part.failed === true) { + sawFailure = true; + } else if (isUsableSummaryPart(part)) { + sawCheckpoint = true; + } else { unfinishedSummary = true; } } else if (part?.type === ContentTypes.ERROR) { - sawSummaryOrError = true; + sawFailure = true; } } - if (sawSummaryOrError && !unfinishedSummary) { - return null; + if (unfinishedSummary) { + return { write: true, content: markAbortedCompactionContent(content, true) }; + } + if (sawFailure) { + return { write: true }; + } + if (sawCheckpoint) { + return { write: false }; } - return { content: markAbortedCompactionContent(content, true) }; + return { write: true, content: markAbortedCompactionContent(content, true) }; } /** From b1ad4fbc68f58f5c8d0ffd82b2bf8ce01de91330 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 15:25:45 +0200 Subject: [PATCH 05/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Consolidate=20the?= =?UTF-8?q?=20Compaction=20Abort=20Policy=20Into=20One=20Injected=20Operat?= =?UTF-8?q?ion?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/controllers/agents/request.js | 32 +++---- .../routes/agents/__tests__/abort.spec.js | 2 +- api/server/routes/agents/index.js | 24 +++-- packages/api/src/agents/compaction.spec.ts | 94 ++++++++++++++++--- packages/api/src/agents/compaction.ts | 63 +++++++++---- 5 files changed, 153 insertions(+), 62 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index e37db4f4f04..e54e51940ae 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -459,28 +459,24 @@ async function saveErrorTurn( } const userId = req.user.id; - /** Full documents: the compaction finalization below reuses the loaded - * row rather than reading it a second time. */ - const existing = await getMessages({ - user: userId, - messageId: errorMessageId, - conversationId, - }); + const existing = await getMessages( + { user: userId, messageId: errorMessageId, conversationId }, + '_id', + ); if (existing.length > 0) { - await finalizeFailedCompactionTurn(req, { - userId, - conversationId, - messageId: errorMessageId, - partialRow: existing[0], - }); + /** No compaction finalization here: this id can normalize back to the + * anchor itself when the anchor ends in `_`, and the anchor is never + * the failed run's row. The run's own snapshot, if any, is checked + * under its distinct live response id below. */ return; } if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { - const partial = await getMessages({ - user: userId, - messageId: liveResponseMessageId, - conversationId, - }); + /** Full documents only where the compaction finalization needs the + * content; ordinary failures keep the id-only projection. */ + const partial = await getMessages( + { user: userId, messageId: liveResponseMessageId, conversationId }, + req.body?.compact === true ? undefined : '_id', + ); if (partial.length > 0) { await finalizeFailedCompactionTurn(req, { userId, diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index d823a76305b..64c72268d8d 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -503,7 +503,7 @@ describe('Agent Abort Endpoint', () => { expect(response.status).toBe(200); expect(mockSaveMessage).not.toHaveBeenCalled(); expect(beforePublishError).toBeInstanceOf(Error); - expect(beforePublishError.message).toContain('anchor was never persisted'); + expect(beforePublishError.message).toContain('anchor unavailable'); }); }); diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index 53e065e9c08..cf716a37134 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -6,7 +6,7 @@ const { TERMINAL_PUBLICATION_RECONNECT_ERROR, hasPersistableAbortContent, announceStoppedReply, - resolveAbortAnchorDecision, + resolveAbortedTurnAnchorDecision, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -796,23 +796,21 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * anchor is never upserted over it; absent (Stop won the race * before the branch loaded), there is nothing to parent the aborted * response onto and nothing is written. */ - let anchorDecision = resolveAbortAnchorDecision(jobData, true); - if (jobData?.compact === true && jobData?.userMessage?.messageId && req?.user?.id) { - const anchorRows = await getMessages({ - user: req.user.id, - messageId: jobData.userMessage.messageId, - conversationId: jobData.conversationId, - }); - anchorDecision = resolveAbortAnchorDecision(jobData, anchorRows.length > 0); - } + /** The compaction anchor policy (never upsert the projected leaf, + * never write a response with nothing to hang it on) lives in + * @librechat/api; this supplies the route's reader. */ + const anchorDecision = await resolveAbortedTurnAnchorDecision(jobData, { + messageExists: (messageId, conversationId) => + getMessages({ user: req?.user?.id, messageId, conversationId }).then( + (rows) => rows.length > 0, + ), + }); if (anchorDecision === 'skip-turn') { /** Throwing here is the contract for "do not publish the normal * FINAL": the manager emits a reconciliation frame instead of * one whose response points at a row deliberately never * persisted. */ - persistenceErrors.push( - new Error('Compaction anchor was never persisted; abort turn withheld'), - ); + persistenceErrors.push(new Error('Compaction anchor unavailable; abort turn withheld')); } if ( diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 32bf57f2a5a..da5afeca227 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -18,7 +18,7 @@ import { markAbortedCompactionContent, markCompactionOutcome, persistFinalizedCompactionTurn, - resolveAbortAnchorDecision, + resolveAbortedTurnAnchorDecision, resolveFailedTurnContent, resolveCheckpointMessage, resolveFinalizedCompactionTurn, @@ -435,22 +435,50 @@ describe('markAbortedCompactionContent', () => { }); }); -describe('resolveAbortAnchorDecision', () => { - it('keeps the prerequisite user write for an ordinary turn', () => { - expect(resolveAbortAnchorDecision({}, true)).toBe('persist'); - expect(resolveAbortAnchorDecision(null, false)).toBe('persist'); +describe('resolveAbortedTurnAnchorDecision', () => { + const reader = (exists: boolean) => jest.fn(async () => exists); + const jobData = { + compact: true, + conversationId: 'conversation-1', + userMessage: { messageId: 'leaf-1' }, + }; + + it('keeps the prerequisite user write for an ordinary turn', async () => { + const messageExists = reader(true); + + await expect(resolveAbortedTurnAnchorDecision({}, { messageExists })).resolves.toBe('persist'); + await expect(resolveAbortedTurnAnchorDecision(null, { messageExists })).resolves.toBe( + 'persist', + ); + expect(messageExists).not.toHaveBeenCalled(); }); /** The compaction's `userMessage` is the persisted leaf projected for * identity only; upserting it would erase the leaf. */ - it('skips the prerequisite write for a compaction anchored on a persisted leaf', () => { - expect(resolveAbortAnchorDecision({ compact: true }, true)).toBe('skip-anchor'); + it('skips the prerequisite write for a compaction anchored on a persisted leaf', async () => { + await expect( + resolveAbortedTurnAnchorDecision(jobData, { messageExists: reader(true) }), + ).resolves.toBe('skip-anchor'); }); /** Stop can win the race before the branch loaded, leaving the projection * with no row behind it: a response written there would be orphaned. */ - it('skips the whole turn when the compaction anchor was never persisted', () => { - expect(resolveAbortAnchorDecision({ compact: true }, false)).toBe('skip-turn'); + it('skips the whole turn when the compaction anchor was never persisted', async () => { + await expect( + resolveAbortedTurnAnchorDecision(jobData, { messageExists: reader(false) }), + ).resolves.toBe('skip-turn'); + }); + + /** A read that throws must not escape past the caller's remaining cleanup: + * nothing is known about the anchor, so nothing is written either. */ + it('skips the whole turn when the anchor read fails', async () => { + const messageExists = jest.fn(async () => { + throw new Error('mongo unavailable'); + }); + + await expect(resolveAbortedTurnAnchorDecision(jobData, { messageExists })).resolves.toBe( + 'skip-turn', + ); }); }); @@ -534,6 +562,32 @@ describe('persistFinalizedCompactionTurn', () => { expect(saveMessage).not.toHaveBeenCalled(); }); + + /** The surrounding failed-turn persistence treats a falsy save as a + * failure, not a settled row. */ + it('fails when the injected save resolves falsy', async () => { + const partialRow = { + content: [ + { + type: ContentTypes.ERROR, + error: 'Summarization failed', + initiatedBy: 'user', + }, + ], + }; + + await expect( + persistFinalizedCompactionTurn( + partialRow, + { compact: true }, + { + messageId: 'response-1', + conversationId: 'conversation-1', + saveMessage: async () => null, + }, + ), + ).rejects.toThrow('Failed compaction turn could not be finalized'); + }); }); describe('resolveFinalizedCompactionTurn', () => { @@ -612,9 +666,27 @@ describe('resolveFinalizedCompactionTurn', () => { }); }); - /** A checkpoint the run completed before failing stands exactly as it is. */ - it('leaves a completed checkpoint untouched', () => { + /** A checkpoint the run completed before failing is preserved as content, + * but a snapshot still flagged unfinished settles its envelope: the + * restored conversation must not keep treating the terminal job as live. */ + it('settles the envelope of an unfinished snapshot holding a completed checkpoint', () => { + const snapshot = { + unfinished: true, + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'A finished checkpoint.' }], + boundary: completedBoundary, + }, + ], + }; + + expect(resolveFinalizedCompactionTurn(snapshot, { compact: true })).toEqual({ write: true }); + }); + + it('leaves an already-settled checkpoint row untouched', () => { const row = { + unfinished: false, content: [ { type: ContentTypes.SUMMARY, diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 07708c63e91..66d762726c4 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -267,22 +267,40 @@ export function markAbortedCompactionContent( export type AbortAnchorDecision = 'persist' | 'skip-anchor' | 'skip-turn'; /** - * A compaction's `userMessage` is the branch leaf projected for identity - * only: when the leaf is already persisted, the projection must never be - * upserted over it (an ordinary prerequisite write would erase a user leaf's - * text or turn an assistant leaf into an empty user row), so only the aborted - * response is written. When the leaf is NOT persisted, Stop won the race - * before the branch loaded and there is nothing to anchor the response onto, - * so nothing is written at all. Ordinary turns keep the prerequisite write. + * Decides how a stopped turn's persistence treats its user row, reading the + * anchor through the caller's database reader. A compaction's `userMessage` + * is the branch leaf projected for identity only: when the leaf is persisted, + * the projection must never be upserted over it (an ordinary prerequisite + * write would erase a user leaf's text or turn an assistant leaf into an + * empty user row), so only the aborted response is written. When the leaf is + * NOT persisted, Stop won the race before the branch loaded and there is + * nothing to anchor the response onto, so nothing is written at all; a read + * that fails says the same thing, without throwing past the caller's + * remaining cleanup. Ordinary turns keep the prerequisite write. */ -export function resolveAbortAnchorDecision( - jobData: { compact?: boolean } | null | undefined, - anchorExists: boolean, -): AbortAnchorDecision { - if (jobData?.compact !== true) { +export async function resolveAbortedTurnAnchorDecision( + jobData: + | { + compact?: boolean; + conversationId?: string; + userMessage?: { messageId?: string } | null; + } + | null + | undefined, + { + messageExists, + }: { messageExists: (messageId: string, conversationId?: string) => Promise }, +): Promise { + const anchorId = jobData?.userMessage?.messageId; + if (jobData?.compact !== true || anchorId == null || anchorId.length === 0) { return 'persist'; } - return anchorExists ? 'skip-anchor' : 'skip-turn'; + try { + const anchorExists = await messageExists(anchorId, jobData.conversationId); + return anchorExists ? 'skip-anchor' : 'skip-turn'; + } catch { + return 'skip-turn'; + } } /** @@ -311,13 +329,18 @@ export async function persistFinalizedCompactionTurn( if (!finalized.write) { return false; } - await saveMessage({ + const saved = await saveMessage({ messageId, conversationId, unfinished: false, error: true, ...(finalized.content != null && { content: finalized.content }), }); + if (saved == null) { + /** The same contract the surrounding failed-turn persistence holds: a + * falsy save is a failure to settle, not a settled row. */ + throw new Error('Failed compaction turn could not be finalized'); + } return true; } @@ -337,12 +360,14 @@ export type FinalizedCompactionTurn = * fires, so when the run then fails that snapshot is the row that stays: a * partial summary is marked failed beside its text, a snapshot with no * summary or error part gets the typed failure, and a snapshot whose parts - * already carry the failure still settles its live-run flags. Only a - * completed checkpoint (the run produced its summary before failing) is left - * untouched, and rows of turns that were not compactions are never written. + * already carry the failure still settles its live-run flags. A completed + * checkpoint is preserved as content, but a snapshot still flagged + * `unfinished` settles its envelope even then, or the restored conversation + * keeps treating the terminal job as live; a row that was already settled is + * left alone. Rows of turns that were not compactions are never written. */ export function resolveFinalizedCompactionTurn( - partialRow: { content?: unknown } | null | undefined, + partialRow: { content?: unknown; unfinished?: boolean } | null | undefined, requestBody: { compact?: boolean } | null | undefined, ): FinalizedCompactionTurn { if (requestBody?.compact !== true) { @@ -377,7 +402,7 @@ export function resolveFinalizedCompactionTurn( return { write: true }; } if (sawCheckpoint) { - return { write: false }; + return partialRow?.unfinished === true ? { write: true } : { write: false }; } return { write: true, content: markAbortedCompactionContent(content, true) }; } From 8b3348c8011bef801f4a70b309852dcd9b040835 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 15:52:01 +0200 Subject: [PATCH 06/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Settle=20the=20Live?= =?UTF-8?q?=20Compaction=20Row=20Past=20an=20Anchor-Shaped=20Collision?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/controllers/agents/request.js | 53 +++++++++++++++--------- 1 file changed, 33 insertions(+), 20 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index e54e51940ae..67193e19e4a 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -459,33 +459,46 @@ async function saveErrorTurn( } const userId = req.user.id; + /** Whether the failed run's own partial row exists under its distinct + * live response id, finalizing it when it does (a no-op for turns that + * were not compactions). Full documents are requested only where the + * compaction finalization needs the content; ordinary failures keep the + * id-only projection. */ + const settleLiveSnapshot = async () => { + if (liveResponseMessageId == null || liveResponseMessageId === errorMessageId) { + return false; + } + const partial = await getMessages( + { user: userId, messageId: liveResponseMessageId, conversationId }, + req.body?.compact === true ? undefined : '_id', + ); + if (partial.length === 0) { + return false; + } + await finalizeFailedCompactionTurn(req, { + userId, + conversationId, + messageId: liveResponseMessageId, + partialRow: partial[0], + }); + return true; + }; const existing = await getMessages( { user: userId, messageId: errorMessageId, conversationId }, '_id', ); if (existing.length > 0) { - /** No compaction finalization here: this id can normalize back to the - * anchor itself when the anchor ends in `_`, and the anchor is never - * the failed run's row. The run's own snapshot, if any, is checked - * under its distinct live response id below. */ + /** This id can normalize back to the compaction anchor itself when the + * anchor ends in `_`: the failed run's row is the distinct live + * response id, so a compaction settles there and never writes the + * error row over whatever matched here. */ + if (req.body?.compact === true) { + await settleLiveSnapshot(); + } return; } - if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { - /** Full documents only where the compaction finalization needs the - * content; ordinary failures keep the id-only projection. */ - const partial = await getMessages( - { user: userId, messageId: liveResponseMessageId, conversationId }, - req.body?.compact === true ? undefined : '_id', - ); - if (partial.length > 0) { - await finalizeFailedCompactionTurn(req, { - userId, - conversationId, - messageId: liveResponseMessageId, - partialRow: partial[0], - }); - return; - } + if (await settleLiveSnapshot()) { + return; } const reqCtx = { From d6bf0ac444f0a06e22682958d21f3e6410e24dd8 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 15:59:14 +0200 Subject: [PATCH 07/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Lift=20the=20Failed?= =?UTF-8?q?-Turn=20Settlement=20Into=20the=20Compaction=20Module?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../request.partialDisconnect.spec.js | 2 + .../__tests__/request.resumeMetadata.spec.js | 2 + api/server/controllers/agents/request.js | 94 +++----------- api/server/routes/agents/index.js | 41 +++--- packages/api/src/agents/compaction.spec.ts | 119 ++++++++++++++++++ packages/api/src/agents/compaction.ts | 101 +++++++++++++++ 6 files changed, 261 insertions(+), 98 deletions(-) diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index 2350045fe48..6fa785ade0c 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -57,6 +57,8 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + settleExistingRowsBeforeErrorTurn: (...args) => + jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), persistedReasoningOverrideFields: jest.requireActual('@librechat/api').persistedReasoningOverrideFields, diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index 6125e4aa7e0..3d976876ada 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -269,6 +269,8 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + settleExistingRowsBeforeErrorTurn: (...args) => + jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), /** Real, because whether a skipped-persistence turn may raise an indicator is under test. */ isAnnounceableReply: jest.requireActual('@librechat/api').isAnnounceableReply, diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 67193e19e4a..7d584ce9d77 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -63,7 +63,7 @@ const { stampPreliminaryPrivateTextMessage, announceReply, announceErrorTurn, - persistFinalizedCompactionTurn, + settleExistingRowsBeforeErrorTurn, markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); @@ -459,48 +459,6 @@ async function saveErrorTurn( } const userId = req.user.id; - /** Whether the failed run's own partial row exists under its distinct - * live response id, finalizing it when it does (a no-op for turns that - * were not compactions). Full documents are requested only where the - * compaction finalization needs the content; ordinary failures keep the - * id-only projection. */ - const settleLiveSnapshot = async () => { - if (liveResponseMessageId == null || liveResponseMessageId === errorMessageId) { - return false; - } - const partial = await getMessages( - { user: userId, messageId: liveResponseMessageId, conversationId }, - req.body?.compact === true ? undefined : '_id', - ); - if (partial.length === 0) { - return false; - } - await finalizeFailedCompactionTurn(req, { - userId, - conversationId, - messageId: liveResponseMessageId, - partialRow: partial[0], - }); - return true; - }; - const existing = await getMessages( - { user: userId, messageId: errorMessageId, conversationId }, - '_id', - ); - if (existing.length > 0) { - /** This id can normalize back to the compaction anchor itself when the - * anchor ends in `_`: the failed run's row is the distinct live - * response id, so a compaction settles there and never writes the - * error row over whatever matched here. */ - if (req.body?.compact === true) { - await settleLiveSnapshot(); - } - return; - } - if (await settleLiveSnapshot()) { - return; - } - const reqCtx = { userId, isTemporary: @@ -511,6 +469,24 @@ async function saveErrorTurn( req?._agentEventBindingRetention?.expiredAt ?? req?.resolvedConversation?.expiredAt, interfaceConfig: req?.config?.interfaceConfig, }; + /** The existing-row settlement (which row a failed turn settles, and + * whether its error row may be written at all) lives in @librechat/api; + * this supplies the caller's reads and write. */ + const coveredByExistingRow = await settleExistingRowsBeforeErrorTurn(req.body, { + userId, + conversationId, + errorMessageId, + liveResponseMessageId, + getMessages, + saveFinalizedTurn: (message) => + saveMessage(reqCtx, message, { + context: 'api/server/controllers/agents/request.js - finalize failed compaction turn', + }), + }); + if (coveredByExistingRow) { + return; + } + const context = 'api/server/controllers/agents/request.js - failed turn'; const endpoint = endpointOption?.endpoint; const model = getAgentResponseModel(req, endpointOption); @@ -629,38 +605,6 @@ async function saveErrorTurn( * marking. The decision lives in @librechat/api; this is the wiring, reusing * the row the caller already loaded. */ -/** - * The disconnect save is marker-only while the run is still live; a failed - * turn is what settles it, so a compaction's partial row is finalized with the - * terminal outcome instead of keeping the snapshot's live-run marking. The - * operation lives in @librechat/api; this wiring supplies the caller's - * persistence and the row saveErrorTurn already loaded. - */ -async function finalizeFailedCompactionTurn( - req, - { userId, conversationId, messageId, partialRow }, -) { - await persistFinalizedCompactionTurn(partialRow, req.body, { - messageId, - conversationId, - saveMessage: (message) => - saveMessage( - { - userId, - isTemporary: - req?._agentEventBindingRetention?.isTemporary ?? - req?.resolvedConversation?.isTemporary ?? - req?.body?.isTemporary, - expiredAt: - req?._agentEventBindingRetention?.expiredAt ?? req?.resolvedConversation?.expiredAt, - interfaceConfig: req?.config?.interfaceConfig, - }, - message, - { context: 'api/server/controllers/agents/request.js - finalize failed compaction turn' }, - ), - }); -} - function classifyScheduledFailure(error, aborted = false) { if (aborted || error?.code === 'SCHEDULE_NO_LONGER_ACTIVE') { return { status: 'interrupted', error: error?.message }; diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index cf716a37134..41bfbbd4f60 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -7,6 +7,7 @@ const { hasPersistableAbortContent, announceStoppedReply, resolveAbortedTurnAnchorDecision, + planAbortedTurnPersistence, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -791,33 +792,27 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * its parent and the preliminary-parent fence correctly rejects it. */ const shouldPersistAbortedTurn = hasPersistableAbortContent(content) || jobData?.createdEventEmitted === true; - /** A compaction anchors on the persisted branch leaf, so its - * existence decides the prerequisite write: present, the projected - * anchor is never upserted over it; absent (Stop won the race - * before the branch loaded), there is nothing to parent the aborted - * response onto and nothing is written. */ - /** The compaction anchor policy (never upsert the projected leaf, - * never write a response with nothing to hang it on) lives in - * @librechat/api; this supplies the route's reader. */ - const anchorDecision = await resolveAbortedTurnAnchorDecision(jobData, { - messageExists: (messageId, conversationId) => - getMessages({ user: req?.user?.id, messageId, conversationId }).then( - (rows) => rows.length > 0, - ), - }); - if (anchorDecision === 'skip-turn') { - /** Throwing here is the contract for "do not publish the normal - * FINAL": the manager emits a reconciliation frame instead of - * one whose response points at a row deliberately never - * persisted. */ - persistenceErrors.push(new Error('Compaction anchor unavailable; abort turn withheld')); + /** The stopped turn's persistence plan (which rows to write, and + * whether the normal FINAL must be withheld for a reconciliation + * frame instead) comes from @librechat/api, decided from the + * compaction anchor this route reads. */ + const abortPersistencePlan = planAbortedTurnPersistence( + await resolveAbortedTurnAnchorDecision(jobData, { + messageExists: (messageId, conversationId) => + getMessages({ user: req?.user?.id, messageId, conversationId }).then( + (rows) => rows.length > 0, + ), + }), + shouldPersistAbortedTurn, + ); + if (abortPersistencePlan.withholdFinal && abortPersistencePlan.withholdReason) { + persistenceErrors.push(new Error(abortPersistencePlan.withholdReason)); } if ( jobData?.userMessage?.messageId && jobData?.responseMessageId && - shouldPersistAbortedTurn && - anchorDecision !== 'skip-turn' + abortPersistencePlan.writeResponseRow ) { const messageContext = { userId: req?.user?.id, @@ -879,7 +874,7 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * operation gets a chance to succeed. A compaction skips the * prerequisite: its anchor is the persisted leaf itself. */ let persistedRequestId; - if (anchorDecision === 'persist') { + if (abortPersistencePlan.writeUserRow) { try { const persistedRequest = await saveAbortedUserMessage( { saveMessage, getPersistedPrivateTextId, getPrivateMessageTexts }, diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index da5afeca227..0cb63af4212 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -18,7 +18,9 @@ import { markAbortedCompactionContent, markCompactionOutcome, persistFinalizedCompactionTurn, + planAbortedTurnPersistence, resolveAbortedTurnAnchorDecision, + settleExistingRowsBeforeErrorTurn, resolveFailedTurnContent, resolveCheckpointMessage, resolveFinalizedCompactionTurn, @@ -978,3 +980,120 @@ describe('unusable summary parts', () => { expect(stripUnusableSummaryParts(payload)).toBe(payload); }); }); + +describe('planAbortedTurnPersistence', () => { + it('writes both rows for an ordinary turn the abort must persist', () => { + expect(planAbortedTurnPersistence('persist', true)).toEqual({ + writeUserRow: true, + writeResponseRow: true, + withholdFinal: false, + }); + }); + + it('writes only the response for a compaction anchored on a persisted leaf', () => { + expect(planAbortedTurnPersistence('skip-anchor', true)).toEqual({ + writeUserRow: false, + writeResponseRow: true, + withholdFinal: false, + }); + }); + + it('withholds the final and every row when the anchor never persisted', () => { + expect(planAbortedTurnPersistence('skip-turn', true)).toEqual({ + writeUserRow: false, + writeResponseRow: false, + withholdFinal: true, + withholdReason: expect.stringContaining('anchor unavailable'), + }); + }); + + it('writes nothing for a turn the abort would not persist', () => { + expect(planAbortedTurnPersistence('persist', false)).toEqual({ + writeUserRow: false, + writeResponseRow: false, + withholdFinal: false, + }); + }); +}); + +describe('settleExistingRowsBeforeErrorTurn', () => { + const partialSummaryRow = () => ({ + messageId: 'live-response', + unfinished: true, + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + }, + ], + }); + const deps = (rowsByMessageId: Record) => { + const saved: Record[] = []; + return { + saved, + deps: { + userId: 'user-1', + conversationId: 'conversation-1', + errorMessageId: 'error-target', + liveResponseMessageId: 'live-response', + getMessages: jest.fn(async ({ messageId }: { messageId: string }) => + (rowsByMessageId[messageId] ?? []).map((row) => row), + ) as never, + saveFinalizedTurn: async (message: Record) => { + saved.push(message); + return message; + }, + }, + }; + }; + + it('settles a compaction snapshot under its live id and blocks the error row', async () => { + const { saved, deps: d } = deps({ 'live-response': [partialSummaryRow()] }); + + await expect(settleExistingRowsBeforeErrorTurn({ compact: true }, d)).resolves.toBe(true); + + expect(saved).toHaveLength(1); + expect(saved[0]).toMatchObject({ + messageId: 'live-response', + unfinished: false, + error: true, + }); + }); + + /** The error id normalizes back to the anchor itself when the anchor ends + * in `_`: the anchor match must not stop the live row from settling, and + * nothing may be written over that match. */ + it('settles the live row past an anchor-shaped collision', async () => { + const { saved, deps: d } = deps({ + 'error-target': [{ messageId: 'error-target', _id: 'anchor-shaped-match' }], + 'live-response': [partialSummaryRow()], + }); + + await expect(settleExistingRowsBeforeErrorTurn({ compact: true }, d)).resolves.toBe(true); + + expect(saved).toHaveLength(1); + expect(saved[0]).toMatchObject({ messageId: 'live-response' }); + }); + + it('blocks the error row for an ordinary turn with an existing row, writing nothing', async () => { + const { saved, deps: d } = deps({ + 'error-target': [{ messageId: 'error-target', _id: 'existing' }], + 'live-response': [{ messageId: 'live-response', _id: 'partial' }], + }); + + await expect(settleExistingRowsBeforeErrorTurn({}, d)).resolves.toBe(true); + + expect(saved).toHaveLength(0); + // The ordinary early return never reads the live row. + expect(d.getMessages).toHaveBeenCalledTimes(1); + }); + + it('lets the error row through when no row covers the turn', async () => { + const { saved, deps: d } = deps({}); + + await expect(settleExistingRowsBeforeErrorTurn({ compact: true }, d)).resolves.toBe(false); + + expect(saved).toHaveLength(0); + }); +}); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 66d762726c4..22278297244 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -303,6 +303,107 @@ export async function resolveAbortedTurnAnchorDecision( } } +/** The abort route's persistence plan for a stopped turn: which rows to write + * and whether the normal FINAL must be withheld (the manager publishes a + * reconciliation frame instead, so the client is never pointed at a response + * that was deliberately never persisted). */ +export interface AbortedTurnPersistencePlan { + writeUserRow: boolean; + writeResponseRow: boolean; + withholdFinal: boolean; + withholdReason?: string; +} + +export function planAbortedTurnPersistence( + anchorDecision: AbortAnchorDecision, + shouldPersistAbortedTurn: boolean, +): AbortedTurnPersistencePlan { + const active = shouldPersistAbortedTurn && anchorDecision !== 'skip-turn'; + return { + writeUserRow: active && anchorDecision === 'persist', + writeResponseRow: active, + withholdFinal: anchorDecision === 'skip-turn', + ...(anchorDecision === 'skip-turn' && { + withholdReason: 'Compaction anchor unavailable; abort turn withheld', + }), + }; +} + +/** A message row as the failed-turn settlement reads it: identity for the + * anchor-shaped check, content and envelope for the live row it finalizes. */ +export type ReadableMessageRow = { + messageId: string; + content?: unknown; + unfinished?: boolean; +}; + +/** + * Settles the rows a failed generation already persisted before its error row + * is written, through the caller's injected reads and write. Returns whether + * an existing row covers the turn, in which case the caller skips the fresh + * error row entirely. + * + * The error id can normalize back to the compaction anchor itself when the + * anchor ends in `_`: a match there never receives the error row, and the + * failed run settles its own distinct live response row instead. Ordinary + * turns keep their existing behavior: a found partial row is preserved as it + * stands and blocks the error row. + */ +export async function settleExistingRowsBeforeErrorTurn( + requestBody: { compact?: boolean } | null | undefined, + { + userId, + conversationId, + errorMessageId, + liveResponseMessageId, + getMessages, + saveFinalizedTurn, + }: { + userId: string; + conversationId: string; + errorMessageId: string; + liveResponseMessageId?: string | null; + getMessages: ( + filter: { user: string; messageId: string; conversationId: string }, + projection?: string, + ) => Promise; + saveFinalizedTurn: (message: Record) => Promise; + }, +): Promise { + const isCompaction = requestBody?.compact === true; + const settleLiveRow = async (): Promise => { + if (liveResponseMessageId == null || liveResponseMessageId === errorMessageId) { + return false; + } + /** Full documents only where the compaction finalization needs the + * content; ordinary failures keep the id-only projection. */ + const partial = await getMessages( + { user: userId, messageId: liveResponseMessageId, conversationId }, + isCompaction ? undefined : '_id', + ); + if (partial.length === 0) { + return false; + } + await persistFinalizedCompactionTurn(partial[0], requestBody, { + messageId: liveResponseMessageId, + conversationId, + saveMessage: saveFinalizedTurn, + }); + return true; + }; + const existing = await getMessages( + { user: userId, messageId: errorMessageId, conversationId }, + '_id', + ); + if (existing.length > 0) { + if (isCompaction) { + await settleLiveRow(); + } + return true; + } + return settleLiveRow(); +} + /** * Finalizes a failed compaction's already-persisted partial row, with the * write injected so the operation runs against whatever persistence the From d208156686cc57d2f0585a1e2a92d36db30012d9 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:39:12 +0200 Subject: [PATCH 08/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Keep=20a=20Removed?= =?UTF-8?q?=20Round's=20Failure=20and=20Guard=20Settled=20Jobs=20From=20Sn?= =?UTF-8?q?apshots?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../request.partialDisconnect.spec.js | 13 ++++++ .../__tests__/request.resumeMetadata.spec.js | 1 + api/server/controllers/agents/request.js | 9 ++++ packages/api/src/agents/compaction.spec.ts | 41 +++++++++++++++++++ packages/api/src/agents/compaction.ts | 26 +++++++++++- 5 files changed, 89 insertions(+), 1 deletion(-) diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index 6fa785ade0c..374f3e64df0 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -57,6 +57,7 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), settleExistingRowsBeforeErrorTurn: (...args) => jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), @@ -344,4 +345,16 @@ describe('ResumableAgentController tenant context', () => { }); expect(savedMessage.content).toHaveLength(1); }); + /** The settling path (completion, error, abort) owns the final row: a + * disconnect snapshot landing after it would reopen the settled turn as + * an unfinished response. */ + it('skips the partial save when the job record has settled', async () => { + await firePartialDisconnect( + { id: 'user-123' }, + { createdAt: 1000, status: 'error' }, + { aggregatedContent: [{ type: 'text', text: 'Partial response' }] }, + ); + + expect(mockSaveMessage).not.toHaveBeenCalled(); + }); }); diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index 3d976876ada..c4aa1af0e73 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -269,6 +269,7 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), settleExistingRowsBeforeErrorTurn: (...args) => jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 7d584ce9d77..00f6ead704b 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -64,6 +64,7 @@ const { announceReply, announceErrorTurn, settleExistingRowsBeforeErrorTurn, + isSettledJobRecord, markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); @@ -1946,6 +1947,14 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit * record is the source, since the client-facing resume snapshot never * carries server-private state. */ const contextMeta = jobRecord?.createdAt === jobCreatedAt ? jobRecord.contextMeta : undefined; + /** A job whose settling path (completion, error, abort) owns the final + * row must not have it reopened as an unfinished snapshot here; the + * guard reads the same record, so the window is the settling path's + * own commit span. */ + if (isSettledJobRecord(jobRecord, jobCreatedAt)) { + logger.debug('[ResumableAgentController] Skipping partial response save for a settled job'); + return; + } try { const partialMessage = { diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 0cb63af4212..4ad98c43330 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -18,6 +18,7 @@ import { markAbortedCompactionContent, markCompactionOutcome, persistFinalizedCompactionTurn, + isSettledJobRecord, planAbortedTurnPersistence, resolveAbortedTurnAnchorDecision, settleExistingRowsBeforeErrorTurn, @@ -402,6 +403,27 @@ describe('markAbortedCompactionContent', () => { ]); }); + /** An earlier round's checkpoint is not the stopped round's outcome: the + * failure lands beside it, or the row reads as the successful compaction + * the checkpoint describes (and as a leaf that can no longer compact). */ + it('records the typed failure beside an earlier checkpoint when the current round streamed nothing', () => { + const parts = [completedSummary('An earlier checkpoint.'), emptySummaryPlaceholder()]; + + markAbortedCompactionContent(parts, true); + + expect(parts).toEqual([ + expect.objectContaining({ + type: ContentTypes.SUMMARY, + initiatedBy: 'user', + }), + { + type: ContentTypes.ERROR, + error: JSON.stringify({ type: ErrorTypes.COMPACTION_FAILED }), + initiatedBy: 'user', + }, + ]); + }); + /** A run stopped before any part streamed still needs an identifiable row: * an empty one reads as an answer to the message it hangs off. */ it('records the typed failure when nothing streamed before the stop', () => { @@ -1097,3 +1119,22 @@ describe('settleExistingRowsBeforeErrorTurn', () => { expect(saved).toHaveLength(0); }); }); + +describe('isSettledJobRecord', () => { + it.each(['complete', 'error', 'aborted'])('treats a %s record as settled', (status) => { + expect(isSettledJobRecord({ createdAt: 1000, status })).toBe(true); + }); + + it('leaves live and missing records unsettled', () => { + expect(isSettledJobRecord({ createdAt: 1000, status: 'running' })).toBe(false); + expect(isSettledJobRecord({ createdAt: 1000, status: 'requires_action' })).toBe(false); + expect(isSettledJobRecord(null)).toBe(false); + expect(isSettledJobRecord(undefined)).toBe(false); + }); + + /** Another epoch's record describes a different generation, not this one. */ + it('ignores a record from another epoch', () => { + expect(isSettledJobRecord({ createdAt: 2000, status: 'error' }, 1000)).toBe(false); + expect(isSettledJobRecord({ createdAt: 1000, status: 'error' }, 1000)).toBe(true); + }); +}); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 22278297244..06dfc3afd74 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -224,6 +224,7 @@ export function markAbortedCompactionContent( return contentParts; } let hasOutcome = false; + let removedUnfinishedRound = false; for (let index = contentParts.length - 1; index >= 0; index -= 1) { const part = contentParts[index]; if (part == null) { @@ -256,13 +257,36 @@ export function markAbortedCompactionContent( continue; } contentParts.splice(index, 1); + removedUnfinishedRound = true; } - if (!hasOutcome && synthesizeFailure) { + /** An earlier round's checkpoint is not this round's outcome: a round the + * run opened but never finished still records the typed failure beside it, + * or the stopped turn reads as the successful compaction the checkpoint + * describes. */ + if ((!hasOutcome || removedUnfinishedRound) && synthesizeFailure) { contentParts.push(...compactionFailureContent()); } return contentParts; } +/** Whether a job record has reached a status whose path owns the turn's final + * row (completion, error, or abort): the disconnect snapshot must not be + * written over it, or the settled row reopens as an unfinished response. + * Only a same-epoch record is trusted. */ +export function isSettledJobRecord( + jobRecord: { createdAt?: number; status?: string } | null | undefined, + jobCreatedAt?: number, +): boolean { + if (jobRecord == null || (jobCreatedAt != null && jobRecord.createdAt !== jobCreatedAt)) { + return false; + } + return ( + jobRecord.status === 'complete' || + jobRecord.status === 'error' || + jobRecord.status === 'aborted' + ); +} + /** How the abort route persists a stopped turn's prerequisite rows. */ export type AbortAnchorDecision = 'persist' | 'skip-anchor' | 'skip-turn'; From 4643ffd08ac6f2e476f3b90478abf363f139d64e Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:58:15 +0200 Subject: [PATCH 09/19] =?UTF-8?q?=E2=9A=A1=20perf:=20Read=20Only=20the=20A?= =?UTF-8?q?nchor=20Id=20in=20the=20Abort=20Existence=20Check?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api/server/routes/agents/index.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index 41bfbbd4f60..e62220dfa34 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -799,7 +799,7 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { const abortPersistencePlan = planAbortedTurnPersistence( await resolveAbortedTurnAnchorDecision(jobData, { messageExists: (messageId, conversationId) => - getMessages({ user: req?.user?.id, messageId, conversationId }).then( + getMessages({ user: req?.user?.id, messageId, conversationId }, '_id').then( (rows) => rows.length > 0, ), }), From bac57d38a748ecae8a9b9e445754fb4d52417fb0 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:44:38 +0200 Subject: [PATCH 10/19] =?UTF-8?q?=F0=9F=93=B4=20fix:=20Withhold=20the=20Ab?= =?UTF-8?q?ort=20Final=20Only=20When=20a=20Row=20Needed=20Writing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/api/src/agents/compaction.spec.ts | 11 +++++++++++ packages/api/src/agents/compaction.ts | 8 ++++++-- 2 files changed, 17 insertions(+), 2 deletions(-) diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 4ad98c43330..784580181d4 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -1036,6 +1036,17 @@ describe('planAbortedTurnPersistence', () => { withholdFinal: false, }); }); + + /** An abort with no persistable content and no created event publishes an + * early-abort FINAL of its own; withholding it would replace that frame + * with a reconciliation one even though no row was ever at stake. */ + it('does not withhold the final when no row needed writing', () => { + expect(planAbortedTurnPersistence('skip-turn', false)).toEqual({ + writeUserRow: false, + writeResponseRow: false, + withholdFinal: false, + }); + }); }); describe('settleExistingRowsBeforeErrorTurn', () => { diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 06dfc3afd74..2307ab9e11c 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -343,11 +343,15 @@ export function planAbortedTurnPersistence( shouldPersistAbortedTurn: boolean, ): AbortedTurnPersistencePlan { const active = shouldPersistAbortedTurn && anchorDecision !== 'skip-turn'; + /** Withholding the FINAL only matters when a row would otherwise have been + * written: an abort with no persistable content and no created event + * publishes an early-abort FINAL of its own, and nothing was withheld. */ + const withhold = shouldPersistAbortedTurn && anchorDecision === 'skip-turn'; return { writeUserRow: active && anchorDecision === 'persist', writeResponseRow: active, - withholdFinal: anchorDecision === 'skip-turn', - ...(anchorDecision === 'skip-turn' && { + withholdFinal: withhold, + ...(withhold && { withholdReason: 'Compaction anchor unavailable; abort turn withheld', }), }; From ff342a6e3e21f3910a90d5bb4f1e19caf763751b Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 09:51:12 +0200 Subject: [PATCH 11/19] test: Cover Stopped Compaction Finalization and the Missing Anchor in E2E --- .../compaction-abort-finalize.spec.ts | 216 ++++++++++++++++++ 1 file changed, 216 insertions(+) create mode 100644 e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts diff --git a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts new file mode 100644 index 00000000000..0ce7d18a731 --- /dev/null +++ b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts @@ -0,0 +1,216 @@ +import { expect, test } from '@playwright/test'; +import type { APIRequestContext, Page } from '@playwright/test'; +import { randomUUID } from 'node:crypto'; +import { getE2EUser } from '../../../setup/user'; +import { deleteConversations, deleteMessagesByConversation, seedMessages, withMongo } from '../db'; +import { messagesView, sendMessage, sendMessageAndWaitForCompletion } from '../helpers'; + +const userEmail = getE2EUser().email; +const LABEL_SERVER = `http://127.0.0.1:${process.env.E2E_LABEL_PORT || '8889'}`; + +type Row = Record; + +async function cleanup(conversationId: string) { + await deleteMessagesByConversation([conversationId]); + await deleteConversations([conversationId]); +} + +function findRow(filter: Row): Promise { + return withMongo((db) => db.collection('messages').findOne(filter)); +} + +/** A finished exchange plus a user leaf seeded under its answer, then a manual + * compaction held open on that leaf until the caller presses Stop. */ +async function startCompactionOnUserLeaf(page: Page, request: APIRequestContext, label: string) { + await page.goto('/c/new'); + await sendMessageAndWaitForCompletion(page, `tell me about ${label}`); + const conversationId = new URL(page.url()).pathname.replace('/c/', ''); + expect(conversationId).not.toBe('new'); + + const answer = await withMongo((db) => + db + .collection('messages') + .findOne({ conversationId, isCreatedByUser: false }, { sort: { createdAt: -1 } }), + ); + expect(answer?.messageId).toBeTruthy(); + + const leafUserId = randomUUID(); + const leafText = `Compact this before answering ${label}`; + await seedMessages(userEmail, conversationId, [ + { + messageId: leafUserId, + parentMessageId: answer?.messageId as string, + text: leafText, + isCreatedByUser: true, + sender: 'User', + }, + ]); + + const behavior = await request.post(`${LABEL_SERVER}/__e2e/behavior`, { + data: { mode: 'ok', delayMs: 60_000 }, + }); + expect(behavior.ok()).toBeTruthy(); + + await page.goto(`/c/${conversationId}`); + await expect(messagesView(page).getByText(leafText)).toBeVisible(); + await page.getByTestId('token-usage').click(); + await page.getByRole('button', { name: 'Compact context' }).click(); + const stop = page.getByTestId('stop-generation-button'); + await expect(stop).toBeVisible({ timeout: 20_000 }); + return { conversationId, leafUserId, leafText, stop }; +} + +test.describe('compaction abort finalize', () => { + test.afterEach(async ({ request }) => { + const response = await request.post(`${LABEL_SERVER}/__e2e/reset`); + expect(response.ok()).toBeTruthy(); + }); + + /* The stopped compaction's anchor is the persisted leaf itself: the abort + writes a settled response under it and leaves the leaf exactly as it was. */ + test('a stopped compaction settles under its anchor and survives a reload @scenario:stopped-compaction-settles-under-its-anchor', async ({ + page, + request, + }) => { + const { conversationId, leafUserId, leafText, stop } = await startCompactionOnUserLeaf( + page, + request, + 'stopped-anchor', + ); + try { + await stop.click(); + + let compaction: Row | null = null; + await expect + .poll( + async () => { + compaction = await findRow({ + conversationId, + parentMessageId: leafUserId, + isCreatedByUser: false, + }); + return compaction != null; + }, + { timeout: 20_000 }, + ) + .toBeTruthy(); + const compactionId = (compaction as Row | null)?.messageId as string; + await expect(stop).toBeHidden({ timeout: 20_000 }); + + /* A live snapshot's shape would leave the turn reading as still running. */ + const settled = await findRow({ conversationId, messageId: compactionId }); + expect(settled?.unfinished).not.toBe(true); + + const leaf = await findRow({ conversationId, messageId: leafUserId }); + expect(leaf?.isCreatedByUser).toBe(true); + expect(leaf?.text).toBe(leafText); + + await page.reload(); + const row = page.locator(`[id="${compactionId}"]`); + await expect(row).toBeVisible({ timeout: 20_000 }); + await expect(row.getByText('Could not compact the context', { exact: false })).toBeVisible(); + await expect(row.getByText('Summarizing...')).toHaveCount(0); + await expect(messagesView(page).getByText(leafText)).toBeVisible(); + await expect( + page.getByRole('navigation', { name: 'Sibling message navigation' }), + ).toHaveCount(0); + } finally { + await cleanup(conversationId); + } + }); + + /* Stop can win before the anchor exists in storage. The response would then + hang off a row that was never written, so the abort persists nothing for + it, and the run still settles instead of waiting on a final. */ + test('a stopped compaction whose anchor is not stored writes no orphaned response @scenario:stopped-compaction-without-stored-anchor-writes-no-orphan', async ({ + page, + request, + }) => { + const { conversationId, leafUserId, stop } = await startCompactionOnUserLeaf( + page, + request, + 'missing-anchor', + ); + try { + await withMongo((db) => + db.collection('messages').deleteOne({ conversationId, messageId: leafUserId }), + ); + + const abort = page.waitForResponse( + (response) => + response.request().method() === 'POST' && + new URL(response.url()).pathname === '/api/agents/chat/abort', + { timeout: 20_000 }, + ); + await stop.click(); + expect((await abort).ok()).toBeTruthy(); + await expect(stop).toBeHidden({ timeout: 20_000 }); + + /* Neither the anchor nor a response parented on it may appear later. */ + await expect + .poll( + () => + withMongo((db) => + db.collection('messages').countDocuments({ + conversationId, + $or: [{ messageId: leafUserId }, { parentMessageId: leafUserId }], + }), + ), + { timeout: 5_000, intervals: [1_000] }, + ) + .toBe(0); + } finally { + await cleanup(conversationId); + } + }); + + /* An ordinary reply carries no compaction anchor: Stop keeps the user turn + and the partial reply under it, as it did before. */ + test('a stopped ordinary reply keeps its turn and partial answer @scenario:stopped-reply-keeps-turn-and-partial-answer', async ({ + page, + }) => { + test.setTimeout(120_000); + const label = `stopped-reply-${randomUUID().slice(0, 8)}`; + const prompt = `E2E_SLOW_REPLY:${label}`; + + await page.goto('/c/new'); + const run = await sendMessage(page, prompt); + expect(run.ok()).toBeTruthy(); + await expect(messagesView(page).getByText('chunk-010')).toBeVisible({ timeout: 15_000 }); + await expect(page).toHaveURL(/\/c\/[0-9a-fA-F-]{36}$/, { timeout: 15_000 }); + const conversationId = new URL(page.url()).pathname.replace('/c/', ''); + + try { + const stop = page.getByTestId('stop-generation-button'); + await stop.click(); + await expect(stop).toBeHidden({ timeout: 20_000 }); + + let reply: Row | null = null; + await expect + .poll( + async () => { + const user = await findRow({ conversationId, isCreatedByUser: true }); + if (!user) { + return false; + } + reply = await findRow({ + conversationId, + parentMessageId: user.messageId, + isCreatedByUser: false, + }); + return reply != null; + }, + { timeout: 20_000 }, + ) + .toBeTruthy(); + expect(reply).not.toBeNull(); + + await page.reload(); + await expect(messagesView(page).getByText(prompt)).toBeVisible({ timeout: 20_000 }); + await expect(messagesView(page).getByText('chunk-010')).toBeVisible(); + await expect(page.getByTestId('stop-generation-button')).toHaveCount(0); + } finally { + await cleanup(conversationId); + } + }); +}); From 9b47aab93a06f5f7fd02f2f46dcd4f3103ad0e88 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 10:10:09 +0200 Subject: [PATCH 12/19] fix: Withhold a Stopped Compaction That Carries No Anchor Id --- packages/api/src/agents/compaction.spec.ts | 11 +++++++++++ packages/api/src/agents/compaction.ts | 8 ++++++-- 2 files changed, 17 insertions(+), 2 deletions(-) diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 784580181d4..5615380b303 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -493,6 +493,17 @@ describe('resolveAbortedTurnAnchorDecision', () => { ).resolves.toBe('skip-turn'); }); + it('skips the whole turn when a compaction carries no anchor id', async () => { + const messageExists = reader(true); + + for (const userMessage of [undefined, null, {}, { messageId: '' }]) { + await expect( + resolveAbortedTurnAnchorDecision({ ...jobData, userMessage }, { messageExists }), + ).resolves.toBe('skip-turn'); + } + expect(messageExists).not.toHaveBeenCalled(); + }); + /** A read that throws must not escape past the caller's remaining cleanup: * nothing is known about the anchor, so nothing is written either. */ it('skips the whole turn when the anchor read fails', async () => { diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 2307ab9e11c..20023d2df7f 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -315,10 +315,14 @@ export async function resolveAbortedTurnAnchorDecision( messageExists, }: { messageExists: (messageId: string, conversationId?: string) => Promise }, ): Promise { - const anchorId = jobData?.userMessage?.messageId; - if (jobData?.compact !== true || anchorId == null || anchorId.length === 0) { + if (jobData?.compact !== true) { return 'persist'; } + /** A compaction with no anchor id has nothing to hang its response on. */ + const anchorId = jobData.userMessage?.messageId; + if (anchorId == null || anchorId.length === 0) { + return 'skip-turn'; + } try { const anchorExists = await messageExists(anchorId, jobData.conversationId); return anchorExists ? 'skip-anchor' : 'skip-turn'; From a4f81f2b4a4161898b2fb4a8b577675661128d34 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 11:21:27 +0200 Subject: [PATCH 13/19] fix: Settle Stopped Compactions, Announce Finalized Rows, and Own the Abort Read in packages/api --- api/server/controllers/agents/request.js | 11 +++ .../routes/agents/__tests__/abort.spec.js | 2 +- api/server/routes/agents/index.js | 19 ++--- .../compaction-abort-finalize.spec.ts | 26 ++++-- packages/api/src/agents/compaction.spec.ts | 83 +++++++++++++++++++ packages/api/src/agents/compaction.ts | 44 +++++++++- 6 files changed, 163 insertions(+), 22 deletions(-) diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 00f6ead704b..06adc3e7b4b 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -483,6 +483,17 @@ async function saveErrorTurn( saveMessage(reqCtx, message, { context: 'api/server/controllers/agents/request.js - finalize failed compaction turn', }), + announceSettledTurn: (messageId) => + announceErrorTurn( + { stampConvoLastResponse }, + { + userId, + conversationId, + messageId, + isTemporary: reqCtx.isTemporary, + context: 'AgentController - finalized failed compaction turn', + }, + ), }); if (coveredByExistingRow) { return; diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index 64c72268d8d..8ba35e57f68 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -439,7 +439,7 @@ describe('Agent Abort Endpoint', () => { expect.objectContaining({ messageId: compactionRowId, parentMessageId: anchorId, - unfinished: true, + unfinished: false, isCreatedByUser: false, }), expect.objectContaining({ context: expect.stringContaining('abort endpoint') }), diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index e62220dfa34..04e22c246a1 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -6,8 +6,7 @@ const { TERMINAL_PUBLICATION_RECONNECT_ERROR, hasPersistableAbortContent, announceStoppedReply, - resolveAbortedTurnAnchorDecision, - planAbortedTurnPersistence, + resolveAbortedTurnPersistence, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -796,18 +795,12 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { * whether the normal FINAL must be withheld for a reconciliation * frame instead) comes from @librechat/api, decided from the * compaction anchor this route reads. */ - const abortPersistencePlan = planAbortedTurnPersistence( - await resolveAbortedTurnAnchorDecision(jobData, { - messageExists: (messageId, conversationId) => - getMessages({ user: req?.user?.id, messageId, conversationId }, '_id').then( - (rows) => rows.length > 0, - ), - }), + const abortPersistencePlan = await resolveAbortedTurnPersistence( + jobData, shouldPersistAbortedTurn, + { userId: req?.user?.id, getMessages }, ); - if (abortPersistencePlan.withholdFinal && abortPersistencePlan.withholdReason) { - persistenceErrors.push(new Error(abortPersistencePlan.withholdReason)); - } + persistenceErrors.push(...abortPersistencePlan.persistenceErrors); if ( jobData?.userMessage?.messageId && @@ -842,7 +835,7 @@ router.post('/chat/abort', chatConfigMiddleware, async (req, res, next) => { endpoint: jobData.endpoint, iconURL: jobData.iconURL, model: jobData.model, - unfinished: true, + unfinished: abortPersistencePlan.responseUnfinished, error: false, isCreatedByUser: false, ...(Array.isArray(jobData.userSubmittedPaths) && diff --git a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts index 0ce7d18a731..8303f0eab1e 100644 --- a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts +++ b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts @@ -3,7 +3,14 @@ import type { APIRequestContext, Page } from '@playwright/test'; import { randomUUID } from 'node:crypto'; import { getE2EUser } from '../../../setup/user'; import { deleteConversations, deleteMessagesByConversation, seedMessages, withMongo } from '../db'; -import { messagesView, sendMessage, sendMessageAndWaitForCompletion } from '../helpers'; +import { + MOCK_ENDPOINTS, + NEW_CHAT_PATH, + messagesView, + sendMessage, + selectMockEndpoint, + sendMessageAndWaitForCompletion, +} from '../helpers'; const userEmail = getE2EUser().email; const LABEL_SERVER = `http://127.0.0.1:${process.env.E2E_LABEL_PORT || '8889'}`; @@ -173,18 +180,21 @@ test.describe('compaction abort finalize', () => { const label = `stopped-reply-${randomUUID().slice(0, 8)}`; const prompt = `E2E_SLOW_REPLY:${label}`; - await page.goto('/c/new'); + await page.goto(NEW_CHAT_PATH); + await selectMockEndpoint(page, MOCK_ENDPOINTS[0]); const run = await sendMessage(page, prompt); expect(run.ok()).toBeTruthy(); await expect(messagesView(page).getByText('chunk-010')).toBeVisible({ timeout: 15_000 }); + + /* Stop while the reply is still streaming; the conversation id is read + once the run has settled. */ + const stop = page.getByRole('button', { name: 'Stop generating' }); + await stop.click({ timeout: 10_000 }); + await expect(stop).toBeHidden({ timeout: 20_000 }); await expect(page).toHaveURL(/\/c\/[0-9a-fA-F-]{36}$/, { timeout: 15_000 }); const conversationId = new URL(page.url()).pathname.replace('/c/', ''); try { - const stop = page.getByTestId('stop-generation-button'); - await stop.click(); - await expect(stop).toBeHidden({ timeout: 20_000 }); - let reply: Row | null = null; await expect .poll( @@ -208,7 +218,9 @@ test.describe('compaction abort finalize', () => { await page.reload(); await expect(messagesView(page).getByText(prompt)).toBeVisible({ timeout: 20_000 }); await expect(messagesView(page).getByText('chunk-010')).toBeVisible(); - await expect(page.getByTestId('stop-generation-button')).toHaveCount(0); + /* Stopped well before the end of the scripted stream. */ + await expect(messagesView(page).getByText('chunk-159')).toHaveCount(0); + await expect(page.getByRole('button', { name: 'Stop generating' })).toHaveCount(0); } finally { await cleanup(conversationId); } diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 5615380b303..8db4d4ebefc 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -20,6 +20,7 @@ import { persistFinalizedCompactionTurn, isSettledJobRecord, planAbortedTurnPersistence, + resolveAbortedTurnPersistence, resolveAbortedTurnAnchorDecision, settleExistingRowsBeforeErrorTurn, resolveFailedTurnContent, @@ -1019,6 +1020,7 @@ describe('planAbortedTurnPersistence', () => { expect(planAbortedTurnPersistence('persist', true)).toEqual({ writeUserRow: true, writeResponseRow: true, + responseUnfinished: true, withholdFinal: false, }); }); @@ -1027,6 +1029,7 @@ describe('planAbortedTurnPersistence', () => { expect(planAbortedTurnPersistence('skip-anchor', true)).toEqual({ writeUserRow: false, writeResponseRow: true, + responseUnfinished: false, withholdFinal: false, }); }); @@ -1035,6 +1038,7 @@ describe('planAbortedTurnPersistence', () => { expect(planAbortedTurnPersistence('skip-turn', true)).toEqual({ writeUserRow: false, writeResponseRow: false, + responseUnfinished: false, withholdFinal: true, withholdReason: expect.stringContaining('anchor unavailable'), }); @@ -1044,6 +1048,7 @@ describe('planAbortedTurnPersistence', () => { expect(planAbortedTurnPersistence('persist', false)).toEqual({ writeUserRow: false, writeResponseRow: false, + responseUnfinished: true, withholdFinal: false, }); }); @@ -1055,11 +1060,68 @@ describe('planAbortedTurnPersistence', () => { expect(planAbortedTurnPersistence('skip-turn', false)).toEqual({ writeUserRow: false, writeResponseRow: false, + responseUnfinished: false, withholdFinal: false, }); }); }); +describe('resolveAbortedTurnPersistence', () => { + const jobData = { + compact: true, + conversationId: 'conversation-1', + userMessage: { messageId: 'leaf-1' }, + }; + + /** Nothing continues a stopped compaction, so its row is written settled. */ + it('writes a settled response under a persisted compaction anchor', async () => { + const getMessages = jest.fn(async () => [{ _id: 'row' }]); + + const plan = await resolveAbortedTurnPersistence(jobData, true, { + userId: 'user-1', + getMessages, + }); + + expect(getMessages).toHaveBeenCalledWith( + { user: 'user-1', messageId: 'leaf-1', conversationId: 'conversation-1' }, + '_id', + ); + expect(plan).toMatchObject({ + writeUserRow: false, + writeResponseRow: true, + responseUnfinished: false, + withholdFinal: false, + persistenceErrors: [], + }); + }); + + it('reports the withheld final when the compaction anchor is missing', async () => { + const plan = await resolveAbortedTurnPersistence(jobData, true, { + userId: 'user-1', + getMessages: jest.fn(async () => []), + }); + + expect(plan.writeResponseRow).toBe(false); + expect(plan.withholdFinal).toBe(true); + expect(plan.persistenceErrors).toHaveLength(1); + expect(plan.persistenceErrors[0].message).toContain('anchor unavailable'); + }); + + it('keeps an ordinary stopped reply unfinished without reading its anchor', async () => { + const getMessages = jest.fn(async () => []); + + const plan = await resolveAbortedTurnPersistence({}, true, { getMessages }); + + expect(getMessages).not.toHaveBeenCalled(); + expect(plan).toMatchObject({ + writeUserRow: true, + writeResponseRow: true, + responseUnfinished: true, + persistenceErrors: [], + }); + }); +}); + describe('settleExistingRowsBeforeErrorTurn', () => { const partialSummaryRow = () => ({ messageId: 'live-response', @@ -1105,6 +1167,27 @@ describe('settleExistingRowsBeforeErrorTurn', () => { }); }); + /** The error row's own path announces the persisted turn; a finalized live + * row stands in for it, so it is announced the same way. */ + it('announces the live row it finalized', async () => { + const announceSettledTurn = jest.fn(async () => undefined); + const { deps: d } = deps({ 'live-response': [partialSummaryRow()] }); + + await settleExistingRowsBeforeErrorTurn({ compact: true }, { ...d, announceSettledTurn }); + + expect(announceSettledTurn).toHaveBeenCalledWith('live-response'); + }); + + it('announces nothing when the existing row needed no write', async () => { + const announceSettledTurn = jest.fn(async () => undefined); + const { saved, deps: d } = deps({ 'live-response': [{ messageId: 'live-response' }] }); + + await settleExistingRowsBeforeErrorTurn({}, { ...d, announceSettledTurn }); + + expect(saved).toHaveLength(0); + expect(announceSettledTurn).not.toHaveBeenCalled(); + }); + /** The error id normalizes back to the anchor itself when the anchor ends * in `_`: the anchor match must not stop the live row from settling, and * nothing may be written over that match. */ diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 20023d2df7f..05fe21c6e10 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -338,6 +338,10 @@ export async function resolveAbortedTurnAnchorDecision( export interface AbortedTurnPersistencePlan { writeUserRow: boolean; writeResponseRow: boolean; + /** An ordinary stopped reply stays `unfinished` so it can be continued; a + * stopped compaction is settled, since nothing continues it and a live + * envelope would keep restored sessions reading it as still running. */ + responseUnfinished: boolean; withholdFinal: boolean; withholdReason?: string; } @@ -354,6 +358,7 @@ export function planAbortedTurnPersistence( return { writeUserRow: active && anchorDecision === 'persist', writeResponseRow: active, + responseUnfinished: anchorDecision === 'persist', withholdFinal: withhold, ...(withhold && { withholdReason: 'Compaction anchor unavailable; abort turn withheld', @@ -361,6 +366,36 @@ export function planAbortedTurnPersistence( }; } +/** + * The abort route's whole persistence decision for a stopped turn: reads the + * compaction anchor through the caller's message reader (id-only), plans the + * rows, and returns the failures the route must report so the manager + * publishes a reconciliation frame instead of a normal FINAL. + */ +export async function resolveAbortedTurnPersistence( + jobData: Parameters[0], + shouldPersistAbortedTurn: boolean, + { + userId, + getMessages, + }: { + userId?: string; + getMessages: ( + filter: { user?: string; messageId: string; conversationId?: string }, + projection?: string, + ) => Promise; + }, +): Promise { + const anchorDecision = await resolveAbortedTurnAnchorDecision(jobData, { + messageExists: async (messageId, conversationId) => + (await getMessages({ user: userId, messageId, conversationId }, '_id')).length > 0, + }); + const plan = planAbortedTurnPersistence(anchorDecision, shouldPersistAbortedTurn); + const persistenceErrors = + plan.withholdFinal && plan.withholdReason ? [new Error(plan.withholdReason)] : []; + return { ...plan, persistenceErrors }; +} + /** A message row as the failed-turn settlement reads it: identity for the * anchor-shaped check, content and envelope for the live row it finalizes. */ export type ReadableMessageRow = { @@ -390,6 +425,7 @@ export async function settleExistingRowsBeforeErrorTurn( liveResponseMessageId, getMessages, saveFinalizedTurn, + announceSettledTurn, }: { userId: string; conversationId: string; @@ -400,6 +436,9 @@ export async function settleExistingRowsBeforeErrorTurn( projection?: string, ) => Promise; saveFinalizedTurn: (message: Record) => Promise; + /** Announces a row this settlement finalized, as the error row's own + * path does, so other devices learn the persisted turn ended. */ + announceSettledTurn?: (messageId: string) => Promise; }, ): Promise { const isCompaction = requestBody?.compact === true; @@ -416,11 +455,14 @@ export async function settleExistingRowsBeforeErrorTurn( if (partial.length === 0) { return false; } - await persistFinalizedCompactionTurn(partial[0], requestBody, { + const finalized = await persistFinalizedCompactionTurn(partial[0], requestBody, { messageId: liveResponseMessageId, conversationId, saveMessage: saveFinalizedTurn, }); + if (finalized) { + await announceSettledTurn?.(liveResponseMessageId); + } return true; }; const existing = await getMessages( From f7c7a1c181c1ea96fba31a08a346ebac15ab1dff Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 11:31:19 +0200 Subject: [PATCH 14/19] test: Accept the Model Spec Query on the Stopped Reply Conversation URL --- e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts index 8303f0eab1e..c797d70749e 100644 --- a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts +++ b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts @@ -191,7 +191,7 @@ test.describe('compaction abort finalize', () => { const stop = page.getByRole('button', { name: 'Stop generating' }); await stop.click({ timeout: 10_000 }); await expect(stop).toBeHidden({ timeout: 20_000 }); - await expect(page).toHaveURL(/\/c\/[0-9a-fA-F-]{36}$/, { timeout: 15_000 }); + await expect(page).toHaveURL(/\/c\/[0-9a-fA-F-]{36}(\?|$)/, { timeout: 15_000 }); const conversationId = new URL(page.url()).pathname.replace('/c/', ''); try { From 1bc8d2046ec1539f96b37c9b73ae8d74a34c6542 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:08:47 +0200 Subject: [PATCH 15/19] fix: Report the Anchor Read Failure Beside a Withheld Compaction Abort --- .../compaction-abort-finalize.spec.ts | 8 +++++++ packages/api/src/agents/compaction.spec.ts | 18 +++++++++++++++ packages/api/src/agents/compaction.ts | 23 +++++++++++++++---- 3 files changed, 45 insertions(+), 4 deletions(-) diff --git a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts index c797d70749e..b7ed81d7a26 100644 --- a/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts +++ b/e2e/specs/mock/scenarios/compaction-abort-finalize.spec.ts @@ -104,6 +104,14 @@ test.describe('compaction abort finalize', () => { const compactionId = (compaction as Row | null)?.messageId as string; await expect(stop).toBeHidden({ timeout: 20_000 }); + /* The live row comes from the stopped run's final event, before any + reload: it must not present the stopped compaction as a reply that + was cut short (the notice renders 250ms after the run settles). */ + const liveRow = page.locator(`[id="${compactionId}"]`); + await expect(liveRow).toBeVisible(); + await page.waitForTimeout(1_000); + await expect(liveRow.getByText('This response stopped before it finished')).toHaveCount(0); + /* A live snapshot's shape would leave the turn reading as still running. */ const settled = await findRow({ conversationId, messageId: compactionId }); expect(settled?.unfinished).not.toBe(true); diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 8db4d4ebefc..284329fb2c8 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -1107,6 +1107,24 @@ describe('resolveAbortedTurnPersistence', () => { expect(plan.persistenceErrors[0].message).toContain('anchor unavailable'); }); + /** The outage itself reaches the caller's error boundary, not only the + * synthetic withheld-turn reason an absent anchor also produces. */ + it('reports the anchor read failure beside the withheld turn', async () => { + const outage = new Error('mongo unavailable'); + + const plan = await resolveAbortedTurnPersistence(jobData, true, { + userId: 'user-1', + getMessages: jest.fn(async () => { + throw outage; + }), + }); + + expect(plan.writeResponseRow).toBe(false); + expect(plan.withholdFinal).toBe(true); + expect(plan.persistenceErrors[0]).toBe(outage); + expect(plan.persistenceErrors[1].message).toContain('anchor unavailable'); + }); + it('keeps an ordinary stopped reply unfinished without reading its anchor', async () => { const getMessages = jest.fn(async () => []); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 05fe21c6e10..95a1c9b27ce 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -386,13 +386,28 @@ export async function resolveAbortedTurnPersistence( ) => Promise; }, ): Promise { + /** A failed read still resolves to skip-turn, so cleanup runs; the failure + * itself is reported beside the withheld turn, keeping an outage + * distinguishable from an absent anchor at the caller's error boundary. */ + let readError: Error | undefined; const anchorDecision = await resolveAbortedTurnAnchorDecision(jobData, { - messageExists: async (messageId, conversationId) => - (await getMessages({ user: userId, messageId, conversationId }, '_id')).length > 0, + messageExists: async (messageId, conversationId) => { + try { + return (await getMessages({ user: userId, messageId, conversationId }, '_id')).length > 0; + } catch (error) { + readError = error instanceof Error ? error : new Error(String(error)); + throw error; + } + }, }); const plan = planAbortedTurnPersistence(anchorDecision, shouldPersistAbortedTurn); - const persistenceErrors = - plan.withholdFinal && plan.withholdReason ? [new Error(plan.withholdReason)] : []; + const persistenceErrors: Error[] = []; + if (readError != null) { + persistenceErrors.push(readError); + } + if (plan.withholdFinal && plan.withholdReason) { + persistenceErrors.push(new Error(plan.withholdReason)); + } return { ...plan, persistenceErrors }; } From dd04da27fb24f13ffe8068865be0e92f1ad6a422 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 14:43:27 +0200 Subject: [PATCH 16/19] fix: Synthesize a Compaction Failure Only for a Round After the Latest Outcome --- packages/api/src/agents/compaction.spec.ts | 15 +++++++++++++++ packages/api/src/agents/compaction.ts | 6 +++++- 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 284329fb2c8..00d94b48285 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -425,6 +425,21 @@ describe('markAbortedCompactionContent', () => { ]); }); + /** Checkpoints are last-summary-wins: an earlier round that never + * finished is superseded by the later usable summary, not a failure. */ + it('synthesizes no failure when a later round completed after an empty one', () => { + const parts = [emptySummaryPlaceholder(), completedSummary('The later checkpoint.')]; + + markAbortedCompactionContent(parts, true); + + expect(parts).toEqual([ + expect.objectContaining({ + type: ContentTypes.SUMMARY, + initiatedBy: 'user', + }), + ]); + }); + /** A run stopped before any part streamed still needs an identifiable row: * an empty one reads as an answer to the message it hangs off. */ it('records the typed failure when nothing streamed before the stop', () => { diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 95a1c9b27ce..9afb7fda1e8 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -257,7 +257,11 @@ export function markAbortedCompactionContent( continue; } contentParts.splice(index, 1); - removedUnfinishedRound = true; + /** Walking backwards, an outcome already seen belongs to a later round: + * this placeholder is an earlier round the later one superseded. */ + if (!hasOutcome) { + removedUnfinishedRound = true; + } } /** An earlier round's checkpoint is not this round's outcome: a round the * run opened but never finished still records the typed failure beside it, From 9cbc80d090ec233d09a33c5c40b0db4eded7b0f8 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Fri, 2 Oct 2026 14:59:37 +0200 Subject: [PATCH 17/19] fix: Skip the Anchor Read When a Stopped Turn Writes No Row --- packages/api/src/agents/compaction.spec.ts | 16 ++++++++++++++++ packages/api/src/agents/compaction.ts | 5 +++++ 2 files changed, 21 insertions(+) diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index 00d94b48285..d53e07cf3df 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -1140,6 +1140,22 @@ describe('resolveAbortedTurnPersistence', () => { expect(plan.persistenceErrors[1].message).toContain('anchor unavailable'); }); + it('reads no anchor and reports nothing when the abort writes no row', async () => { + const getMessages = jest.fn(async () => { + throw new Error('mongo unavailable'); + }); + + const plan = await resolveAbortedTurnPersistence(jobData, false, { getMessages }); + + expect(getMessages).not.toHaveBeenCalled(); + expect(plan).toMatchObject({ + writeUserRow: false, + writeResponseRow: false, + withholdFinal: false, + persistenceErrors: [], + }); + }); + it('keeps an ordinary stopped reply unfinished without reading its anchor', async () => { const getMessages = jest.fn(async () => []); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 9afb7fda1e8..bc2f787be00 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -390,6 +390,11 @@ export async function resolveAbortedTurnPersistence( ) => Promise; }, ): Promise { + /** No row to write means no anchor to verify: the read is skipped, so an + * outage cannot turn an early abort's own FINAL into a reconciliation. */ + if (!shouldPersistAbortedTurn) { + return { ...planAbortedTurnPersistence('persist', false), persistenceErrors: [] }; + } /** A failed read still resolves to skip-turn, so cleanup runs; the failure * itself is reported beside the withheld turn, keeping an outage * distinguishable from an absent anchor at the caller's error boundary. */ From f4c7a3b486b4339509b204fdd0abb7416936c28e Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Sat, 3 Oct 2026 15:35:09 +0200 Subject: [PATCH 18/19] fix: Decide the Disconnect Snapshot Gate in packages/api The last-subscriber partial save now asks resolveDisconnectSnapshotMode for the settled-job decision instead of branching on the job record in the CJS controller. --- .../request.partialDisconnect.spec.js | 3 ++- .../__tests__/request.resumeMetadata.spec.js | 3 ++- api/server/controllers/agents/request.js | 8 ++----- packages/api/src/agents/compaction.spec.ts | 15 +++++++++++++ packages/api/src/agents/compaction.ts | 21 +++++++++++++++++++ 5 files changed, 42 insertions(+), 8 deletions(-) diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index 374f3e64df0..85d5ec045e1 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -57,7 +57,8 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), - isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), + resolveDisconnectSnapshotMode: (...args) => + jest.requireActual('@librechat/api').resolveDisconnectSnapshotMode(...args), settleExistingRowsBeforeErrorTurn: (...args) => jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index c4aa1af0e73..1d0790a95aa 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -269,7 +269,8 @@ jest.mock('@librechat/api', () => ({ resolveResumableRetention: jest.requireActual('@librechat/api').resolveResumableRetention, markAbortedCompactionContent: (...args) => jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), - isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), + resolveDisconnectSnapshotMode: (...args) => + jest.requireActual('@librechat/api').resolveDisconnectSnapshotMode(...args), settleExistingRowsBeforeErrorTurn: (...args) => jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 06adc3e7b4b..8eb3bbd2c9d 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -64,7 +64,7 @@ const { announceReply, announceErrorTurn, settleExistingRowsBeforeErrorTurn, - isSettledJobRecord, + resolveDisconnectSnapshotMode, markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); @@ -1958,11 +1958,7 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit * record is the source, since the client-facing resume snapshot never * carries server-private state. */ const contextMeta = jobRecord?.createdAt === jobCreatedAt ? jobRecord.contextMeta : undefined; - /** A job whose settling path (completion, error, abort) owns the final - * row must not have it reopened as an unfinished snapshot here; the - * guard reads the same record, so the window is the settling path's - * own commit span. */ - if (isSettledJobRecord(jobRecord, jobCreatedAt)) { + if (resolveDisconnectSnapshotMode(jobRecord, jobCreatedAt) === 'skip') { logger.debug('[ResumableAgentController] Skipping partial response save for a settled job'); return; } diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index d53e07cf3df..c212dff52d4 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -19,6 +19,7 @@ import { markCompactionOutcome, persistFinalizedCompactionTurn, isSettledJobRecord, + resolveDisconnectSnapshotMode, planAbortedTurnPersistence, resolveAbortedTurnPersistence, resolveAbortedTurnAnchorDecision, @@ -1292,3 +1293,17 @@ describe('isSettledJobRecord', () => { expect(isSettledJobRecord({ createdAt: 1000, status: 'error' }, 1000)).toBe(true); }); }); + +describe('resolveDisconnectSnapshotMode', () => { + it.each(['complete', 'error', 'aborted'])('withholds the snapshot of a %s job', (status) => { + expect(resolveDisconnectSnapshotMode({ createdAt: 1000, status }, 1000)).toBe('skip'); + }); + + it('writes the snapshot for a live, missing, or other-epoch record', () => { + expect(resolveDisconnectSnapshotMode({ createdAt: 1000, status: 'running' }, 1000)).toBe( + 'live', + ); + expect(resolveDisconnectSnapshotMode(null, 1000)).toBe('live'); + expect(resolveDisconnectSnapshotMode({ createdAt: 2000, status: 'error' }, 1000)).toBe('live'); + }); +}); diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index bc2f787be00..3e7cb3ef424 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -291,6 +291,27 @@ export function isSettledJobRecord( ); } +/** How the last-subscriber disconnect may persist this turn's snapshot. */ +export type DisconnectSnapshotMode = + /** The run is still live: the snapshot is written as the fallback row. */ + | 'live' + /** A settling path owns the final row: the snapshot is withheld so it + * cannot reopen the settled turn as an unfinished response. */ + | 'skip'; + +/** + * How the last-subscriber disconnect may persist this turn's snapshot, read + * from the same-epoch job record the caller already loaded. The guard reads + * the record the settling path writes, so the remaining window is that + * path's own commit span. + */ +export function resolveDisconnectSnapshotMode( + jobRecord: { createdAt?: number; status?: string } | null | undefined, + jobCreatedAt?: number, +): DisconnectSnapshotMode { + return isSettledJobRecord(jobRecord, jobCreatedAt) ? 'skip' : 'live'; +} + /** How the abort route persists a stopped turn's prerequisite rows. */ export type AbortAnchorDecision = 'persist' | 'skip-anchor' | 'skip-turn'; From b8f118c803629fdc0af81af56689fe55f8d84f16 Mon Sep 17 00:00:00 2001 From: Marco Beretta <81851188+berry-13@users.noreply.github.com> Date: Sun, 4 Oct 2026 20:05:38 +0200 Subject: [PATCH 19/19] fix: Leave an Already Settled Compaction Row to the Path That Settled It --- packages/api/src/agents/compaction.spec.ts | 15 +++++++++++++++ packages/api/src/agents/compaction.ts | 5 +++++ 2 files changed, 20 insertions(+) diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index c212dff52d4..fa5bc41c02a 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -736,6 +736,21 @@ describe('resolveFinalizedCompactionTurn', () => { expect(resolveFinalizedCompactionTurn(snapshot, { compact: true })).toEqual({ write: true }); }); + /** A stopped compaction is written settled by the abort route, and an + * error row is settled by its own write: a late failure leaves both. */ + it.each([ + ['a settled error row', [{ type: ContentTypes.ERROR, error: 'Summarization failed' }], true], + [ + 'a settled stopped row', + [{ type: ContentTypes.SUMMARY, content: [], summarizing: true, initiatedBy: 'user' }], + false, + ], + ])('leaves %s untouched', (_label, content, error) => { + const row = { unfinished: false, error, content: content as TMessageContentParts[] }; + + expect(resolveFinalizedCompactionTurn(row, { compact: true })).toEqual({ write: false }); + }); + it('leaves an already-settled checkpoint row untouched', () => { const row = { unfinished: false, diff --git a/packages/api/src/agents/compaction.ts b/packages/api/src/agents/compaction.ts index 3e7cb3ef424..529b49754d5 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -593,6 +593,11 @@ export function resolveFinalizedCompactionTurn( if (requestBody?.compact !== true) { return { write: false }; } + /** A settled row belongs to the path that settled it, whatever its + * parts hold: a late failure must not rewrite or re-announce it. */ + if (partialRow?.unfinished === false) { + return { write: false }; + } const content = Array.isArray(partialRow?.content) ? (partialRow.content as TMessageContentParts[]) : [];