Skip to content

Commit 48a13cb

Browse files
authored
improvement(knowledge): scan processing recovery per connector and prune finished outbox events (#8330)
- Processing recovery discovered candidates by walking the global `(uploaded_at, id)` recovery index, joining each document to its knowledge base and connector, and rejecting it only afterwards. Retained inputs of paused or dormant sources sit at the front of that order, so every call read through all of them first. Follow-up batches (`NOT IN` attempted connectors) walked the same range again - Discovery now starts from the eligible connectors and reads each one's oldest candidates through a new partial index `doc_connector_processing_recovery_idx (connector_id, uploaded_at, id)` with a `CROSS JOIN LATERAL`, then takes the oldest overall. Paused sources cost nothing. Batch size, the follow-up batches for refused connectors, blocked knowledge bases, the lock-time recheck, and liveness are unchanged, so recovery picks the same candidates as before - The old `doc_processing_recovery_idx` stays for now and is dropped in a later contract migration - Outbox retention: `outbox_event` never deleted completed rows. The outbox processor now prunes `completed` rows older than 7 days in bounded, `SKIP LOCKED` batches, only for `knowledge.document.processing.recover` and `knowledge.document.storage.cleanup`. Both have random ids, and nothing reads their completed rows. Every other event type is kept, including idempotency-keyed ones and the checkpoint expiry events, whose completed rows are read back by id, as are pending, processing, and dead-letter rows
1 parent 6d77ae1 commit 48a13cb

11 files changed

Lines changed: 29817 additions & 32 deletions

File tree

‎apps/sim/lib/core/outbox/processor.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,10 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
55
const mocks = vi.hoisted(() => ({
66
recover: vi.fn(),
77
reap: vi.fn(),
8+
prune: vi.fn(),
89
}))
910
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)
11+
vi.mock('@/lib/core/outbox/retention', () => ({ pruneCompletedOutboxEvents: mocks.prune }))
1012
vi.mock('@/lib/knowledge/documents/processing-recovery', () => ({
1113
recoverKnowledgeDocumentProcessing: mocks.recover,
1214
}))
@@ -64,6 +66,7 @@ describe('outbox processor recovery', () => {
6466
mockProcessOutboxEvents.mockResolvedValue(result)
6567
mocks.recover.mockResolvedValue(2)
6668
mocks.reap.mockResolvedValue(3)
69+
mocks.prune.mockResolvedValue(4)
6770
})
6871
afterEach(() => vi.useRealTimers())
6972

@@ -73,6 +76,7 @@ describe('outbox processor recovery', () => {
7376
result,
7477
recoveredDocuments: 0,
7578
reapedBackgroundWork: 3,
79+
prunedEvents: 4,
7680
})
7781
})
7882

@@ -82,6 +86,17 @@ describe('outbox processor recovery', () => {
8286
result,
8387
recoveredDocuments: 2,
8488
reapedBackgroundWork: 0,
89+
prunedEvents: 4,
90+
})
91+
})
92+
93+
it('retains delivery, recovery and reap results when completed-event pruning fails', async () => {
94+
mocks.prune.mockRejectedValueOnce(new Error('statement timeout'))
95+
await expect(runOutboxProcessor()).resolves.toEqual({
96+
result,
97+
recoveredDocuments: 2,
98+
reapedBackgroundWork: 3,
99+
prunedEvents: 0,
85100
})
86101
})
87102

@@ -94,6 +109,7 @@ describe('outbox processor recovery', () => {
94109
result,
95110
recoveredDocuments: 0,
96111
reapedBackgroundWork: 3,
112+
prunedEvents: 4,
97113
})
98114
expect(mocks.recover).not.toHaveBeenCalled()
99115
})

‎apps/sim/lib/core/outbox/processor.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import {
1111
OUTBOX_PROCESSOR_MAX_RUNTIME_MS,
1212
OUTBOX_PROCESSOR_RECOVERY_CUTOFF_MS,
1313
} from '@/lib/core/outbox/constants'
14+
import { pruneCompletedOutboxEvents } from '@/lib/core/outbox/retention'
1415
import { type ProcessOutboxResult, processOutboxEvents } from '@/lib/core/outbox/service'
1516
import { DeadlineExceededError } from '@/lib/core/utils/deadline'
1617
import { directGrantOutboxHandlers } from '@/lib/invitations/direct-grant'
@@ -56,6 +57,7 @@ export interface OutboxProcessorResult {
5657
result: ProcessOutboxResult
5758
recoveredDocuments: number
5859
reapedBackgroundWork: number
60+
prunedEvents: number
5961
}
6062

6163
/** Processes one bounded batch and its recovery work in either the worker or self-hosted cron. */
@@ -92,11 +94,19 @@ export async function runOutboxProcessor(): Promise<OutboxProcessorResult> {
9294
logger.error('Background-work reap failed', { error: toError(error).message })
9395
}
9496

95-
const output = { result, reapedBackgroundWork, recoveredDocuments }
97+
let prunedEvents = 0
98+
try {
99+
prunedEvents = await pruneCompletedOutboxEvents()
100+
} catch (error) {
101+
logger.error('Completed outbox pruning failed', { error: toError(error).message })
102+
}
103+
104+
const output = { result, reapedBackgroundWork, recoveredDocuments, prunedEvents }
96105
logger.info('Outbox processing completed', {
97106
...result,
98107
reapedBackgroundWork,
99108
recoveredDocuments,
109+
prunedEvents,
100110
durationMs: Date.now() - startedAt,
101111
})
102112
return output
Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
1+
/** Real PostgreSQL retention removes only finished events that nothing reads again. */
2+
3+
import { outboxEvent } from '@sim/db/schema'
4+
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
5+
import { withUtcTimestamps } from '@sim/db/timestamps'
6+
import { generateId } from '@sim/utils/id'
7+
import { sql } from 'drizzle-orm'
8+
import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js'
9+
import postgres from 'postgres'
10+
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
11+
12+
const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase | undefined }))
13+
14+
vi.mock('@sim/db', () => ({
15+
get db() {
16+
if (!database.current) throw new Error('Outbox PostgreSQL test database is not initialized')
17+
return database.current
18+
},
19+
}))
20+
21+
import {
22+
COMPLETED_OUTBOX_RETENTION_MS,
23+
OUTBOX_PRUNE_BATCH_SIZE,
24+
pruneCompletedOutboxEvents,
25+
} from '@/lib/core/outbox/retention'
26+
27+
const RECOVER = 'knowledge.document.processing.recover'
28+
const STORAGE_CLEANUP = 'knowledge.document.storage.cleanup'
29+
const ADMIN_OPERATION = 'admin.organization-member-operation'
30+
const OCR_CHECKPOINT_EXPIRY = 'knowledge.document.ocr-checkpoint.expire'
31+
32+
describe('completed outbox retention in PostgreSQL', () => {
33+
const schemaName = `outbox_retention_${generateId().replaceAll('-', '')}`
34+
const connection = postgres(
35+
readTestDatabaseUrl(),
36+
withUtcTimestamps({
37+
max: 2,
38+
prepare: false,
39+
fetch_types: false,
40+
connection: { search_path: schemaName },
41+
onnotice: () => {},
42+
})
43+
)
44+
const expired = () => new Date(Date.now() - COMPLETED_OUTBOX_RETENTION_MS - 60_000)
45+
46+
beforeAll(async () => {
47+
await connection`CREATE SCHEMA ${connection(schemaName)}`
48+
await connection`CREATE TABLE outbox_event (LIKE public.outbox_event INCLUDING ALL)`
49+
database.current = drizzle(connection)
50+
})
51+
52+
afterEach(async () => {
53+
await connection`TRUNCATE outbox_event`
54+
})
55+
56+
afterAll(async () => {
57+
try {
58+
await connection`DROP SCHEMA ${connection(schemaName)} CASCADE`
59+
} finally {
60+
await connection.end()
61+
database.current = undefined
62+
}
63+
})
64+
65+
async function seed(eventType: string, status: string, createdAt: Date) {
66+
const id = generateId()
67+
await database.current!.insert(outboxEvent).values({
68+
id,
69+
eventType,
70+
payload: {},
71+
status,
72+
createdAt,
73+
availableAt: createdAt,
74+
})
75+
return id
76+
}
77+
78+
async function remainingIds() {
79+
const rows = await database.current!.select({ id: outboxEvent.id }).from(outboxEvent)
80+
return new Set(rows.map((row) => row.id))
81+
}
82+
83+
it('deletes only expired completed events of prunable types', async () => {
84+
const pruned = [
85+
await seed(RECOVER, 'completed', expired()),
86+
await seed(STORAGE_CLEANUP, 'completed', expired()),
87+
]
88+
const kept = [
89+
await seed(RECOVER, 'completed', new Date(Date.now() - 60_000)),
90+
await seed(STORAGE_CLEANUP, 'pending', expired()),
91+
await seed(RECOVER, 'processing', expired()),
92+
await seed(STORAGE_CLEANUP, 'dead_letter', expired()),
93+
await seed(ADMIN_OPERATION, 'completed', expired()),
94+
await seed(OCR_CHECKPOINT_EXPIRY, 'completed', expired()),
95+
]
96+
97+
expect(await pruneCompletedOutboxEvents()).toBe(pruned.length)
98+
expect(await remainingIds()).toEqual(new Set(kept))
99+
})
100+
101+
it('deletes at most one oldest batch per type per run, and the next run continues', async () => {
102+
const createdAt = expired().toISOString()
103+
for (const [prefix, eventType] of [
104+
['recover', RECOVER],
105+
['storage', STORAGE_CLEANUP],
106+
]) {
107+
await database.current!.execute(sql`
108+
INSERT INTO outbox_event (id, event_type, payload, status, available_at, created_at)
109+
SELECT ${prefix} || ':' || n, ${eventType}, '{}'::json, 'completed',
110+
${createdAt}::timestamp, ${createdAt}::timestamp - n * interval '1 millisecond'
111+
FROM generate_series(0, ${OUTBOX_PRUNE_BATCH_SIZE}::integer) AS n
112+
`)
113+
}
114+
const pending = await seed(RECOVER, 'pending', expired())
115+
116+
expect(await pruneCompletedOutboxEvents()).toBe(2 * OUTBOX_PRUNE_BATCH_SIZE)
117+
expect(await remainingIds()).toEqual(new Set(['recover:0', 'storage:0', pending]))
118+
expect(await pruneCompletedOutboxEvents()).toBe(2)
119+
expect(await remainingIds()).toEqual(new Set([pending]))
120+
})
121+
})
Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
import { db } from '@sim/db'
2+
import { outboxEvent } from '@sim/db/schema'
3+
import { createLogger } from '@sim/logger'
4+
import { and, asc, eq, inArray, lt, sql } from 'drizzle-orm'
5+
import { KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT } from '@/lib/knowledge/documents/processing-recovery'
6+
import { KNOWLEDGE_STORAGE_CLEANUP_EVENT } from '@/lib/knowledge/documents/storage-cleanup'
7+
8+
const logger = createLogger('OutboxRetention')
9+
10+
/** How long a completed event stays readable for operators after it was enqueued. */
11+
export const COMPLETED_OUTBOX_RETENTION_MS = 7 * 24 * 60 * 60_000
12+
/**
13+
* Rows deleted per type per run. The processor runs once a minute, so each type drains at most
14+
* 1,000 rows × 1,440 runs = 1.44M rows a day: a steady trickle whose WAL and dead tuples
15+
* autovacuum absorbs, yet five times what recovery can enqueue (200 per run × 1,440 runs).
16+
*/
17+
export const OUTBOX_PRUNE_BATCH_SIZE = 1_000
18+
19+
/**
20+
* Event types whose completed rows nothing reads again: each carries a fresh random id and is
21+
* looked up only while pending or processing. Operation records, idempotency keys and
22+
* deterministic-id dedupe gates, such as the checkpoint expiry events, must outlive completion.
23+
*/
24+
const PRUNABLE_OUTBOX_EVENT_TYPES = [
25+
KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT,
26+
KNOWLEDGE_STORAGE_CLEANUP_EVENT,
27+
] as const
28+
29+
async function pruneCompletedBatch(eventType: string, cutoff: Date): Promise<number> {
30+
return db.transaction(async (tx) => {
31+
await tx.execute(sql`SELECT set_config('statement_timeout', '10000', true)`)
32+
const expired = tx
33+
.select({ id: outboxEvent.id })
34+
.from(outboxEvent)
35+
.where(
36+
and(
37+
eq(outboxEvent.eventType, eventType),
38+
lt(outboxEvent.createdAt, cutoff),
39+
eq(outboxEvent.status, 'completed')
40+
)
41+
)
42+
.orderBy(asc(outboxEvent.createdAt))
43+
.limit(OUTBOX_PRUNE_BATCH_SIZE)
44+
.for('update', { skipLocked: true })
45+
const deleted = await tx.delete(outboxEvent).where(inArray(outboxEvent.id, expired))
46+
return deleted.count
47+
})
48+
}
49+
50+
/**
51+
* Deletes one oldest-first batch of completed prunable events per type, enqueued before the
52+
* retention window, through the type/creation index. One bounded batch per run keeps a large
53+
* backlog from turning into a burst of deletes; overlapping runs skip each other's locked rows.
54+
* Pending, processing and dead-letter rows are never deleted.
55+
*/
56+
export async function pruneCompletedOutboxEvents(now = new Date()): Promise<number> {
57+
const cutoff = new Date(now.getTime() - COMPLETED_OUTBOX_RETENTION_MS)
58+
let pruned = 0
59+
for (const eventType of PRUNABLE_OUTBOX_EVENT_TYPES) {
60+
pruned += await pruneCompletedBatch(eventType, cutoff)
61+
}
62+
if (pruned > 0) logger.info('Pruned completed outbox events', { pruned })
63+
return pruned
64+
}

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

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@ import {
9797
DOCUMENT_RECOVERY_BATCH_SIZE,
9898
KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT,
9999
recoverKnowledgeDocumentProcessing,
100+
recoveryCandidatesQuery,
100101
} from '@/lib/knowledge/documents/processing-recovery'
101102
import {
102103
processDocumentAsync,
@@ -107,6 +108,11 @@ import { MAX_PROCESSING_ATTEMPTS, QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowled
107108
import { searchScopedKnowledge } from '@/lib/sim-search/indexed/search/scoped-search'
108109
import type { SyncResult } from '@/connectors/types'
109110

111+
interface QueryPlan {
112+
'Shared Hit Blocks': number
113+
'Shared Read Blocks': number
114+
}
115+
110116
const fixtures: ReturnType<typeof createKnowledgeAclFixtureIds>[] = []
111117
const old = () => new Date(Date.now() - QUEUED_DISPATCH_GRACE_MS - 60_000)
112118
async function seed() {
@@ -956,6 +962,80 @@ describe('independent recovery of retained connector documents', () => {
956962
}
957963
})
958964

965+
it('reads a bounded page however many older documents belong to paused sources', async () => {
966+
const paused = await seed()
967+
const file = await failedFile(paused)
968+
const [original] = await db.select().from(document).where(eq(document.id, file.documentId))
969+
const backlogStart = original.uploadedAt.getTime() - 60 * 60_000
970+
for (let offset = 0; offset < 5_000; offset += 1_000) {
971+
await db.insert(document).values(
972+
Array.from({ length: 1_000 }, (_, index) => ({
973+
...original,
974+
id: generateId(),
975+
externalId: generateId(),
976+
secretProvenanceVersion: null,
977+
uploadedAt: new Date(backlogStart + offset + index),
978+
}))
979+
)
980+
}
981+
await db
982+
.update(knowledgeConnector)
983+
.set({ status: 'paused' })
984+
.where(eq(knowledgeConnector.id, paused.connectorId))
985+
const healthy = await seed()
986+
const recoverable = await failedFile(healthy)
987+
await db.execute(sql`ANALYZE ${document}`)
988+
989+
const plans = await db.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(sql`
990+
EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${recoveryCandidatesQuery(db, new Date(), {
991+
limit: DOCUMENT_RECOVERY_BATCH_SIZE,
992+
attemptedConnectors: new Set(),
993+
blockedKnowledgeBases: new Set(),
994+
})}
995+
`)
996+
const plan = plans[0]['QUERY PLAN'][0].Plan
997+
/** Buffer work, unlike wall-clock time, catches a walk through the paused backlog. */
998+
expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(1_000)
999+
1000+
expect(await recoverKnowledgeDocumentProcessing()).toBe(1)
1001+
const [event] = await eventsFor(healthy)
1002+
expect(event.payload).toMatchObject({ documentId: recoverable.documentId })
1003+
await db
1004+
.update(knowledgeConnector)
1005+
.set({ status: 'paused' })
1006+
.where(eq(knowledgeConnector.id, healthy.connectorId))
1007+
})
1008+
1009+
it('recovers a sibling source in the same call when a full batch of older work is live', async () => {
1010+
const live = await seed()
1011+
const file = await failedFile(live)
1012+
const [original] = await db.select().from(document).where(eq(document.id, file.documentId))
1013+
await db.insert(document).values(
1014+
Array.from({ length: DOCUMENT_RECOVERY_BATCH_SIZE - 1 }, () => ({
1015+
...original,
1016+
id: generateId(),
1017+
externalId: generateId(),
1018+
secretProvenanceVersion: null,
1019+
}))
1020+
)
1021+
const sibling = await seed()
1022+
const siblingFile = await failedFile(sibling)
1023+
fixture.useTrigger = true
1024+
fixture.listRuns.mockImplementation(async ({ tag }: { tag: string }) => ({
1025+
data: tag === `documentId:${siblingFile.documentId}` ? [] : [{ status: 'QUEUED' }],
1026+
hasNextPage: () => false,
1027+
}))
1028+
1029+
expect(await recoverKnowledgeDocumentProcessing()).toBe(1)
1030+
expect(await eventsFor(live)).toHaveLength(0)
1031+
const [event] = await eventsFor(sibling)
1032+
expect(event.payload).toMatchObject({ documentId: siblingFile.documentId })
1033+
await db
1034+
.update(knowledgeConnector)
1035+
.set({ status: 'paused' })
1036+
.where(inArray(knowledgeConnector.id, [live.connectorId, sibling.connectorId]))
1037+
})
1038+
9591039
it('recovers another document while an index transaction holds the same KB foreign-key lock', async () => {
9601040
const ids = await seed()
9611041
const file = await failedFile(ids)

0 commit comments

Comments
 (0)