Skip to content

Commit f211b90

Browse files
committed
fix(knowledge): bound recovery database checks and sync staging
2 parents e27c6b5 + f359031 commit f211b90

12 files changed

Lines changed: 436 additions & 85 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: 104 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@ vi.mock('@trigger.dev/sdk', async (importOriginal) => {
4343
const original = await importOriginal<typeof import('@trigger.dev/sdk')>()
4444
return {
4545
...original,
46-
runs: { list: fixture.listRuns },
4746
tasks: { ...original.tasks, batchTrigger: fixture.batchTrigger },
4847
}
4948
})
@@ -96,6 +95,7 @@ import {
9695
retryDocumentProcessing,
9796
} from '@/lib/knowledge/documents/service'
9897
import { MAX_PROCESSING_ATTEMPTS, QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowledge/documents/types'
98+
import type { SyncResult } from '@/connectors/types'
9999

100100
const fixtures: ReturnType<typeof createKnowledgeAclFixtureIds>[] = []
101101
const old = () => new Date(Date.now() - QUEUED_DISPATCH_GRACE_MS - 60_000)
@@ -179,7 +179,110 @@ afterAll(async () => {
179179
await db.$client.end()
180180
})
181181

182+
async function recoverFixture(
183+
ids: ReturnType<typeof createKnowledgeAclFixtureIds>,
184+
mode: 'independent' | 'connector'
185+
) {
186+
if (mode === 'independent') return recoverKnowledgeDocumentProcessing()
187+
const result: SyncResult = {
188+
docsAdded: 0,
189+
docsUpdated: 0,
190+
docsDeleted: 0,
191+
docsUnchanged: 0,
192+
docsSkipped: 0,
193+
docsFailed: 0,
194+
processingDispatch: { requested: 0, accepted: 0, failed: 0 },
195+
}
196+
await sweepStuckDocuments({
197+
connectorId: ids.connectorId,
198+
knowledgeBaseId: ids.knowledgeBaseId,
199+
syncStartedAt: new Date(),
200+
retryCutoff: new Date(Date.now() - 7 * 24 * 60 * 60_000),
201+
billingAttribution: await resolveSystemBillingAttribution(ids.workspaceId),
202+
result,
203+
lease: createContentSyncLease(ids.connectorId, ids.lockId),
204+
})
205+
return result.processingDispatch.requested
206+
}
207+
182208
describe('independent recovery of retained connector documents', () => {
209+
it.each(['independent', 'connector'] as const)(
210+
'%s recovery preserves a job queued beyond the grace period',
211+
async (mode) => {
212+
const ids = await seed()
213+
const file = await failedFile(ids)
214+
await db
215+
.update(document)
216+
.set({ processingStatus: 'pending' })
217+
.where(eq(document.id, file.documentId))
218+
fixture.useTrigger = true
219+
fixture.listRuns.mockResolvedValue({
220+
data: [{ id: 'run-queued', status: 'QUEUED' }],
221+
hasNextPage: () => false,
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('old-fixture-generation')
227+
expect(row.processingRecoveryAfter).not.toBeNull()
228+
expect(await eventsFor(ids)).toHaveLength(0)
229+
await db
230+
.update(knowledgeConnector)
231+
.set({ status: 'paused' })
232+
.where(eq(knowledgeConnector.id, ids.connectorId))
233+
}
234+
)
235+
236+
it.each(['independent', 'connector'] as const)(
237+
'%s recovery rechecks the generation after its remote lookup',
238+
async (mode) => {
239+
const ids = await seed()
240+
const file = await failedFile(ids)
241+
fixture.useTrigger = true
242+
fixture.listRuns.mockImplementation(async () => {
243+
await db
244+
.update(document)
245+
.set({ processingQueueToken: 'replacement-generation' })
246+
.where(eq(document.id, file.documentId))
247+
return { data: [], hasNextPage: () => false }
248+
})
249+
expect(await recoverFixture(ids, mode)).toBe(0)
250+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
251+
expect(row.processingAttempts).toBe(1)
252+
expect(row.processingQueueToken).toBe('replacement-generation')
253+
expect(await eventsFor(ids)).toHaveLength(0)
254+
await db
255+
.update(knowledgeConnector)
256+
.set({ status: 'paused' })
257+
.where(eq(knowledgeConnector.id, ids.connectorId))
258+
}
259+
)
260+
261+
it.each(['pending', 'processing'])(
262+
'does not replace an aged %s outbox continuation',
263+
async (status) => {
264+
const ids = await seed()
265+
const file = await failedFile(ids)
266+
const token = generateId()
267+
await db
268+
.update(document)
269+
.set({ processingQueueToken: token })
270+
.where(eq(document.id, file.documentId))
271+
await db.insert(outboxEvent).values({
272+
id: token,
273+
eventType: 'knowledge.document.processing.resume',
274+
payload: { knowledgeBaseId: ids.knowledgeBaseId, documentId: file.documentId },
275+
status,
276+
availableAt: old(),
277+
})
278+
expect(await recoverFixture(ids, 'independent')).toBe(0)
279+
expect(await recoverFixture(ids, 'connector')).toBe(0)
280+
const [row] = await db.select().from(document).where(eq(document.id, file.documentId))
281+
expect(row.processingAttempts).toBe(1)
282+
expect(row.processingQueueToken).toBe(token)
283+
}
284+
)
285+
183286
it.each(['manual', 'redelivery'])(
184287
'respects a concurrent liveness cooldown before %s replacement',
185288
async (path) => {

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

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,13 @@ import {
2727
persistSkippedRetryHashes,
2828
updateDocument,
2929
} from '@/lib/knowledge/connectors/sync-persistence'
30+
import { documentProcessingRecoveryCondition } from '@/lib/knowledge/documents/processing-recovery-policy'
3031
import {
3132
DOCUMENT_LIVENESS_BATCH_SIZE,
3233
documentProcessingSnapshotCondition,
3334
findAbandonedDocumentProcessing,
3435
processingSnapshotColumns,
35-
} from '@/lib/knowledge/documents/processing-liveness'
36-
import { documentProcessingRecoveryCondition } from '@/lib/knowledge/documents/processing-recovery-policy'
36+
} from '@/lib/knowledge/documents/processing-recovery-queue'
3737
import { DOCUMENT_PROCESSING_STALE_THRESHOLD_MS } from '@/lib/knowledge/documents/processing-timeouts.server'
3838
import type { DocumentData } from '@/lib/knowledge/documents/service'
3939
import { isTriggerAvailable, processDocumentsWithQueue } from '@/lib/knowledge/documents/service'
@@ -1341,7 +1341,6 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
13411341
filename: document.filename,
13421342
fileSize: document.fileSize,
13431343
mimeType: document.mimeType,
1344-
uploadedAt: document.uploadedAt,
13451344
})
13461345
.from(document)
13471346
.where(
@@ -1403,7 +1402,6 @@ export async function sweepStuckDocuments(input: SweepStuckDocumentsInput): Prom
14031402
filename: document.filename,
14041403
fileSize: document.fileSize,
14051404
mimeType: document.mimeType,
1406-
uploadedAt: document.uploadedAt,
14071405
})
14081406
.from(document)
14091407
.where(

‎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),

apps/sim/lib/knowledge/documents/processing-liveness-transport.test.ts renamed to apps/sim/lib/knowledge/documents/processing-recovery-queue-transport.test.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,11 @@ import { resetInsideTriggerRunForTests } from '@/lib/core/config/trigger-runtime
1111
import {
1212
type DocumentProcessingSnapshot,
1313
findAbandonedDocumentProcessing,
14-
} from '@/lib/knowledge/documents/processing-liveness'
14+
} from '@/lib/knowledge/documents/processing-recovery-queue'
1515

1616
const snapshot: DocumentProcessingSnapshot = {
1717
id: 'doc-1',
18+
uploadedAt: new Date('2026-09-01T00:00:00Z'),
1819
processingStatus: 'pending',
1920
processingQueueToken: 'generation-1',
2021
processingQueuedAt: new Date('2026-09-01T00:00:00Z'),
@@ -92,6 +93,9 @@ describe('document liveness SDK transport', () => {
9293
expect(fetch).toHaveBeenCalledOnce()
9394
const url = new URL(fetch.mock.calls[0][0])
9495
expect(url.searchParams.get('page[size]')).toBe('1')
96+
expect(url.searchParams.get('filter[createdAt][from]')).toBe(
97+
String(new Date('2026-08-31T20:00:00Z').getTime())
98+
)
9599
expect(url.searchParams.get('filter[tag]')).toBe('documentId:doc-1')
96100
}
97101
)

0 commit comments

Comments
 (0)