diff --git a/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts index 98fc26882b8..9da06ae10c6 100644 --- a/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts @@ -95,4 +95,68 @@ describe('dormant Search document processing', () => { .where(eq(document.id, documentId)) expect(stored.status).toBe('pending') }) + + it('withdraws its own queued generation and refunds the charged attempt once', async () => { + const queuedAt = new Date('2026-09-29T00:00:00.000Z') + const token = generateId() + await db + .update(document) + .set({ processingQueueToken: token, processingQueuedAt: queuedAt, processingAttempts: 2 }) + .where(eq(document.id, documentId)) + const attempt = { + chargedAtDispatch: true, + processingQueueToken: token, + processingQueuedAt: queuedAt, + } + for (const _ of [1, 2]) { + expect( + await processDocumentAsync( + ids.knowledgeBaseId, + documentId, + source, + {}, + undefined, + 'pass', + attempt + ) + ).toEqual({ outcome: 'skipped', reason: 'unavailable' }) + } + const [stored] = await db + .select({ + status: document.processingStatus, + token: document.processingQueueToken, + queuedAt: document.processingQueuedAt, + attempts: document.processingAttempts, + }) + .from(document) + .where(eq(document.id, documentId)) + expect(stored).toEqual({ status: 'pending', token, queuedAt: null, attempts: 1 }) + }) + + it('leaves a newer queued generation untouched, even under a reused token', async () => { + const queuedAt = new Date('2026-09-29T01:00:00.000Z') + const newer = generateId() + await db + .update(document) + .set({ processingQueueToken: newer, processingQueuedAt: queuedAt, processingAttempts: 1 }) + .where(eq(document.id, documentId)) + await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', { + chargedAtDispatch: true, + processingQueueToken: generateId(), + }) + await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', { + chargedAtDispatch: true, + processingQueueToken: newer, + processingQueuedAt: new Date('2026-09-29T00:30:00.000Z'), + }) + const [stored] = await db + .select({ + token: document.processingQueueToken, + queuedAt: document.processingQueuedAt, + attempts: document.processingAttempts, + }) + .from(document) + .where(eq(document.id, documentId)) + expect(stored).toEqual({ token: newer, queuedAt, attempts: 1 }) + }) }) diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 82c09359367..c767cd0b3ca 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -1673,6 +1673,34 @@ export async function processDocumentAsync( .limit(1) if (contextRows[0] && !requiresConnectorIndexing(contextRows[0].isSearchIndex)) { + /** + * A generation queued before its KB went dormant (e.g. legacy index adoption) gives back its + * stamp and charged attempt, as `clearDocumentsQueued` does; the token stays as the owner. + */ + if (attemptContext?.processingQueueToken || attemptContext?.processingQueuedAt) { + await db + .update(document) + .set({ + processingQueuedAt: null, + ...(attemptContext.chargedAtDispatch + ? { processingAttempts: sql`GREATEST(${document.processingAttempts} - 1, 0)` } + : {}), + }) + .where( + and( + eq(document.id, documentId), + eq(document.processingStatus, 'pending'), + /** + * Only the exact stamp this payload was queued with: a duplicate of an already + * withdrawn generation, or a newer stamp under a reused token, is left alone. + */ + attemptContext.processingQueuedAt + ? eq(document.processingQueuedAt, attemptContext.processingQueuedAt) + : isNotNull(document.processingQueuedAt), + ...queueGenerationConditions(attemptContext) + ) + ) + } return { outcome: 'skipped', reason: 'unavailable' } } if (contextRows.length === 0) {