From 5c1f21eb2130c4dda0f7e3faaebf70121a860312 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 19:09:47 -0700 Subject: [PATCH 1/3] fix(knowledge): withdraw a queued generation when its Search KB turns dormant --- .../dormant-search-processing.integration.ts | 48 +++++++++++++++++++ apps/sim/lib/knowledge/documents/service.ts | 21 ++++++++ 2 files changed, 69 insertions(+) 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..854ccea876c 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,52 @@ 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', async () => { + const queuedAt = new Date('2026-09-29T00:00:00.000Z') + const token = generateId() + await db + .update(document) + .set({ processingQueueToken: token, processingQueuedAt: queuedAt, processingAttempts: 1 }) + .where(eq(document.id, documentId)) + expect( + await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', { + chargedAtDispatch: true, + processingQueueToken: token, + processingQueuedAt: queuedAt, + }) + ).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: 0 }) + }) + + it('leaves a newer queued generation untouched', 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(), + }) + 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..3b463a3de56 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -1673,6 +1673,27 @@ 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'), + ...queueGenerationConditions(attemptContext) + ) + ) + } return { outcome: 'skipped', reason: 'unavailable' } } if (contextRows.length === 0) { From ccce52b6bd0ae5bb569496071c9e779b23fb9881 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 19:16:52 -0700 Subject: [PATCH 2/3] fix(knowledge): refund a dormant queued generation only once --- .../dormant-search-processing.integration.ts | 31 +++++++++++++------ apps/sim/lib/knowledge/documents/service.ts | 2 ++ 2 files changed, 23 insertions(+), 10 deletions(-) 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 854ccea876c..455046d4031 100644 --- a/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts @@ -96,20 +96,31 @@ describe('dormant Search document processing', () => { expect(stored.status).toBe('pending') }) - it('withdraws its own queued generation and refunds the charged attempt', async () => { + 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: 1 }) + .set({ processingQueueToken: token, processingQueuedAt: queuedAt, processingAttempts: 2 }) .where(eq(document.id, documentId)) - expect( - await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', { - chargedAtDispatch: true, - processingQueueToken: token, - processingQueuedAt: queuedAt, - }) - ).toEqual({ outcome: 'skipped', reason: 'unavailable' }) + 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, @@ -119,7 +130,7 @@ describe('dormant Search document processing', () => { }) .from(document) .where(eq(document.id, documentId)) - expect(stored).toEqual({ status: 'pending', token, queuedAt: null, attempts: 0 }) + expect(stored).toEqual({ status: 'pending', token, queuedAt: null, attempts: 1 }) }) it('leaves a newer queued generation untouched', async () => { diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 3b463a3de56..483d81375e6 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -1690,6 +1690,8 @@ export async function processDocumentAsync( and( eq(document.id, documentId), eq(document.processingStatus, 'pending'), + /** A duplicate of an already-withdrawn generation must not refund again. */ + isNotNull(document.processingQueuedAt), ...queueGenerationConditions(attemptContext) ) ) From 25a4c222bc6056582b43d877223736638e599aa9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 29 Sep 2026 19:30:19 -0700 Subject: [PATCH 3/3] fix(knowledge): withdraw only the exact dormant queue stamp --- .../dormant-search-processing.integration.ts | 7 ++++++- apps/sim/lib/knowledge/documents/service.ts | 9 +++++++-- 2 files changed, 13 insertions(+), 3 deletions(-) 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 455046d4031..9da06ae10c6 100644 --- a/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts @@ -133,7 +133,7 @@ describe('dormant Search document processing', () => { expect(stored).toEqual({ status: 'pending', token, queuedAt: null, attempts: 1 }) }) - it('leaves a newer queued generation untouched', async () => { + 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 @@ -144,6 +144,11 @@ describe('dormant Search document processing', () => { 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, diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 483d81375e6..c767cd0b3ca 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -1690,8 +1690,13 @@ export async function processDocumentAsync( and( eq(document.id, documentId), eq(document.processingStatus, 'pending'), - /** A duplicate of an already-withdrawn generation must not refund again. */ - isNotNull(document.processingQueuedAt), + /** + * 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) ) )