diff --git a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js index 43e0f23daa6..27774e71869 100644 --- a/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js +++ b/api/server/controllers/agents/__tests__/request.partialDisconnect.spec.js @@ -49,6 +49,15 @@ jest.mock('@librechat/data-schemas', () => ({ jest.mock('@librechat/api', () => ({ getAgentErrorMetadata: (...args) => jest.requireActual('@librechat/api').getAgentErrorMetadata(...args), + markAbortedCompactionContent: (...args) => + jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), + resolveDisconnectSnapshotMode: (...args) => + jest.requireActual('@librechat/api').resolveDisconnectSnapshotMode(...args), + resolveReconciledSnapshotEnvelope: (...args) => + jest.requireActual('@librechat/api').resolveReconciledSnapshotEnvelope(...args), + settleExistingRowsBeforeErrorTurn: (...args) => + jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), isScheduleFireRequest: jest.fn(() => false), exemptFromConcurrencyLimiter: jest.fn(() => false), @@ -163,6 +172,7 @@ describe('ResumableAgentController tenant context', () => { const firePartialDisconnect = async ( user, jobRecord = { createdAt: 1000, contextMeta: partialContextMeta }, + { body = {}, aggregatedContent = [{ type: 'text', text: 'Partial response' }] } = {}, ) => { let allSubscribersLeftHandler; mockGenerationJobManager.getJobStore.mockReturnValue({ @@ -210,6 +220,7 @@ describe('ResumableAgentController tenant context', () => { endpoint: 'agents', modelOptions: { model: 'gpt-4.1' }, }, + ...body, }, config: {}, }; @@ -222,7 +233,7 @@ describe('ResumableAgentController tenant context', () => { await AgentController(req, res, jest.fn(), initializeClient, null); expect(allSubscribersLeftHandler).toEqual(expect.any(Function)); - await allSubscribersLeftHandler([{ type: 'text', text: 'Partial response' }]); + await allSubscribersLeftHandler(aggregatedContent); return tenantSeenBySave; }; @@ -268,4 +279,180 @@ describe('ResumableAgentController tenant context', () => { expect(tenantSeenBySave).toBeUndefined(); expect(mockSaveMessage).toHaveBeenCalledTimes(1); }); + + /** A cancelled compaction's partial row is built here, not by sendCompletion, + * so it carries no marker unless the disconnect path stamps one: without it + * the row reads as an answer to the message it hangs off and keeps that + * message's rerun controls. */ + it('stamps a partial response saved on disconnect with the compaction identity', async () => { + await firePartialDisconnect( + { id: 'user-123' }, + { createdAt: 1000 }, + { + body: { compact: true }, + aggregatedContent: [ + { + type: 'summary', + content: [{ type: 'text', text: 'Half a summary' }], + summarizing: true, + }, + ], + }, + ); + + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ + messageId: 'response-message', + unfinished: true, + error: false, + content: [{ type: 'summary', summarizing: true, initiatedBy: 'user' }], + }); + }); + + /** The disconnect save runs while the generation is still live and the + * completing run overwrites the row, so it must not report a failure that + * has not happened: no typed failure is invented for a compaction whose + * snapshot carries no summary or error part. */ + it('saves a non-outcome compaction partial on disconnect without a synthesized failure', async () => { + await firePartialDisconnect( + { id: 'user-123' }, + { createdAt: 1000 }, + { + body: { compact: true }, + aggregatedContent: [{ type: 'think', think: 'Picking what to summarize' }], + }, + ); + + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ + unfinished: true, + error: false, + content: [{ type: 'think', think: 'Picking what to summarize' }], + }); + expect(savedMessage.content).toHaveLength(1); + }); + /** The settling path (completion, error, abort) owns the compaction's final + * row: a disconnect snapshot landing after it would reopen the settled + * turn as an unfinished response. */ + it('skips the compaction partial save when the job record has settled', async () => { + await firePartialDisconnect( + { id: 'user-123' }, + { createdAt: 1000, status: 'error' }, + { + body: { compact: true }, + aggregatedContent: [{ type: 'text', text: 'Partial response' }], + }, + ); + + expect(mockSaveMessage).not.toHaveBeenCalled(); + }); + + /** Ordinary turns keep the pre-change behavior exactly: their snapshot is + * the fallback row even when the job record has settled, because the + * terminal row write may still fail. */ + it('still persists an ordinary 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).toHaveBeenCalledTimes(1); + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ unfinished: true, error: false }); + }); + /** The terminal write failed and settled for a reconciliation frame: the + * snapshot is the turn's only row, so it persists with the terminal + * outcome (the typed failure for content with no summary) and envelope. */ + it('promotes a reconciled compaction snapshot to the terminal row', async () => { + mockGetMessages.mockResolvedValue([{ _id: 'existing-anchor' }]); + await firePartialDisconnect( + { id: 'user-123' }, + { + createdAt: 1000, + status: 'error', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }, + { + body: { compact: true }, + aggregatedContent: [{ type: 'think', think: 'Picking what to summarize' }], + }, + ); + + expect(mockSaveMessage).toHaveBeenCalledTimes(1); + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ unfinished: false, error: true }); + expect(savedMessage.content).toEqual([ + { type: 'think', think: 'Picking what to summarize' }, + expect.objectContaining({ type: 'error', initiatedBy: 'user' }), + ]); + }); + /** Without a persisted anchor the promotion is withheld: the fallback must + * not recreate the orphan the absent-anchor abort suppressed. */ + it('withholds a reconciled snapshot whose anchor was never persisted', async () => { + await firePartialDisconnect( + { id: 'user-123' }, + { + createdAt: 1000, + status: 'error', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }, + { + body: { compact: true }, + aggregatedContent: [{ type: 'text', text: 'Partial response' }], + }, + ); + + expect(mockSaveMessage).not.toHaveBeenCalled(); + }); + + /** An empty snapshot on a reconciled compaction still records the turn: + * the terminal synthesis appends the typed failure from nothing. */ + it('synthesizes the terminal outcome for an empty reconciled compaction snapshot', async () => { + mockGetMessages.mockResolvedValue([{ _id: 'existing-anchor' }]); + await firePartialDisconnect( + { id: 'user-123' }, + { + createdAt: 1000, + status: 'error', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }, + { body: { compact: true }, aggregatedContent: [] }, + ); + + expect(mockSaveMessage).toHaveBeenCalledTimes(1); + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ unfinished: false, error: true }); + expect(savedMessage.content).toEqual([ + expect.objectContaining({ type: 'error', initiatedBy: 'user' }), + ]); + }); + + /** A completed claim's promoted snapshot settles as a finished row, not an + * error. */ + it('settles a completed reconciled snapshot as a finished row', async () => { + mockGetMessages.mockResolvedValue([{ _id: 'existing-anchor' }]); + await firePartialDisconnect( + { id: 'user-123' }, + { + createdAt: 1000, + status: 'complete', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }, + { + body: { compact: true }, + aggregatedContent: [ + { + type: 'summary', + content: [{ type: 'text', text: 'A finished checkpoint.' }], + boundary: { messageId: 'm', contentIndex: 0 }, + }, + ], + }, + ); + + expect(mockSaveMessage).toHaveBeenCalledTimes(1); + const [, savedMessage] = mockSaveMessage.mock.calls[0]; + expect(savedMessage).toMatchObject({ unfinished: false, error: false }); + }); }); diff --git a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js index ef482d04aec..015d1d5e01f 100644 --- a/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js +++ b/api/server/controllers/agents/__tests__/request.resumeMetadata.spec.js @@ -257,6 +257,15 @@ jest.mock('@librechat/api', () => ({ ).getSteerRecoveryFailure, getAgentErrorMetadata: (...args) => jest.requireActual('@librechat/api').getAgentErrorMetadata(...args), + markAbortedCompactionContent: (...args) => + jest.requireActual('@librechat/api').markAbortedCompactionContent(...args), + isSettledJobRecord: (...args) => jest.requireActual('@librechat/api').isSettledJobRecord(...args), + resolveDisconnectSnapshotMode: (...args) => + jest.requireActual('@librechat/api').resolveDisconnectSnapshotMode(...args), + resolveReconciledSnapshotEnvelope: (...args) => + jest.requireActual('@librechat/api').resolveReconciledSnapshotEnvelope(...args), + settleExistingRowsBeforeErrorTurn: (...args) => + jest.requireActual('@librechat/api').settleExistingRowsBeforeErrorTurn(...args), sendEvent: jest.fn(), logAgentMemorySnapshot: jest.fn(), isScheduleFireRequest: (...args) => mockIsScheduleFireRequest(...args), diff --git a/api/server/controllers/agents/request.js b/api/server/controllers/agents/request.js index 7135ca4113d..b851444349f 100644 --- a/api/server/controllers/agents/request.js +++ b/api/server/controllers/agents/request.js @@ -52,6 +52,10 @@ const { resolvePersistableCodeEnvironmentDecision, getFailedTurnTraceFields, resolveFailedTurnContent, + settleExistingRowsBeforeErrorTurn, + resolveDisconnectSnapshotMode, + resolveReconciledSnapshotEnvelope, + markAbortedCompactionContent, } = require('@librechat/api'); const { disposeClient } = require('~/server/cleanup'); const { @@ -426,23 +430,6 @@ async function saveErrorTurn( } const userId = req.user.id; - const existing = await getMessages( - { user: userId, messageId: errorMessageId, conversationId }, - '_id', - ); - if (existing.length > 0) { - return; - } - if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { - const partial = await getMessages( - { user: userId, messageId: liveResponseMessageId, conversationId }, - '_id', - ); - if (partial.length > 0) { - return; - } - } - const reqCtx = { userId, isTemporary: @@ -453,6 +440,29 @@ 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 settlement = 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 (settlement.covered) { + return; + } + /** The anchor-shaped collision redirects the error row to the failed + * run's own response id, so it can never overwrite the anchor. */ + if (settlement.errorRowMessageId != null) { + errorMessageId = settlement.errorRowMessageId; + } + const context = 'api/server/controllers/agents/request.js - failed turn'; const endpoint = endpointOption?.endpoint; const model = getAgentResponseModel(req, endpointOption); @@ -1663,6 +1673,9 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit * no user message of its own, the response parented onto an * existing message. A reconnecting client rebuilds it that way. */ ...((isRegenerate || isCompaction) && { isRegenerate: true }), + /** The job record is where the abort paths learn the turn was a + * compaction: they run after the request that created the job. */ + ...(isCompaction && { compact: true }), ...(scheduleId ? { scheduleId, @@ -1849,13 +1862,10 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit * overwrite this with the complete response using the same messageId pattern. */ job.emitter.on('allSubscribersLeft', async (aggregatedContent) => { - if (partialResponseSaved || !aggregatedContent || aggregatedContent.length === 0) { - return; - } - - const persistableContent = filterPersistableAbortContent(aggregatedContent); - if (persistableContent.length === 0) { - logger.debug('[ResumableAgentController] No persistable content to save partial response'); + /** Empty content is rejected only after the snapshot mode is known: a + * reconciled compaction synthesizes its terminal outcome from nothing, + * while a live run has nothing persistable. */ + if (partialResponseSaved || !aggregatedContent) { return; } @@ -1868,6 +1878,47 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit return; } + /** How this snapshot may persist is decided in @librechat/api: live + * runs keep the marker-only shape, a failed terminal write (settled + * for a reconciliation frame) promotes the snapshot to the turn's + * terminal row, and a durably settled compaction withholds it. */ + const snapshotMode = await resolveDisconnectSnapshotMode( + isCompaction, + jobRecord, + jobCreatedAt, + { + anchorExists: async () => + ( + await getMessages( + { + user: userId, + messageId: resumeState.userMessage.messageId, + conversationId, + }, + '_id', + ) + ).length > 0, + }, + ); + if (snapshotMode === 'skip') { + logger.debug( + '[ResumableAgentController] Skipping compaction partial save for a settled job', + ); + return; + } + if (snapshotMode !== 'terminal' && aggregatedContent.length === 0) { + return; + } + const persistableContent = markAbortedCompactionContent( + filterPersistableAbortContent(aggregatedContent), + isCompaction, + { synthesizeFailure: snapshotMode === 'terminal' }, + ); + if (persistableContent.length === 0) { + logger.debug('[ResumableAgentController] No persistable content to save partial response'); + return; + } + partialResponseSaved = true; const responseConversationId = resumeState.conversationId || conversationId; /** The run publishes its calibration and fading tiers onto the job; a @@ -1885,8 +1936,12 @@ const ResumableAgentController = async (req, res, next, initializeClient, addTit parentMessageId: resumeState.userMessage.messageId, sender: client?.sender ?? 'AI', content: persistableContent, - unfinished: true, - error: false, + /** A snapshot promoted to the terminal row settles with the + * envelope its reconciled claim's status dictates; a live-run + * snapshot keeps the live shape. */ + ...(snapshotMode === 'terminal' + ? resolveReconciledSnapshotEnvelope(jobRecord?.status) + : { unfinished: true, error: false }), isCreatedByUser: false, user: userId, endpoint: endpointOption.endpoint, diff --git a/api/server/routes/agents/__tests__/abort.spec.js b/api/server/routes/agents/__tests__/abort.spec.js index e31bdc95c7d..2682e75d878 100644 --- a/api/server/routes/agents/__tests__/abort.spec.js +++ b/api/server/routes/agents/__tests__/abort.spec.js @@ -24,6 +24,7 @@ const mockGenerationJobManager = { }; const mockSaveMessage = jest.fn(); +const mockGetMessages = jest.fn(async () => [{ _id: 'existing-anchor' }]); const mockRecordScheduleOutcome = jest.fn(); const mockBeginScheduledStop = jest.fn(); @@ -47,6 +48,7 @@ jest.mock('@librechat/api', () => ({ jest.mock('~/models', () => ({ saveMessage: (...args) => mockSaveMessage(...args), + getMessages: (...args) => mockGetMessages(...args), })); jest.mock('~/server/services/Schedules', () => ({ @@ -358,6 +360,129 @@ describe('Agent Abort Endpoint', () => { }), ).resolves.toBe(false); }); + + /** A compaction's `userMessage` is the persisted leaf projected for + * identity only (`projectCompactionAnchor`), so upserting it would + * erase a user leaf's text or turn an assistant leaf into an empty + * user row: Stop writes only the aborted response. */ + it('skips the anchor prerequisite when an aborted compaction persists its row', async () => { + const jobStreamId = 'test-stream-compact'; + const anchorId = 'persisted-leaf-1'; + const compactionRowId = 'compaction-response-1'; + + mockGenerationJobManager.getJob.mockResolvedValue({ + metadata: { userId: 'test-user-123', generationProtocolVersion: 2 }, + }); + + const abortResult = { + success: true, + jobData: { + compact: true, + createdEventEmitted: true, + userMessage: { + messageId: anchorId, + parentMessageId: 'older-response', + conversationId: jobStreamId, + text: '', + }, + responseMessageId: compactionRowId, + 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).toHaveBeenCalledTimes(1); + expect(mockSaveMessage).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + messageId: compactionRowId, + parentMessageId: anchorId, + unfinished: true, + isCreatedByUser: false, + }), + 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: '', + }; + let beforePublishError; + mockGenerationJobManager.abortJob.mockImplementation(async (_streamId, options) => { + 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; + }); + + 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(); + expect(beforePublishError).toBeInstanceOf(Error); + expect(beforePublishError.message).toContain('anchor unavailable'); + }); }); describe('Partial Response Saving', () => { diff --git a/api/server/routes/agents/index.js b/api/server/routes/agents/index.js index 52a4c34eeb1..1ce203864d6 100644 --- a/api/server/routes/agents/index.js +++ b/api/server/routes/agents/index.js @@ -5,6 +5,8 @@ const { GenerationJobManager, TERMINAL_PUBLICATION_RECONNECT_ERROR, hasPersistableAbortContent, + resolveAbortedTurnAnchorDecision, + planAbortedTurnPersistence, buildAbortedResponseMetadata, isPendingActionStale, toClientPendingAction, @@ -49,7 +51,7 @@ const { getServerGenerationProtocol, negotiateExistingGenerationProtocol, } = require('~/server/controllers/agents/protocol'); -const { getFiles, saveMessage } = require('~/models'); +const { getFiles, getMessages, saveMessage } = require('~/models'); const { recordScheduleOutcome, beginScheduledStop, @@ -764,11 +766,27 @@ router.post('/chat/abort', configMiddleware, async (req, res, next) => { * its parent and the preliminary-parent fence correctly rejects it. */ const shouldPersistAbortedTurn = hasPersistableAbortContent(content) || jobData?.createdEventEmitted === true; + /** 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 }, '_id').then( + (rows) => rows.length > 0, + ), + }), + shouldPersistAbortedTurn, + ); + if (abortPersistencePlan.withholdFinal && abortPersistencePlan.withholdReason) { + persistenceErrors.push(new Error(abortPersistencePlan.withholdReason)); + } if ( jobData?.userMessage?.messageId && jobData?.responseMessageId && - shouldPersistAbortedTurn + abortPersistencePlan.writeResponseRow ) { const messageContext = { userId: req?.user?.id, @@ -827,16 +845,19 @@ router.post('/chat/abort', configMiddleware, async (req, res, next) => { * with neither row stored. Both writes are idempotent upserts; * await the user prerequisite first, but still attempt the child * write and checkpoint cleanup so every independently useful - * operation gets a chance to succeed. */ - try { - const persistedRequest = await saveMessage(messageContext, requestMessage, { - context: 'api/server/routes/agents/index.js - abort user prerequisite', - }); - if (!persistedRequest) { - throw new Error('Abort user prerequisite was not persisted'); + * operation gets a chance to succeed. A compaction skips the + * prerequisite: its anchor is the persisted leaf itself. */ + if (abortPersistencePlan.writeUserRow) { + try { + const persistedRequest = await saveMessage(messageContext, requestMessage, { + context: 'api/server/routes/agents/index.js - abort user prerequisite', + }); + if (!persistedRequest) { + throw new Error('Abort user prerequisite was not persisted'); + } + } catch (error) { + persistenceErrors.push(error); } - } catch (error) { - persistenceErrors.push(error); } try { diff --git a/e2e/specs/mock/scenarios/compaction-rerun-controls.spec.ts b/e2e/specs/mock/scenarios/compaction-rerun-controls.spec.ts index ef0bb10477a..d66c18b697e 100644 --- a/e2e/specs/mock/scenarios/compaction-rerun-controls.spec.ts +++ b/e2e/specs/mock/scenarios/compaction-rerun-controls.spec.ts @@ -117,6 +117,78 @@ async function compactWithEmptySummarizer(page: Page, request: APIRequestContext return { conversationId, compactionId: compactionId as string }; } +/** + * A real cancelled compaction: the fixture summarizer is held silent so the run + * is still summarizing when the composer's Stop fires, and the branch ends in + * a user message, the shape whose Regenerate would otherwise answer that user + * turn instead of redoing the compaction. Returns the conversation and the + * turn the cancelled compaction persisted. + * + * The conversation is created through the composer (a seeded one carries no + * model, so its compact submission never starts) and the user leaf is seeded + * onto the answer: seeded rows are newer, so the leaf reads as latest. + */ +async function cancelCompactionOnUserLeaf(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 answerId = await withMongo(async (db) => { + const row = await db + .collection('messages') + .findOne({ conversationId, isCreatedByUser: false }, { sort: { createdAt: -1 } }); + return row?.messageId as string | undefined; + }); + expect(answerId).toBeTruthy(); + + const leafUserId = randomUUID(); + await seedMessages(userEmail, conversationId, [ + { + messageId: leafUserId, + parentMessageId: answerId as string, + text: `Compact this before answering ${label}`, + 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(`Compact this before answering ${label}`), + ).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 }); + await stop.click(); + + /** The Stop route persists the aborted turn before it publishes the final + * event, so the row is in storage once the stop settles. */ + let compactionId: string | undefined; + await expect + .poll( + () => + withMongo(async (db) => { + const row = await db.collection('messages').findOne({ + conversationId, + parentMessageId: leafUserId, + isCreatedByUser: false, + }); + compactionId = row?.messageId as string | undefined; + return compactionId != null; + }), + { timeout: 20_000 }, + ) + .toBeTruthy(); + return { conversationId, compactionId: compactionId as string }; +} + test.describe('compaction rerun controls', () => { /** `compactWithEmptySummarizer` switches the shared fixture summarizer to * blank output before it returns, so a failure inside it would leave every @@ -315,6 +387,61 @@ test.describe('compaction rerun controls', () => { } }); + /* A real cancelled run, not a seeded row: the abort path owns the stopped + turn, so the marker has to be stamped where the aborted content is + assembled or the row keeps the user turn's rerun controls live. */ + test('a cancelled compaction on a user leaf offers no rerun controls @scenario:cancelled-compaction-on-user-turn-offers-no-rerun-controls', async ({ + page, + request, + }) => { + const { conversationId, compactionId } = await cancelCompactionOnUserLeaf( + page, + request, + 'cancelled-live', + ); + try { + const row = page.locator(`[id="${compactionId}"]`); + await expect(row).toBeVisible(); + await row.hover(); + + /* Replaying the user turn behind it would answer that message again rather + than redo the compaction, so the marker withholds the controls here too. + The stored outcome text is the reload scenario's to check: the live row + renders the turn it streamed, which stopped before its first delta. */ + await expect(page.locator(`[id="edit-${compactionId}"]`)).toHaveCount(0); + await expect(page.getByTestId('regenerate-generation-button')).toHaveCount(0); + await expect(page.getByTestId('continue-generation-button')).toHaveCount(0); + /* The redo path a compaction keeps is the indicator's own action. */ + await page.getByTestId('token-usage').click(); + await expect(page.getByRole('button', { name: 'Compact context' })).toBeEnabled(); + } finally { + await cleanup(conversationId); + } + }); + + /* The same cancelled turn read back from storage: the row the abort path + persisted has to carry the marker a reload rebuilds it from. */ + test('a cancelled compaction on a user leaf stays free of rerun controls after a reload @scenario:cancelled-compaction-on-user-turn-survives-reload-without-rerun-controls', async ({ + page, + request, + }) => { + const { conversationId, compactionId } = await cancelCompactionOnUserLeaf( + page, + request, + 'cancelled-reload', + ); + try { + const row = await openRow(page, conversationId, compactionId); + + await expect(row.getByText('Could not compact the context', { exact: false })).toBeVisible(); + await expect(page.locator(`[id="edit-${compactionId}"]`)).toHaveCount(0); + await expect(page.getByTestId('regenerate-generation-button')).toHaveCount(0); + await expect(page.getByTestId('continue-generation-button')).toHaveCount(0); + } finally { + await cleanup(conversationId); + } + }); + test('a turn that only auto-summarized keeps its rerun controls @scenario:auto-summarized-turn-keeps-rerun-controls', async ({ page, }) => { diff --git a/packages/api/src/agents/compaction.spec.ts b/packages/api/src/agents/compaction.spec.ts index e5460ca7f77..351d4963733 100644 --- a/packages/api/src/agents/compaction.spec.ts +++ b/packages/api/src/agents/compaction.spec.ts @@ -15,8 +15,17 @@ import { dropUnusableSummaryParts, findCheckpointSummaryPart, getSummaryPartText, + markAbortedCompactionContent, markCompactionOutcome, + persistFinalizedCompactionTurn, + isSettledJobRecord, + resolveDisconnectSnapshotMode, + resolveReconciledSnapshotEnvelope, + planAbortedTurnPersistence, + resolveAbortedTurnAnchorDecision, + settleExistingRowsBeforeErrorTurn, resolveFailedTurnContent, + resolveFinalizedCompactionTurn, restoreCompactionSemanticIndex, restoreCompactionSemanticIndexSnapshot, stripUnusableSummaryParts, @@ -246,6 +255,476 @@ describe('resolveFailedTurnContent', () => { }); }); +describe('markAbortedCompactionContent', () => { + const partialSummary = (text: string): TMessageContentParts => ({ + type: ContentTypes.SUMMARY, + /** Streamed deltas never carry a boundary; a stopped round keeps them. */ + content: [{ type: ContentTypes.TEXT, text }], + summarizing: true, + }); + + /** The summarizer opens the part when its round starts, so a stop can land + * between that and the first delta. */ + const emptySummaryPlaceholder = (): TMessageContentParts => ({ + type: ContentTypes.SUMMARY, + content: [], + summarizing: true, + }); + + const completedSummary = (text: string): TMessageContentParts => ({ + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text }], + boundary: completedBoundary, + }); + + /** The abort path owns a cancelled run's row: its partial summary must still + * carry the marker, or on a branch ending in a user message the row keeps a + * Regenerate that answers that user turn instead of redoing the compaction. + * The truncated prefix is kept but marked failed, or its label presents it + * as a finished checkpoint. */ + it('marks the partial summary a stopped compaction had streamed as failed', () => { + const parts = [partialSummary('Half a summary')]; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked).toHaveLength(1); + expect(marked[0]).toMatchObject({ initiatedBy: 'user', failed: true, summarizing: true }); + /** The aggregated parts belong to the live run: the input is untouched. */ + expect(parts[0]).not.toHaveProperty('initiatedBy'); + }); + + /** A round that finished before the Stop landed is a real checkpoint: the + * race is not a failure. */ + it('marks a summary that completed before the stop without failing it', () => { + const parts = [completedSummary('Finished before the stop.')]; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked).toHaveLength(1); + expect(marked[0]).toMatchObject({ initiatedBy: 'user' }); + expect(marked[0]).not.toHaveProperty('failed'); + }); + + /** Every part that can carry the marker gets it: the row's identity must not + * depend on which of its parts a reader inspects first. */ + it('marks an error part the stopped run had already recorded', () => { + const parts: TMessageContentParts[] = [ + partialSummary('Half a summary'), + { type: ContentTypes.ERROR, error: 'Something else failed first' }, + ]; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked[0]).toMatchObject({ initiatedBy: 'user', failed: true }); + expect(marked[1]).toMatchObject({ initiatedBy: 'user' }); + }); + + /** A summary placeholder with no text is not an outcome: nothing of the + * round survived to show, so the typed failure is the row's whole + * outcome. */ + it('replaces a summary placeholder that streamed nothing with the typed failure', () => { + const parts = [emptySummaryPlaceholder()]; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked).toEqual([ + { + type: ContentTypes.ERROR, + error: JSON.stringify({ type: ErrorTypes.COMPACTION_FAILED }), + initiatedBy: 'user', + }, + ]); + }); + + /** 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()]; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked).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', () => { + const parts: TMessageContentParts[] = []; + + const marked = markAbortedCompactionContent(parts, true); + + expect(marked).toEqual([ + { + type: ContentTypes.ERROR, + error: JSON.stringify({ type: ErrorTypes.COMPACTION_FAILED }), + initiatedBy: 'user', + }, + ]); + }); + + /** The disconnect save runs while the generation is still live and the + * completion path overwrites the row: it stamps identity and rewrites + * nothing else, not even the failure flag of a still-streaming part. */ + it('marks a non-terminal snapshot without failing or synthesizing anything', () => { + const parts: TMessageContentParts[] = [partialSummary('Half a summary')]; + + const marked = markAbortedCompactionContent(parts, true, { synthesizeFailure: false }); + + expect(marked).toEqual([{ ...partialSummary('Half a summary'), initiatedBy: 'user' }]); + expect(parts[0]).not.toHaveProperty('initiatedBy'); + }); + + it('returns content from a turn that was not a compaction unchanged', () => { + const parts = [partialSummary('An automatic detour partial')]; + + expect(markAbortedCompactionContent(parts, false)).toBe(parts); + expect(parts[0]).not.toHaveProperty('initiatedBy'); + }); +}); + +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', 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', 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', + ); + }); +}); + +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 the envelope with the marking reapplied 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, + content: [{ type: ContentTypes.ERROR, error: 'Summarization failed', 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(); + }); + + /** 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', () => { + 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, {})).toEqual({ write: false }); + }); + + /** 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({ + write: true, + 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({ + write: true, + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, + initiatedBy: 'user', + failed: true, + }, + ], + }); + }); + + /** The parts already carry the failure, but the snapshot's live-run flags + * are still unsettled: the write settles the envelope with the marking + * reapplied (idempotent on marked parts). */ + it('settles a row whose parts already carry the failure', () => { + 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 })).toEqual({ + write: true, + content: failedSummary.content, + }); + expect(resolveFinalizedCompactionTurn(recordedFailure, { compact: true })).toEqual({ + write: true, + content: recordedFailure.content, + }); + }); + + /** Legacy and imported rows can carry failure parts that predate the + * identity marker: the settle write stamps them, or the restored turn + * keeps the wrong rerun target. */ + it('stamps legacy failure parts that never carried the marker', () => { + const legacyRow = { + unfinished: true, + content: [{ type: ContentTypes.ERROR, error: 'Summarization failed' }], + }; + + expect(resolveFinalizedCompactionTurn(legacyRow, { compact: true })).toEqual({ + write: true, + content: [{ type: ContentTypes.ERROR, error: 'Summarization failed', initiatedBy: 'user' }], + }); + }); + + /** 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, + content: [ + { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'A finished checkpoint.' }], + boundary: completedBoundary, + initiatedBy: 'user', + }, + ], + }); + }); + + it('leaves an already-settled checkpoint row untouched', () => { + const row = { + unfinished: false, + 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 + * 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({ + write: true, + 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', () => { const legacySummary = { type: ContentTypes.SUMMARY, text: 'Summary of conversation' }; @@ -485,3 +964,288 @@ 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, + }); + }); + + /** 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', () => { + 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.toEqual({ + covered: 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.toEqual({ + covered: true, + }); + + expect(saved).toHaveLength(1); + expect(saved[0]).toMatchObject({ messageId: 'live-response' }); + }); + + /** The collision with no live row saved must not swallow the failure: the + * error row proceeds under the failed run's own response id, where it can + * never overwrite the anchor. */ + it('redirects the error row to the live id when only the anchor matched', async () => { + const { saved, deps: d } = deps({ + 'error-target': [{ messageId: 'error-target', _id: 'anchor-shaped-match' }], + }); + + await expect(settleExistingRowsBeforeErrorTurn({ compact: true }, d)).resolves.toEqual({ + covered: false, + errorRowMessageId: 'live-response', + }); + + expect(saved).toHaveLength(0); + }); + + /** A failure before the response id was allocated leaves nowhere safe to + * write the error row: writing it under the error id would overwrite the + * anchor, so it is withheld instead. */ + it('withholds the error row when the collision has no live id to redirect to', async () => { + const { saved, deps: d } = deps({ + 'error-target': [{ messageId: 'error-target', _id: 'anchor-shaped-match' }], + }); + const { liveResponseMessageId: _omitted, ...dWithoutLiveId } = d; + + await expect( + settleExistingRowsBeforeErrorTurn({ compact: true }, dWithoutLiveId), + ).resolves.toEqual({ covered: true }); + + expect(saved).toHaveLength(0); + }); + + 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.toEqual({ covered: 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.toEqual({ + covered: false, + }); + + expect(saved).toHaveLength(0); + }); +}); + +describe('resolveDisconnectSnapshotMode', () => { + it('keeps an ordinary turn writing its fallback row, settled or not', async () => { + await expect( + resolveDisconnectSnapshotMode(false, { createdAt: 1000, status: 'error' }, 1000), + ).resolves.toBe('live'); + await expect(resolveDisconnectSnapshotMode(false, null, undefined)).resolves.toBe('live'); + }); + + it('keeps a live compaction writing its snapshot', async () => { + await expect( + resolveDisconnectSnapshotMode(true, { createdAt: 1000, status: 'running' }, 1000), + ).resolves.toBe('live'); + }); + + it('withholds a settled compaction snapshot', async () => { + await expect( + resolveDisconnectSnapshotMode(true, { createdAt: 1000, status: 'aborted' }, 1000), + ).resolves.toBe('skip'); + }); + + const reconciled = () => ({ + createdAt: 1000, + status: 'error', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }); + + /** The terminal write failed and settled for a reconciliation frame: the + * snapshot is promoted to the turn's terminal row, because no other row + * will ever be persisted for it. */ + it('promotes a reconciled compaction snapshot to the terminal row', async () => { + await expect(resolveDisconnectSnapshotMode(true, reconciled(), 1000)).resolves.toBe('terminal'); + }); + + /** The promotion must not recreate the orphan an absent-anchor abort + * deliberately withheld: without a persisted anchor there is nothing to + * hang the terminal row on. */ + it('withholds a reconciled snapshot whose anchor was never persisted', async () => { + await expect( + resolveDisconnectSnapshotMode(true, reconciled(), 1000, { + anchorExists: async () => false, + }), + ).resolves.toBe('skip'); + }); +}); + +describe('resolveReconciledSnapshotEnvelope', () => { + it('keeps the abort row shape for an aborted claim', () => { + expect(resolveReconciledSnapshotEnvelope('aborted')).toEqual({ + unfinished: true, + error: false, + }); + }); + + it('settles a completed claim as a finished row', () => { + expect(resolveReconciledSnapshotEnvelope('complete')).toEqual({ + unfinished: false, + error: false, + }); + }); + + it('settles everything else with the error envelope', () => { + expect(resolveReconciledSnapshotEnvelope('error')).toEqual({ + unfinished: false, + error: true, + }); + expect(resolveReconciledSnapshotEnvelope(undefined)).toEqual({ + unfinished: false, + error: true, + }); + }); +}); + +describe('isSettledJobRecord', () => { + it.each(['complete', 'error', 'aborted'])('treats a %s record as settled', (status) => { + expect(isSettledJobRecord({ createdAt: 1000, status })).toBe(true); + }); + + /** The terminal claim precedes its row write, so the status alone does not + * prove the row is durable: the snapshot stays the fallback until the + * pending marker clears. */ + it('treats a record with terminal persistence still pending as unsettled', () => { + expect( + isSettledJobRecord({ createdAt: 1000, status: 'error', terminalPersistencePending: true }), + ).toBe(false); + }); + + /** A failed terminal write clears the marker while publishing a + * reconciliation frame: no row was persisted, so the streamed snapshot + * stays the turn's only fallback. */ + it('treats a record whose durable final event is a reconciliation frame as unsettled', () => { + const reconciled = { + createdAt: 1000, + status: 'aborted', + finalEvent: JSON.stringify({ final: true, reconcile: true }), + }; + const settled = { + createdAt: 1000, + status: 'aborted', + finalEvent: JSON.stringify({ final: true }), + }; + + expect(isSettledJobRecord(reconciled)).toBe(false); + expect(isSettledJobRecord(settled)).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 e89311178f6..23ded5c8216 100644 --- a/packages/api/src/agents/compaction.ts +++ b/packages/api/src/agents/compaction.ts @@ -157,6 +157,468 @@ export function resolveFailedTurnContent( return { content: compactionFailureContent(errorText) }; } +/** + * The content an aborted compaction persists: the run's stream-aggregated + * parts, carrying the marker that keeps the turn identifiable as a compaction. + * The abort path owns a cancelled run's row (a stopped turn is unfinished, not + * failed) and nothing else on that path knows the request was a compaction, so + * without this the row reads as an answer to the message it hangs off and keeps + * that message's rerun controls: on a branch ending in a user message, + * Regenerate would answer the user turn behind the compaction instead of + * redoing it. + * + * A terminal abort (Stop) settles the turn, so it applies the completed run's + * outcome rules: a usable summary is marked as the outcome; a partial one + * keeps its text but is marked `failed`, or its label would present the + * truncated prefix as a finished checkpoint; a placeholder that never streamed + * text goes, leaving the typed failure as the row's outcome. A non-terminal + * snapshot (`synthesizeFailure: false`, the disconnect save the run may still + * complete and overwrite) marks what is there and rewrites nothing else. + * Content from a turn that was not a compaction is returned unchanged, and + * the parts are never edited in place: the aggregated parts belong to the + * still-live run on the disconnect path, so every stamped part is a copy. + */ +export function markAbortedCompactionContent( + contentParts: TMessageContentParts[], + isCompaction: boolean, + { synthesizeFailure = true }: { synthesizeFailure?: boolean } = {}, +): TMessageContentParts[] { + if (!isCompaction) { + return contentParts; + } + const marked: TMessageContentParts[] = []; + let hasOutcome = false; + let removedUnfinishedRound = false; + for (const part of contentParts) { + if (part == null) { + marked.push(part); + continue; + } + if (part.type === ContentTypes.ERROR) { + marked.push({ ...part, initiatedBy: 'user' as const }); + hasOutcome = true; + continue; + } + if (part.type !== ContentTypes.SUMMARY) { + marked.push(part); + continue; + } + /** The usability predicate's false side narrows the part's type away, so + * the reference is taken before it runs. */ + const summary = part; + if (isUsableSummaryPart(part)) { + marked.push({ ...summary, initiatedBy: 'user' as const }); + hasOutcome = true; + continue; + } + if (!synthesizeFailure) { + marked.push({ ...summary, initiatedBy: 'user' as const }); + continue; + } + if (isSummaryPartWithText(summary)) { + marked.push({ ...summary, initiatedBy: 'user' as const, failed: true }); + hasOutcome = true; + continue; + } + 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, + * or the stopped turn reads as the successful compaction the checkpoint + * describes. */ + if ((!hasOutcome || removedUnfinishedRound) && synthesizeFailure) { + marked.push(...compactionFailureContent()); + } + return marked; +} + +/** Whether a job record's durable final event is a reconciliation frame: the + * conservative substitute published when the terminal row write failed, so + * no message row backs the terminal claim. */ +function hasDurableReconcileFrame(finalEvent: unknown): boolean { + if (typeof finalEvent !== 'string' || finalEvent.length === 0) { + return false; + } + try { + const parsed = JSON.parse(finalEvent) as { reconcile?: unknown } | null; + return parsed?.reconcile === true; + } catch { + return false; + } +} + +/** 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, and only one whose terminal write + * actually landed. */ +export function isSettledJobRecord( + jobRecord: + | { + createdAt?: number; + status?: string; + terminalPersistencePending?: boolean; + finalEvent?: string; + } + | null + | undefined, + jobCreatedAt?: number, +): boolean { + if (jobRecord == null || (jobCreatedAt != null && jobRecord.createdAt !== jobCreatedAt)) { + return false; + } + if (jobRecord.terminalPersistencePending === true) { + /** The terminal claim precedes its row write: the status alone does not + * prove the row is durable, and the snapshot is still the fallback if + * that write fails. */ + return false; + } + if (hasDurableReconcileFrame(jobRecord.finalEvent)) { + /** The terminal write failed and a reconciliation frame was published in + * its place: nothing was persisted for the turn, so the streamed + * snapshot remains its only row. */ + return false; + } + return ( + jobRecord.status === 'complete' || + jobRecord.status === 'error' || + jobRecord.status === 'aborted' + ); +} + +/** How a disconnect may persist this turn's snapshot. */ +export type DisconnectSnapshotMode = + /** The run is still live: the snapshot keeps the live shape. */ + | 'live' + /** The terminal write failed and settled for a reconciliation frame: the + * snapshot is the turn's only row, so it persists with the terminal + * outcome and envelope. */ + | 'terminal' + /** A settled terminal row exists: the snapshot is withheld so it cannot + * reopen the settled turn. */ + | 'skip'; + +/** The row flags a promoted terminal snapshot settles with, keyed by the + * reconciled claim's status: an aborted run keeps the abort row's shape, a + * completed run a finished row, and everything else the error envelope. */ +export function resolveReconciledSnapshotEnvelope(status: unknown): { + unfinished: boolean; + error: boolean; +} { + if (status === 'aborted') { + return { unfinished: true, error: false }; + } + if (status === 'complete') { + return { unfinished: false, error: false }; + } + return { unfinished: false, error: true }; +} + +/** + * How the last-subscriber disconnect may persist this turn's snapshot. A + * compaction whose settling path (completion, error, abort) durably owns the + * final row must not have it reopened as an unfinished snapshot; a compaction + * whose terminal write settled for a reconciliation frame has no row at all, + * so its snapshot is promoted to the terminal row, anchored on the persisted + * leaf the injected reader confirms (the promotion must not recreate the + * orphan an absent-anchor abort deliberately withheld); ordinary turns keep + * writing their fallback row exactly as before, because their terminal row + * write may still fail. + */ +export async function resolveDisconnectSnapshotMode( + isCompaction: boolean, + jobRecord: + | { + createdAt?: number; + status?: string; + terminalPersistencePending?: boolean; + finalEvent?: string; + } + | null + | undefined, + jobCreatedAt: number | undefined, + { anchorExists = async () => true }: { anchorExists?: () => Promise } = {}, +): Promise { + if (!isCompaction) { + return 'live'; + } + if (isSettledJobRecord(jobRecord, jobCreatedAt)) { + return 'skip'; + } + if (!hasDurableReconcileFrame(jobRecord?.finalEvent)) { + return 'live'; + } + return (await anchorExists()) ? 'terminal' : 'skip'; +} + +/** How the abort route persists a stopped turn's prerequisite rows. */ +export type AbortAnchorDecision = 'persist' | 'skip-anchor' | 'skip-turn'; + +/** + * 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 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'; + } + try { + const anchorExists = await messageExists(anchorId, jobData.conversationId); + return anchorExists ? 'skip-anchor' : 'skip-turn'; + } catch { + return 'skip-turn'; + } +} + +/** 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'; + /** 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: withhold, + ...(withhold && { + 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; +}; + +/** How a failed generation's fresh error row proceeds after its existing rows + * are settled. */ +export type ErrorTurnSettlement = + /** An existing row covers the turn; the caller skips the error row. */ + | { covered: true } + /** The error row is written, under the live response id when the + * anchor-shaped collision makes the error id unusable for it. */ + | { covered: false; errorRowMessageId?: string }; + +/** + * Settles the rows a failed generation already persisted before its error row + * is written, through the caller's injected reads and write. + * + * 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. When no live + * row exists either, the error row is still written, redirected to the live + * response id so it can never overwrite the anchor. 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) { + return { covered: true }; + } + if (await settleLiveRow()) { + return { covered: true }; + } + /** The match is the anchor itself: the error row goes to the failed + * run's own response id when one exists, and is withheld entirely when + * none does (a failure before the id was allocated), because writing it + * under the error id would overwrite the anchor. */ + if (liveResponseMessageId != null && liveResponseMessageId !== errorMessageId) { + return { covered: false, errorRowMessageId: liveResponseMessageId }; + } + return { covered: true }; + } + return { covered: await settleLiveRow() }; +} + +/** + * 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.write) { + return false; + } + 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; +} + +/** 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 on a row that was already settled. */ + | { write: false } + /** The terminal marking is applied to the parts (a legacy or snapshot row + * may carry failure parts that never got the identity marker) and the row + * settles with the terminal envelope. */ + | { write: true; content: TMessageContentParts[] }; + +/** + * 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 has the terminal marking reapplied (idempotent + * for marked parts, stamping legacy parts that predate the marker) beside + * its settled envelope. 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; unfinished?: boolean } | null | undefined, + requestBody: { compact?: boolean } | null | undefined, +): FinalizedCompactionTurn { + if (requestBody?.compact !== true) { + return { write: false }; + } + 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 sawFailure = false; + let sawCheckpoint = false; + let unfinishedSummary = false; + for (const part of content) { + if (part?.type === ContentTypes.SUMMARY) { + if (part.failed === true) { + sawFailure = true; + } else if (isUsableSummaryPart(part)) { + sawCheckpoint = true; + } else { + unfinishedSummary = true; + } + } else if (part?.type === ContentTypes.ERROR) { + sawFailure = true; + } + } + if (unfinishedSummary || sawFailure) { + return { write: true, content: markAbortedCompactionContent(content, true) }; + } + if (sawCheckpoint) { + return partialRow?.unfinished === true + ? { write: true, content: markAbortedCompactionContent(content, true) } + : { write: false }; + } + return { write: true, 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 diff --git a/packages/api/src/stream/GenerationJobManager.ts b/packages/api/src/stream/GenerationJobManager.ts index a54ee838beb..fbafc7262ae 100644 --- a/packages/api/src/stream/GenerationJobManager.ts +++ b/packages/api/src/stream/GenerationJobManager.ts @@ -90,6 +90,7 @@ import { } from './internal/timing'; import { filterPersistableAbortContent } from './abortContent'; import { toClientPendingAction } from '~/agents/hitl/policy'; +import { markAbortedCompactionContent } from '~/agents/compaction'; import { ApprovalLifecycle, pausePersistenceActionId } from './ApprovalLifecycle'; import { projectPendingMCPOAuthPrompts } from '~/mcp/oauth/resume'; import { sanitizeJobMetadata } from './metadata'; @@ -4764,6 +4765,13 @@ class GenerationJobManagerClass { // Filter only after the transform so sparse/empty/OAuth parts cannot // shift a retained ID-less ask answer onto a different tool call. abortContent = filterPersistableAbortContent(content); + // A stopped compaction is unfinished rather than failed, so the row keeps + // the abort shape, plus the marker that keeps it identifiable as the + // compaction's own turn instead of an answer to its parent. + abortContent = markAbortedCompactionContent( + abortContent as TMessageContentParts[], + jobData.compact === true, + ); shouldPersistAbortContent = abortContent.length > 0; text = shouldPersistAbortContent ? parseTextParts(abortContent as TMessageContentParts[], false, { diff --git a/packages/api/src/stream/__tests__/RedisJobStore.spec.ts b/packages/api/src/stream/__tests__/RedisJobStore.spec.ts index 72b873051c8..86a0c0a615a 100644 --- a/packages/api/src/stream/__tests__/RedisJobStore.spec.ts +++ b/packages/api/src/stream/__tests__/RedisJobStore.spec.ts @@ -335,6 +335,7 @@ describe('RedisJobStore', () => { iconURL: 'https://example.com/icon.png', model: 'test-model', agent_id: 'agent-1', + compact: true, isTemporary: false, retentionExpiresAt: '2030-01-01T00:00:00.000Z', agentEventDeliveryKey: 'completion-delivery-1', @@ -393,6 +394,9 @@ describe('RedisJobStore', () => { * degrading to ordinary steering in every Redis deployment. */ expect(job.preemptCapable).toBe(true); + /** The abort paths read this flag from a reloaded job, so it has to + * survive the Redis round trip for a stopped compaction to be stamped. */ + expect(job.compact).toBe(true); expect(job.steerQuotesExecutionId).toBe('exec-1'); expect(job.generationProtocolVersion).toBe(2); expect(job.checkpointNamespace).toEqual(expect.any(String)); diff --git a/packages/api/src/stream/__tests__/abortCompactionIdentity.spec.ts b/packages/api/src/stream/__tests__/abortCompactionIdentity.spec.ts new file mode 100644 index 00000000000..12da6035a79 --- /dev/null +++ b/packages/api/src/stream/__tests__/abortCompactionIdentity.spec.ts @@ -0,0 +1,148 @@ +/** + * A Stop persists the aborted turn from job data alone, so the row only knows + * the run was a compaction through the job's `compact` metadata. Without the + * marker the stopped row reads as an answer to the message it hangs off and + * keeps that message's rerun controls: on a branch ending in a user message, + * Regenerate would answer the user turn behind the compaction. + */ +import { ContentTypes, ErrorTypes } from 'librechat-data-provider'; +import type { Agents } from 'librechat-data-provider'; + +/** Suppress winston Console transport output (survives jest.resetModules) */ +jest.spyOn(console, 'log').mockImplementation(); + +const COMPACTION_FAILED_ERROR = JSON.stringify({ type: ErrorTypes.COMPACTION_FAILED }); + +/** Streamed deltas never carry a boundary; a stopped round keeps them. */ +const partialSummary: Agents.MessageContentComplex = { + type: ContentTypes.SUMMARY, + content: [{ type: ContentTypes.TEXT, text: 'Half a summary' }], + summarizing: true, +}; + +const partialText: Agents.MessageContentComplex = { type: ContentTypes.TEXT, text: 'Partial' }; + +/** The abort result's `finalEvent` is `unknown` to the interface; the fields + * these tests read are the ones the abort FINAL always carries. */ +type AbortFinalEvent = { + responseMessage?: { unfinished?: boolean; error?: boolean } | null; + earlyAbort?: boolean; +}; + +async function configureManager() { + const { GenerationJobManager } = await import('../GenerationJobManager'); + const { InMemoryJobStore } = await import('../implementations/InMemoryJobStore'); + const { InMemoryEventTransport } = await import('../implementations/InMemoryEventTransport'); + + const jobStore = new InMemoryJobStore(); + GenerationJobManager.configure({ + jobStore, + eventTransport: new InMemoryEventTransport(), + isRedis: false, + cleanupOnComplete: false, + }); + GenerationJobManager.initialize(); + return { manager: GenerationJobManager, jobStore }; +} + +describe('abortJob compaction identity', () => { + beforeEach(() => { + jest.resetModules(); + }); + + it('stamps the partial summary a stopped compaction had streamed as failed', async () => { + const { manager, jobStore } = await configureManager(); + const streamId = 'abort-compaction-partial'; + const job = await manager.createJob(streamId, 'user-1', 'conversation-1', { + initialMetadata: { compact: true }, + }); + jobStore.setContentParts(streamId, [partialSummary], job.createdAt); + + const result = await manager.abortJob(streamId); + const finalEvent = result.finalEvent as AbortFinalEvent; + + expect(result.success).toBe(true); + /** The row keeps its unfinished shape; the partial summary keeps its text + * but reads as failed, or its label presents the truncated prefix as a + * finished checkpoint. */ + expect(finalEvent.responseMessage).toMatchObject({ unfinished: true, error: false }); + expect(result.content).toEqual([ + { ...partialSummary, initiatedBy: 'user', failed: true } as Agents.MessageContentComplex, + ]); + + await manager.destroy(); + }); + + /** 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 a stopped compaction streamed nothing', async () => { + const { manager } = await configureManager(); + const streamId = 'abort-compaction-empty'; + await manager.createJob(streamId, 'user-1', 'conversation-1', { + initialMetadata: { compact: true }, + }); + + const result = await manager.abortJob(streamId); + const finalEvent = result.finalEvent as AbortFinalEvent; + const expectedContent: Agents.MessageContentComplex[] = [ + { + type: ContentTypes.ERROR, + error: COMPACTION_FAILED_ERROR, + initiatedBy: 'user', + }, + ]; + + expect(result.success).toBe(true); + expect(result.content).toEqual(expectedContent); + /** The typed failure makes the row persistable, so this is not an early + * abort: the compaction's turn exists and must reach storage. */ + expect(finalEvent.responseMessage).not.toBeNull(); + expect(finalEvent.earlyAbort).not.toBe(true); + + await manager.destroy(); + }); + + /** A placeholder the summarizer opened but never streamed text into carries + * nothing to show: the typed failure replaces it as the row's outcome. */ + it('replaces an empty summary placeholder with the typed failure', async () => { + const { manager, jobStore } = await configureManager(); + const streamId = 'abort-compaction-placeholder'; + const placeholder: Agents.MessageContentComplex = { + type: ContentTypes.SUMMARY, + content: [], + summarizing: true, + }; + const job = await manager.createJob(streamId, 'user-1', 'conversation-1', { + initialMetadata: { compact: true }, + }); + jobStore.setContentParts(streamId, [placeholder], job.createdAt); + + const result = await manager.abortJob(streamId); + + expect(result.success).toBe(true); + expect(result.content).toEqual([ + { + type: ContentTypes.ERROR, + error: COMPACTION_FAILED_ERROR, + initiatedBy: 'user', + }, + ]); + + await manager.destroy(); + }); + + it('leaves a stopped ordinary turn without the marker', async () => { + const { manager, jobStore } = await configureManager(); + const streamId = 'abort-ordinary-turn'; + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + jobStore.setContentParts(streamId, [partialText], job.createdAt); + + const result = await manager.abortJob(streamId); + + expect(result.success).toBe(true); + expect(result.content).toEqual([partialText]); + expect(result.content[0]).not.toHaveProperty('initiatedBy'); + + await manager.destroy(); + }); +}); diff --git a/packages/api/src/stream/implementations/RedisJobStore.ts b/packages/api/src/stream/implementations/RedisJobStore.ts index 36848f35225..f548a488cc2 100644 --- a/packages/api/src/stream/implementations/RedisJobStore.ts +++ b/packages/api/src/stream/implementations/RedisJobStore.ts @@ -5318,6 +5318,7 @@ export class RedisJobStore implements IJobStoreV2 { userMessage: data.userMessage ? JSON.parse(data.userMessage) : undefined, responseMessageId: data.responseMessageId || undefined, isRegenerate: data.isRegenerate != null ? data.isRegenerate === '1' : undefined, + compact: data.compact != null ? data.compact === '1' : undefined, mcpRequestBody: data.mcpRequestBody ? JSON.parse(data.mcpRequestBody) : undefined, userSubmittedPaths: data.userSubmittedPaths ? JSON.parse(data.userSubmittedPaths) : undefined, userSubmittedMessageFieldPaths: data.userSubmittedMessageFieldPaths diff --git a/packages/api/src/stream/interfaces/IJobStore.ts b/packages/api/src/stream/interfaces/IJobStore.ts index 2e7d471b360..3952a2b9ff6 100644 --- a/packages/api/src/stream/interfaces/IJobStore.ts +++ b/packages/api/src/stream/interfaces/IJobStore.ts @@ -169,6 +169,9 @@ export interface SerializableJobData { /** Whether this generation replaces an existing assistant branch. */ isRegenerate?: boolean; + /** Whether this generation is a manual context compaction; the abort paths + * read it to stamp the stopped row with the compaction's identity. */ + compact?: boolean; /** Exact normalized MCP placeholder identity for this turn. */ mcpRequestBody?: MCPRuntimeRequestBody; /** Exact assistant-message fields authored by the user during this running job. */ @@ -448,6 +451,7 @@ export type JobMetadataPatch = Partial< SerializableJobData, | 'responseMessageId' | 'isRegenerate' + | 'compact' | 'mcpRequestBody' | 'userSubmittedPaths' | 'userSubmittedMessageFieldPaths' diff --git a/packages/api/src/stream/metadata.ts b/packages/api/src/stream/metadata.ts index a3bdbcbf041..c1a01f168a2 100644 --- a/packages/api/src/stream/metadata.ts +++ b/packages/api/src/stream/metadata.ts @@ -9,6 +9,9 @@ export function sanitizeJobMetadata(metadata: Partial): J if (metadata.isRegenerate !== undefined) { patch.isRegenerate = metadata.isRegenerate; } + if (metadata.compact !== undefined) { + patch.compact = metadata.compact; + } if (metadata.mcpRequestBody) { patch.mcpRequestBody = metadata.mcpRequestBody; } diff --git a/packages/api/src/types/stream.ts b/packages/api/src/types/stream.ts index d3096cb1c82..8105192f7ab 100644 --- a/packages/api/src/types/stream.ts +++ b/packages/api/src/types/stream.ts @@ -31,6 +31,9 @@ export interface GenerationJobMetadata { responseMessageId?: string; /** Whether this generation replaces an existing assistant branch. */ isRegenerate?: boolean; + /** Whether this generation is a manual context compaction; the abort paths + * read it to stamp the stopped row with the compaction's identity. */ + compact?: boolean; /** Exact normalized MCP placeholder identity for this turn. Persisted so HITL * resume does not reconstruct a different parent or overridden conversation. */ mcpRequestBody?: MCPRuntimeRequestBody;