Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 })
})
})
28 changes: 28 additions & 0 deletions apps/sim/lib/knowledge/documents/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
)
Comment thread
waleedlatif1 marked this conversation as resolved.
)
}
return { outcome: 'skipped', reason: 'unavailable' }
}
if (contextRows.length === 0) {
Expand Down
Loading