Skip to content

Commit f359031

Browse files
authored
fix(knowledge): preserve live document jobs during recovery (#7994)
1 parent 3f245d7 commit f359031

10 files changed

Lines changed: 484 additions & 32 deletions

‎apps/sim/background/knowledge-processing.test.ts‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,7 @@ describe('knowledge processing worker', () => {
119119
}
120120
return value
121121
})
122-
mockProcessDocumentAsync.mockResolvedValue(undefined)
122+
mockProcessDocumentAsync.mockResolvedValue({ outcome: 'indexed' })
123123
mockResolveTriggerRegion.mockResolvedValue('us-east-1')
124124
mockTrigger.mockResolvedValue({ id: 'quota-continuation-run' })
125125
})
@@ -128,6 +128,27 @@ describe('knowledge processing worker', () => {
128128
vi.restoreAllMocks()
129129
})
130130

131+
it('reports indexed only when the document service committed the index', async () => {
132+
expect(await runDocumentProcessing(WORKSPACE_PAYLOAD)).toMatchObject({
133+
success: true,
134+
outcome: 'indexed',
135+
documentId: WORKSPACE_PAYLOAD.documentId,
136+
})
137+
})
138+
139+
it.each(['unavailable', 'not_claimed', 'superseded'] as const)(
140+
'reports a harmless %s skip without turning it into a task failure or an indexed success',
141+
async (reason) => {
142+
mockProcessDocumentAsync.mockResolvedValue({ outcome: 'skipped', reason })
143+
expect(await runDocumentProcessing(WORKSPACE_PAYLOAD)).toMatchObject({
144+
success: false,
145+
outcome: 'skipped',
146+
reason,
147+
})
148+
expect(mockTrigger).not.toHaveBeenCalled()
149+
}
150+
)
151+
131152
it('rejects workspace work without attribution before document processing starts', async () => {
132153
await expect(
133154
runDocumentProcessing({

‎apps/sim/background/knowledge-processing.ts‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ export async function runDocumentProcessing(
5858
logger.info(`[${requestId}] Starting Trigger.dev processing for document: ${docData.filename}`)
5959

6060
try {
61-
await processDocumentAsync(
61+
const result = await processDocumentAsync(
6262
knowledgeBaseId,
6363
documentId,
6464
docData,
@@ -90,10 +90,11 @@ export async function runDocumentProcessing(
9090
}
9191
)
9292

93-
logger.info(`[${requestId}] Successfully processed document: ${docData.filename}`)
93+
logger.info(`[${requestId}] Document processing finished`, { documentId, ...result })
9494

9595
return {
96-
success: true,
96+
success: result.outcome === 'indexed',
97+
...result,
9798
documentId,
9899
filename: docData.filename,
99100
processingTime: Date.now() - startedAt,

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

Lines changed: 121 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,21 @@ import {
1717
} from '@sim/db/schema'
1818
import { generateId } from '@sim/utils/id'
1919
import { and, eq, inArray, sql } from 'drizzle-orm'
20-
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
20+
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
2121

22-
const fixture = vi.hoisted(() => ({ root: '', embeddingCalls: 0 }))
22+
const fixture = vi.hoisted(() => ({
23+
root: '',
24+
embeddingCalls: 0,
25+
queueEnabled: false,
26+
listRuns: vi.fn(),
27+
}))
28+
vi.mock('@/lib/core/config/trigger-runtime', () => ({
29+
isInsideTriggerRun: () => fixture.queueEnabled,
30+
}))
31+
vi.mock('@trigger.dev/sdk', async (original) => ({
32+
...(await original<typeof import('@trigger.dev/sdk')>()),
33+
runs: { list: fixture.listRuns },
34+
}))
2335
vi.mock('@/lib/uploads/core/setup.server', () => ({
2436
get UPLOAD_DIR_SERVER() {
2537
return fixture.root
@@ -52,6 +64,7 @@ import { searchKnowledge } from '@/lib/knowledge/application/search'
5264
import { searchScopedKnowledge } from '@/lib/knowledge/application/workspace-search'
5365
import { createContentSyncLease } from '@/lib/knowledge/connectors/sync-lock'
5466
import { addDocument } from '@/lib/knowledge/connectors/sync-persistence'
67+
import { sweepStuckDocuments } from '@/lib/knowledge/connectors/sync-primitives'
5568
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
5669
import {
5770
DOCUMENT_RECOVERY_BATCH_SIZE,
@@ -60,6 +73,7 @@ import {
6073
} from '@/lib/knowledge/documents/processing-recovery'
6174
import { processDocumentAsync } from '@/lib/knowledge/documents/service'
6275
import { MAX_PROCESSING_ATTEMPTS, QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowledge/documents/types'
76+
import type { SyncResult } from '@/connectors/types'
6377

6478
const fixtures: ReturnType<typeof createKnowledgeAclFixtureIds>[] = []
6579
const old = () => new Date(Date.now() - QUEUED_DISPATCH_GRACE_MS - 60_000)
@@ -119,6 +133,11 @@ async function failedFile(
119133
return file
120134
}
121135

136+
afterEach(() => {
137+
fixture.queueEnabled = false
138+
fixture.listRuns.mockReset()
139+
})
140+
122141
beforeAll(() => {
123142
fixture.root = mkdtempSync(path.join(tmpdir(), 'sim-stored-recovery-'))
124143
})
@@ -137,7 +156,107 @@ afterAll(async () => {
137156
await db.$client.end()
138157
})
139158

159+
async function recoverFixture(
160+
ids: ReturnType<typeof createKnowledgeAclFixtureIds>,
161+
mode: 'independent' | 'connector'
162+
) {
163+
if (mode === 'independent') return recoverKnowledgeDocumentProcessing()
164+
const result: SyncResult = {
165+
docsAdded: 0,
166+
docsUpdated: 0,
167+
docsDeleted: 0,
168+
docsUnchanged: 0,
169+
docsSkipped: 0,
170+
docsFailed: 0,
171+
processingDispatch: { requested: 0, accepted: 0, failed: 0 },
172+
}
173+
await sweepStuckDocuments({
174+
connectorId: ids.connectorId,
175+
knowledgeBaseId: ids.knowledgeBaseId,
176+
syncStartedAt: new Date(),
177+
retryCutoff: new Date(Date.now() - 7 * 24 * 60 * 60_000),
178+
billingAttribution: await resolveSystemBillingAttribution(ids.workspaceId),
179+
result,
180+
lease: createContentSyncLease(ids.connectorId, ids.lockId),
181+
})
182+
return result.processingDispatch.requested
183+
}
184+
140185
describe('independent recovery of retained connector documents', () => {
186+
it.each(['independent', 'connector'] as const)(
187+
'%s recovery preserves a job queued beyond the grace period',
188+
async (mode) => {
189+
const ids = await seed()
190+
const file = await failedFile(ids)
191+
await db
192+
.update(document)
193+
.set({ processingStatus: 'pending' })
194+
.where(eq(document.id, file.documentId))
195+
fixture.queueEnabled = true
196+
fixture.listRuns.mockResolvedValue({ data: [{ id: 'run-queued', status: 'QUEUED' }] })
197+
expect(await recoverFixture(ids, mode)).toBe(0)
198+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
199+
expect(row.processingAttempts).toBe(1)
200+
expect(row.processingQueueToken).toBe('old-fixture-generation')
201+
expect(row.processingRecoveryAfter).not.toBeNull()
202+
expect(await eventsFor(ids)).toHaveLength(0)
203+
await db
204+
.update(knowledgeConnector)
205+
.set({ status: 'paused' })
206+
.where(eq(knowledgeConnector.id, ids.connectorId))
207+
}
208+
)
209+
210+
it.each(['independent', 'connector'] as const)(
211+
'%s recovery rechecks the generation after its remote lookup',
212+
async (mode) => {
213+
const ids = await seed()
214+
const file = await failedFile(ids)
215+
fixture.queueEnabled = true
216+
fixture.listRuns.mockImplementation(async () => {
217+
await db
218+
.update(document)
219+
.set({ processingQueueToken: 'replacement-generation' })
220+
.where(eq(document.id, file.documentId))
221+
return { data: [] }
222+
})
223+
expect(await recoverFixture(ids, mode)).toBe(0)
224+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
225+
expect(row.processingAttempts).toBe(1)
226+
expect(row.processingQueueToken).toBe('replacement-generation')
227+
expect(await eventsFor(ids)).toHaveLength(0)
228+
await db
229+
.update(knowledgeConnector)
230+
.set({ status: 'paused' })
231+
.where(eq(knowledgeConnector.id, ids.connectorId))
232+
}
233+
)
234+
235+
it.each(['pending', 'processing'])(
236+
'does not replace an aged %s outbox continuation',
237+
async (status) => {
238+
const ids = await seed()
239+
const file = await failedFile(ids)
240+
const token = generateId()
241+
await db
242+
.update(document)
243+
.set({ processingQueueToken: token })
244+
.where(eq(document.id, file.documentId))
245+
await db.insert(outboxEvent).values({
246+
id: token,
247+
eventType: 'knowledge.document.processing.resume',
248+
payload: { knowledgeBaseId: ids.knowledgeBaseId, documentId: file.documentId },
249+
status,
250+
availableAt: old(),
251+
})
252+
expect(await recoverFixture(ids, 'independent')).toBe(0)
253+
expect(await recoverFixture(ids, 'connector')).toBe(0)
254+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
255+
expect(row.processingAttempts).toBe(1)
256+
expect(row.processingQueueToken).toBe(token)
257+
}
258+
)
259+
141260
it('uses the organization owner and preserves Search visibility during source backoff', async () => {
142261
const ids = await seed()
143262
await db.insert(member).values({

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

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ import {
99
import { createLogger } from '@sim/logger'
1010
import { toError } from '@sim/utils/errors'
1111
import { generateId } from '@sim/utils/id'
12-
import { and, asc, desc, eq, inArray, isNotNull, isNull, lt, ne, sql } from 'drizzle-orm'
12+
import { and, asc, desc, eq, inArray, isNotNull, isNull, lt, ne, or, sql } from 'drizzle-orm'
1313
import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution'
1414
import { ProviderCapacityDeferredError } from '@/lib/core/rate-limiter/provider-capacity-error'
1515
import { withDatabaseReadRetry } from '@/lib/db/read-retry'
@@ -28,6 +28,10 @@ import {
2828
updateDocument,
2929
} from '@/lib/knowledge/connectors/sync-persistence'
3030
import { documentProcessingRecoveryCondition } from '@/lib/knowledge/documents/processing-recovery-policy'
31+
import {
32+
documentRecoveryGenerationCondition,
33+
filterAbandonedDocumentProcessing,
34+
} from '@/lib/knowledge/documents/processing-recovery-queue'
3135
import { DOCUMENT_PROCESSING_STALE_THRESHOLD_MS } from '@/lib/knowledge/documents/processing-timeouts.server'
3236
import type { DocumentData } from '@/lib/knowledge/documents/service'
3337
import { isTriggerAvailable, processDocumentsWithQueue } from '@/lib/knowledge/documents/service'
@@ -1336,6 +1340,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
13361340
fileSize: document.fileSize,
13371341
mimeType: document.mimeType,
13381342
processingStatus: document.processingStatus,
1343+
processingQueueToken: document.processingQueueToken,
13391344
processingQueuedAt: document.processingQueuedAt,
13401345
processingStartedAt: document.processingStartedAt,
13411346
processingDeferredUntil: document.processingDeferredUntil,
@@ -1361,7 +1366,8 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
13611366
asc(document.id)
13621367
)
13631368
.limit(STUCK_RETRY_MAX_CANDIDATES_PER_SYNC)
1364-
const stuckDocs = sweepCandidates.filter(
1369+
const abandoned = await filterAbandonedDocumentProcessing(sweepCandidates)
1370+
const stuckDocs = abandoned.filter(
13651371
(row): row is typeof row & { processingStatus: DocumentProcessingStatus } =>
13661372
isDocumentProcessingStatus(row.processingStatus)
13671373
)
@@ -1402,6 +1408,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
14021408
fileSize: document.fileSize,
14031409
mimeType: document.mimeType,
14041410
processingStatus: document.processingStatus,
1411+
processingQueueToken: document.processingQueueToken,
14051412
processingQueuedAt: document.processingQueuedAt,
14061413
processingStartedAt: document.processingStartedAt,
14071414
processingDeferredUntil: document.processingDeferredUntil,
@@ -1412,6 +1419,7 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
14121419
.where(
14131420
and(
14141421
inArray(document.id, stuckDocIds),
1422+
or(...stuckDocs.map(documentRecoveryGenerationCondition)),
14151423
eq(document.connectorId, connectorId),
14161424
documentProcessingRecoveryCondition(sweepEvaluatedAt, retryCutoff),
14171425
lt(document.uploadedAt, syncStartedAt)

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

Lines changed: 37 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1460,11 +1460,46 @@ describe('processDocumentAsync write guards', () => {
14601460
expect(schedule).not.toHaveBeenCalled()
14611461
})
14621462

1463+
it('reports an unavailable document without parsing or indexing it', async () => {
1464+
dbChainMockFns.limit.mockResolvedValueOnce([])
1465+
expect(
1466+
await processDocumentAsync(
1467+
'knowledge-base-1',
1468+
'document-1',
1469+
PERSISTED_CONTEXT,
1470+
{},
1471+
BILLING_ATTRIBUTION
1472+
)
1473+
).toEqual({ outcome: 'skipped', reason: 'unavailable' })
1474+
expect(mockProcessDocument).not.toHaveBeenCalled()
1475+
})
1476+
1477+
it('reports discarded output when the generation changes before the index commit', async () => {
1478+
armProviderSource()
1479+
dbChainMockFns.limit.mockReset()
1480+
dbChainMockFns.limit
1481+
.mockResolvedValueOnce([PERSISTED_CONTEXT])
1482+
.mockResolvedValueOnce([PERSISTED_PROVENANCE_ROW])
1483+
.mockResolvedValueOnce([])
1484+
expect(
1485+
await processDocumentAsync(
1486+
'knowledge-base-1',
1487+
'document-1',
1488+
PERSISTED_CONTEXT,
1489+
{},
1490+
BILLING_ATTRIBUTION
1491+
)
1492+
).toEqual({ outcome: 'skipped', reason: 'superseded' })
1493+
expect(
1494+
dbChainMockFns.set.mock.calls.some(([value]) => value.processingStatus === 'completed')
1495+
).toBe(false)
1496+
})
1497+
14631498
it('does not parse or reschedule a superseded provider continuation', async () => {
14641499
armProviderSource()
14651500
dbChainMockFns.returning.mockResolvedValueOnce([])
14661501
const schedule = vi.fn()
1467-
await processDocumentAsync(
1502+
const result = await processDocumentAsync(
14681503
'knowledge-base-1',
14691504
'document-1',
14701505
PERSISTED_CONTEXT,
@@ -1477,6 +1512,7 @@ describe('processDocumentAsync write guards', () => {
14771512
scheduleProviderContinuation: schedule,
14781513
}
14791514
)
1515+
expect(result).toEqual({ outcome: 'skipped', reason: 'not_claimed' })
14801516
expect(mockProcessDocument).not.toHaveBeenCalled()
14811517
expect(schedule).not.toHaveBeenCalled()
14821518
})

‎apps/sim/lib/knowledge/documents/processing-recovery-policy.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { document } from '@sim/db/schema'
1+
import { document, outboxEvent } from '@sim/db/schema'
22
import { and, eq, gt, isNotNull, isNull, lt, lte, or, sql } from 'drizzle-orm'
33
import { DOCUMENT_PROCESSING_STALE_THRESHOLD_MS } from '@/lib/knowledge/documents/processing-timeouts.server'
44
import { MAX_PROCESSING_ATTEMPTS, QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowledge/documents/types'
@@ -15,6 +15,11 @@ export function documentProcessingRecoveryCondition(
1515
return and(
1616
sql`${document.processingStatus} IN ('pending', 'processing', 'failed')`,
1717
isNotNull(document.connectorId),
18+
sql`NOT EXISTS (
19+
SELECT 1 FROM ${outboxEvent}
20+
WHERE ${outboxEvent.id} = ${document.processingQueueToken}
21+
AND ${outboxEvent.status} IN ('pending', 'processing')
22+
)`,
1823
isNotNull(document.contentHash),
1924
isNotNull(document.storageKey),
2025
eq(document.userExcluded, false),

0 commit comments

Comments
 (0)