Skip to content

Commit 74140e0

Browse files
committed
fix(knowledge): stop recovery from re-dispatching documents whose runs are still queued
1 parent 3f245d7 commit 74140e0

15 files changed

Lines changed: 269 additions & 14 deletions

‎apps/sim/lib/knowledge/__integration__/connector-lifecycle-locks.integration.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -259,4 +259,37 @@ describe('source lifecycle KB guards', () => {
259259
.where(eq(document.id, retryDocumentId))
260260
expect(unchanged.status).toBe('failed')
261261
})
262+
263+
it('moves an expired queued generation charge to its replacement instead of adding one', async () => {
264+
const stampedAt = new Date(Date.now() - 24 * 60 * 60_000)
265+
await db
266+
.update(document)
267+
.set({
268+
processingStatus: 'pending',
269+
processingAttempts: 2,
270+
processingQueuedAt: stampedAt,
271+
processingQueueToken: 'expired-generation',
272+
processingCompletedAt: null,
273+
})
274+
.where(eq(document.id, retryDocumentId))
275+
await run('recover')
276+
const [row] = await db
277+
.select({
278+
status: document.processingStatus,
279+
attempts: document.processingAttempts,
280+
token: document.processingQueueToken,
281+
})
282+
.from(document)
283+
.where(eq(document.id, retryDocumentId))
284+
expect(row).toEqual({ status: 'pending', attempts: 1, token: null })
285+
})
286+
287+
it('keeps the charge of an attempt that reached a worker', async () => {
288+
await run('recover')
289+
const [row] = await db
290+
.select({ attempts: document.processingAttempts })
291+
.from(document)
292+
.where(eq(document.id, retryDocumentId))
293+
expect(row.attempts).toBe(1)
294+
})
262295
})

‎apps/sim/lib/knowledge/__integration__/stored-document-recovery.integration.ts‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -244,6 +244,49 @@ describe('independent recovery of retained connector documents', () => {
244244
}
245245
})
246246

247+
it('replaces expired queued generations without spending the budget, then indexes once', async () => {
248+
const ids = await seed()
249+
const file = await failedFile(ids)
250+
const expire = (token: string) =>
251+
db
252+
.update(document)
253+
.set({
254+
processingStatus: 'pending',
255+
processingCompletedAt: null,
256+
processingError: null,
257+
processingQueuedAt: old(),
258+
processingQueueToken: token,
259+
})
260+
.where(eq(document.id, file.documentId))
261+
let generation = 'old-fixture-generation'
262+
await expire(generation)
263+
for (let cycle = 0; cycle <= MAX_PROCESSING_ATTEMPTS; cycle++) {
264+
expect(await recoverKnowledgeDocumentProcessing()).toBe(1)
265+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
266+
expect(row.processingAttempts).toBe(1)
267+
expect(row.processingQueueToken).not.toBe(generation)
268+
generation = row.processingQueueToken!
269+
await expire(generation)
270+
}
271+
await db
272+
.update(document)
273+
.set({ processingQueuedAt: new Date() })
274+
.where(eq(document.id, file.documentId))
275+
expect(await recoverKnowledgeDocumentProcessing()).toBe(0)
276+
277+
const events = await eventsFor(ids)
278+
expect(events).toHaveLength(MAX_PROCESSING_ATTEMPTS + 1)
279+
const before = fixture.embeddingCalls
280+
for (const event of events) {
281+
expect(
282+
await outbox.processOutboxEventById(event.id, knowledgeDocumentProcessingOutboxHandlers)
283+
).toBe('completed')
284+
}
285+
expect(fixture.embeddingCalls).toBe(before + 1)
286+
const [indexed] = await db.select().from(document).where(eq(document.id, file.documentId))
287+
expect(indexed.processingStatus, indexed.processingError ?? undefined).toBe('completed')
288+
})
289+
247290
it('keeps permanent, excluded, paused, fresh, deferred, and expired work out of automatic recovery', async () => {
248291
const ids = await seed()
249292
const file = await failedFile(ids)

‎apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
/** @vitest-environment node */
22
import {
33
dbChainMockFns,
4+
hasMockCondition,
5+
type MockCondition,
46
queueTableRows,
57
resetDbChainMock as resetDatabaseMock,
68
schemaMock,
@@ -643,6 +645,28 @@ describe('content pass checkpoint intent', () => {
643645
)
644646
})
645647

648+
describe('resurrecting verified listed documents', () => {
649+
it('only writes documents that are actually tombstoned', async () => {
650+
sourceBody = { value: '<p>Current content</p>' }
651+
await runPass({ access: 'admin' })
652+
const resurrect = dbChainMockFns.set.mock.invocationCallOrder.find((_order, index) => {
653+
const [values] = dbChainMockFns.set.mock.calls[index]
654+
return Object.keys(values).length === 1 && values.deletedAt === null
655+
})
656+
expect(resurrect).toBeDefined()
657+
const whereIndex = dbChainMockFns.where.mock.invocationCallOrder.findIndex(
658+
(order) => order > resurrect!
659+
)
660+
expect(
661+
hasMockCondition(
662+
dbChainMockFns.where.mock.calls[whereIndex][0],
663+
(node: MockCondition) =>
664+
node.type === 'isNotNull' && node.column === schemaMock.document.deletedAt
665+
)
666+
).toBe(true)
667+
})
668+
})
669+
646670
describe('permission refresh through the shared content pass', () => {
647671
const current: StoredPage = {
648672
...EXISTING,

‎apps/sim/lib/knowledge/connectors/sync-content-pass.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -290,6 +290,7 @@ export async function runConnectorContentPass(input: ContentPassInput) {
290290
and(
291291
eq(document.connectorId, input.connectorId),
292292
inArray(document.externalId, verified.slice(offset, offset + 500)),
293+
isNotNull(document.deletedAt),
293294
isNotNull(document.contentHash),
294295
isNull(document.archivedAt)
295296
)

‎apps/sim/lib/knowledge/connectors/sync-primitives.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,10 @@ import {
2727
persistSkippedRetryHashes,
2828
updateDocument,
2929
} from '@/lib/knowledge/connectors/sync-persistence'
30-
import { documentProcessingRecoveryCondition } from '@/lib/knowledge/documents/processing-recovery-policy'
30+
import {
31+
documentProcessingRecoveryCondition,
32+
releaseUnclaimedDispatchAttempt,
33+
} from '@/lib/knowledge/documents/processing-recovery-policy'
3134
import { DOCUMENT_PROCESSING_STALE_THRESHOLD_MS } from '@/lib/knowledge/documents/processing-timeouts.server'
3235
import type { DocumentData } from '@/lib/knowledge/documents/service'
3336
import { isTriggerAvailable, processDocumentsWithQueue } from '@/lib/knowledge/documents/service'
@@ -1446,6 +1449,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
14461449
processingDeferredUntil: null,
14471450
processingCompletedAt: null,
14481451
processingError: null,
1452+
processingAttempts: releaseUnclaimedDispatchAttempt(sweepEvaluatedAt),
14491453
chunkCount: 0,
14501454
tokenCount: 0,
14511455
characterCount: 0,

‎apps/sim/lib/knowledge/documents/processing-continuation-dispatch.test.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import type {
1919
DocumentProcessingLane,
2020
DocumentProcessingPayload,
2121
} from '@/lib/knowledge/documents/processing-payload'
22+
import { QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowledge/documents/types'
2223

2324
const BILLING_ATTRIBUTION = {
2425
actorUserId: 'user-1',
@@ -73,4 +74,20 @@ describe('dispatchDocumentProcessingContinuation', () => {
7374
concurrencyKey: 'organization:org-1',
7475
})
7576
})
77+
78+
/** Trigger.dev starts a delayed run's TTL when the delay ends; the deadline follows the stamp. */
79+
it('expires a continuation unstarted before recovery may replace it', async () => {
80+
const deferredUntil = new Date(Date.now() + 15 * 60_000)
81+
await dispatchDocumentProcessingContinuation(
82+
{ ...payload('backfill'), processingQueuedAt: deferredUntil.toISOString() },
83+
deferredUntil,
84+
'continuation-key',
85+
true
86+
)
87+
88+
const { ttl } = mockTrigger.mock.calls[0][2]
89+
expect(mockTrigger.mock.calls[0][2]).toMatchObject({ delay: deferredUntil })
90+
expect(ttl).toBeGreaterThan(0)
91+
expect(ttl * 1000).toBeLessThan(QUEUED_DISPATCH_GRACE_MS)
92+
})
7693
})

‎apps/sim/lib/knowledge/documents/processing-continuation-dispatch.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ import { resolveTriggerRegion } from '@/lib/core/async-jobs/region'
55
import { env } from '@/lib/core/config/env'
66
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
77
import { isInsideTriggerRun } from '@/lib/core/config/trigger-runtime'
8-
import { documentProcessingQueueOptions } from '@/lib/knowledge/documents/processing-lane'
8+
import { documentProcessingRunOptions } from '@/lib/knowledge/documents/processing-lane'
99
import type { DocumentProcessingPayload } from '@/lib/knowledge/documents/processing-payload'
1010

1111
export interface DocumentProcessingContinuation {
@@ -37,7 +37,7 @@ export async function dispatchDocumentProcessingContinuation(
3737
* deferred on quota resume as interactive work and escape the tenant's
3838
* bulk ceiling — the retry path would become the way around the limit.
3939
*/
40-
...documentProcessingQueueOptions(payload),
40+
...documentProcessingRunOptions(payload, deferredUntil),
4141
region,
4242
})
4343
return

‎apps/sim/lib/knowledge/documents/processing-lane.test.ts‎

Lines changed: 54 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,19 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { describe, expect, it } from 'vitest'
4+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
55
import {
66
BACKFILL_PROCESSING_QUEUE_NAME,
77
documentProcessingQueueOptions,
8+
documentProcessingRunExpiry,
89
INTERACTIVE_PROCESSING_QUEUE_NAME,
910
} from '@/lib/knowledge/documents/processing-lane'
1011
import type { DocumentProcessingPayload } from '@/lib/knowledge/documents/processing-payload'
1112
import { resolveDocumentProcessingLane } from '@/lib/knowledge/documents/processing-payload'
13+
import {
14+
QUEUED_DISPATCH_GRACE_MS,
15+
QUEUED_DISPATCH_START_DEADLINE_MS,
16+
} from '@/lib/knowledge/documents/types'
1217

1318
const BILLING_ATTRIBUTION = {
1419
actorUserId: 'user-1',
@@ -130,3 +135,51 @@ describe('documentProcessingQueueOptions', () => {
130135
)
131136
})
132137
})
138+
139+
describe('documentProcessingRunExpiry', () => {
140+
const NOW = new Date('2026-09-18T12:00:00.000Z')
141+
142+
beforeEach(() => {
143+
vi.useFakeTimers()
144+
vi.setSystemTime(NOW)
145+
})
146+
afterEach(() => {
147+
vi.useRealTimers()
148+
})
149+
150+
it('expires an undelayed run at the start deadline, before recovery may replace it', () => {
151+
expect(QUEUED_DISPATCH_START_DEADLINE_MS).toBeLessThan(QUEUED_DISPATCH_GRACE_MS)
152+
expect(documentProcessingRunExpiry({ processingQueuedAt: NOW.toISOString() })).toEqual({
153+
ttl: QUEUED_DISPATCH_START_DEADLINE_MS / 1000,
154+
})
155+
})
156+
157+
it('gives a late relay only the time its generation has left', () => {
158+
const stampedAt = new Date(NOW.getTime() - 60 * 60_000)
159+
expect(documentProcessingRunExpiry({ processingQueuedAt: stampedAt.toISOString() })).toEqual({
160+
ttl: (QUEUED_DISPATCH_START_DEADLINE_MS - 60 * 60_000) / 1000,
161+
})
162+
})
163+
164+
it('expires a generation already past its deadline at the minimum', () => {
165+
const stampedAt = new Date(NOW.getTime() - QUEUED_DISPATCH_GRACE_MS)
166+
expect(documentProcessingRunExpiry({ processingQueuedAt: stampedAt.toISOString() })).toEqual({
167+
ttl: 1,
168+
})
169+
})
170+
171+
/** Trigger.dev starts a delayed run's TTL when its delay ends, not when it is triggered. */
172+
it('counts a delayed run from the end of its delay', () => {
173+
const deferredUntil = new Date(NOW.getTime() + 30 * 60_000)
174+
expect(
175+
documentProcessingRunExpiry(
176+
{ processingQueuedAt: deferredUntil.toISOString() },
177+
deferredUntil
178+
)
179+
).toEqual({ ttl: QUEUED_DISPATCH_START_DEADLINE_MS / 1000 })
180+
})
181+
182+
it('leaves a payload without a queue stamp on the queue default', () => {
183+
expect(documentProcessingRunExpiry({})).toEqual({})
184+
})
185+
})

‎apps/sim/lib/knowledge/documents/processing-lane.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import type {
33
DocumentProcessingBillingContext,
44
DocumentProcessingPayload,
55
} from '@/lib/knowledge/documents/processing-payload'
6+
import { QUEUED_DISPATCH_START_DEADLINE_MS } from '@/lib/knowledge/documents/types'
67

78
/**
89
* Queue backing the interactive lane.
@@ -55,3 +56,35 @@ export function documentProcessingQueueOptions(payload: DocumentProcessingPayloa
5556
concurrencyKey: documentProcessingTenantKey(payload),
5657
}
5758
}
59+
60+
/**
61+
* Trigger.dev `ttl` expiring a run unstarted at its generation's start deadline
62+
* (queue stamp + {@link QUEUED_DISPATCH_START_DEADLINE_MS}). Trigger.dev counts `ttl`
63+
* from enqueue, which for a delayed run is `notBefore`. A generation already past its
64+
* deadline gets the one-second minimum; one without a stamp keeps the queue default.
65+
*/
66+
export function documentProcessingRunExpiry(
67+
payload: Pick<DocumentProcessingPayload, 'processingQueuedAt'>,
68+
notBefore?: Date
69+
): { ttl?: number } {
70+
if (!payload.processingQueuedAt) return {}
71+
const deadline =
72+
new Date(payload.processingQueuedAt).getTime() + QUEUED_DISPATCH_START_DEADLINE_MS
73+
const enqueuedAt = Math.max(Date.now(), notBefore?.getTime() ?? 0)
74+
return { ttl: Math.max(1, Math.floor((deadline - enqueuedAt) / 1000)) }
75+
}
76+
77+
/**
78+
* Every Trigger.dev option a `knowledge-process-document` dispatch needs: its lane's
79+
* per-tenant queue and its start deadline. Dispatch sites use this rather than the
80+
* parts so none can enqueue a run that outlives recovery's grace.
81+
*/
82+
export function documentProcessingRunOptions(
83+
payload: DocumentProcessingPayload,
84+
notBefore?: Date
85+
): { queue: string; concurrencyKey: string; ttl?: number } {
86+
return {
87+
...documentProcessingQueueOptions(payload),
88+
...documentProcessingRunExpiry(payload, notBefore),
89+
}
90+
}

‎apps/sim/lib/knowledge/documents/processing-queue.test.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -368,6 +368,23 @@ describe('processDocumentsWithQueue dispatch backend', () => {
368368
})
369369
})
370370

371+
it('expires each queued run before recovery may replace its generation', async () => {
372+
markInsideTriggerRun()
373+
374+
await processDocumentsWithQueue(
375+
[DOCUMENT],
376+
'knowledge-base-1',
377+
{},
378+
'request-1',
379+
BILLING_ATTRIBUTION,
380+
'backfill'
381+
)
382+
383+
const [item] = mockBatchTrigger.mock.calls[0][1]
384+
expect(item.options.ttl).toBeGreaterThan(0)
385+
expect(item.options.ttl * 1000).toBeLessThan(QUEUED_DISPATCH_GRACE_MS)
386+
})
387+
371388
/**
372389
* The starvation this split exists to prevent: connector backfill must not be
373390
* able to occupy the queue a person's upload is admitted through.

0 commit comments

Comments
 (0)