Skip to content

Commit f7d0259

Browse files
authored
fix(knowledge): withdraw a queued generation when its Search KB turns dormant (#8445)
* fix(knowledge): withdraw a queued generation when its Search KB turns dormant * fix(knowledge): refund a dormant queued generation only once * fix(knowledge): withdraw only the exact dormant queue stamp
1 parent 2f52353 commit f7d0259

2 files changed

Lines changed: 92 additions & 0 deletions

File tree

‎apps/sim/lib/knowledge/__integration__/dormant-search-processing.integration.ts‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,4 +95,68 @@ describe('dormant Search document processing', () => {
9595
.where(eq(document.id, documentId))
9696
expect(stored.status).toBe('pending')
9797
})
98+
99+
it('withdraws its own queued generation and refunds the charged attempt once', async () => {
100+
const queuedAt = new Date('2026-09-29T00:00:00.000Z')
101+
const token = generateId()
102+
await db
103+
.update(document)
104+
.set({ processingQueueToken: token, processingQueuedAt: queuedAt, processingAttempts: 2 })
105+
.where(eq(document.id, documentId))
106+
const attempt = {
107+
chargedAtDispatch: true,
108+
processingQueueToken: token,
109+
processingQueuedAt: queuedAt,
110+
}
111+
for (const _ of [1, 2]) {
112+
expect(
113+
await processDocumentAsync(
114+
ids.knowledgeBaseId,
115+
documentId,
116+
source,
117+
{},
118+
undefined,
119+
'pass',
120+
attempt
121+
)
122+
).toEqual({ outcome: 'skipped', reason: 'unavailable' })
123+
}
124+
const [stored] = await db
125+
.select({
126+
status: document.processingStatus,
127+
token: document.processingQueueToken,
128+
queuedAt: document.processingQueuedAt,
129+
attempts: document.processingAttempts,
130+
})
131+
.from(document)
132+
.where(eq(document.id, documentId))
133+
expect(stored).toEqual({ status: 'pending', token, queuedAt: null, attempts: 1 })
134+
})
135+
136+
it('leaves a newer queued generation untouched, even under a reused token', async () => {
137+
const queuedAt = new Date('2026-09-29T01:00:00.000Z')
138+
const newer = generateId()
139+
await db
140+
.update(document)
141+
.set({ processingQueueToken: newer, processingQueuedAt: queuedAt, processingAttempts: 1 })
142+
.where(eq(document.id, documentId))
143+
await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', {
144+
chargedAtDispatch: true,
145+
processingQueueToken: generateId(),
146+
})
147+
await processDocumentAsync(ids.knowledgeBaseId, documentId, source, {}, undefined, 'pass', {
148+
chargedAtDispatch: true,
149+
processingQueueToken: newer,
150+
processingQueuedAt: new Date('2026-09-29T00:30:00.000Z'),
151+
})
152+
const [stored] = await db
153+
.select({
154+
token: document.processingQueueToken,
155+
queuedAt: document.processingQueuedAt,
156+
attempts: document.processingAttempts,
157+
})
158+
.from(document)
159+
.where(eq(document.id, documentId))
160+
expect(stored).toEqual({ token: newer, queuedAt, attempts: 1 })
161+
})
98162
})

‎apps/sim/lib/knowledge/documents/service.ts‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1673,6 +1673,34 @@ export async function processDocumentAsync(
16731673
.limit(1)
16741674

16751675
if (contextRows[0] && !requiresConnectorIndexing(contextRows[0].isSearchIndex)) {
1676+
/**
1677+
* A generation queued before its KB went dormant (e.g. legacy index adoption) gives back its
1678+
* stamp and charged attempt, as `clearDocumentsQueued` does; the token stays as the owner.
1679+
*/
1680+
if (attemptContext?.processingQueueToken || attemptContext?.processingQueuedAt) {
1681+
await db
1682+
.update(document)
1683+
.set({
1684+
processingQueuedAt: null,
1685+
...(attemptContext.chargedAtDispatch
1686+
? { processingAttempts: sql`GREATEST(${document.processingAttempts} - 1, 0)` }
1687+
: {}),
1688+
})
1689+
.where(
1690+
and(
1691+
eq(document.id, documentId),
1692+
eq(document.processingStatus, 'pending'),
1693+
/**
1694+
* Only the exact stamp this payload was queued with: a duplicate of an already
1695+
* withdrawn generation, or a newer stamp under a reused token, is left alone.
1696+
*/
1697+
attemptContext.processingQueuedAt
1698+
? eq(document.processingQueuedAt, attemptContext.processingQueuedAt)
1699+
: isNotNull(document.processingQueuedAt),
1700+
...queueGenerationConditions(attemptContext)
1701+
)
1702+
)
1703+
}
16761704
return { outcome: 'skipped', reason: 'unavailable' }
16771705
}
16781706
if (contextRows.length === 0) {

0 commit comments

Comments
 (0)