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
16 changes: 16 additions & 0 deletions apps/sim/lib/core/outbox/processor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,10 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
const mocks = vi.hoisted(() => ({
recover: vi.fn(),
reap: vi.fn(),
prune: vi.fn(),
}))
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)
vi.mock('@/lib/core/outbox/retention', () => ({ pruneCompletedOutboxEvents: mocks.prune }))
vi.mock('@/lib/knowledge/documents/processing-recovery', () => ({
recoverKnowledgeDocumentProcessing: mocks.recover,
}))
Expand Down Expand Up @@ -64,6 +66,7 @@ describe('outbox processor recovery', () => {
mockProcessOutboxEvents.mockResolvedValue(result)
mocks.recover.mockResolvedValue(2)
mocks.reap.mockResolvedValue(3)
mocks.prune.mockResolvedValue(4)
})
afterEach(() => vi.useRealTimers())

Expand All @@ -73,6 +76,7 @@ describe('outbox processor recovery', () => {
result,
recoveredDocuments: 0,
reapedBackgroundWork: 3,
prunedEvents: 4,
})
})

Expand All @@ -82,6 +86,17 @@ describe('outbox processor recovery', () => {
result,
recoveredDocuments: 2,
reapedBackgroundWork: 0,
prunedEvents: 4,
})
})

it('retains delivery, recovery and reap results when completed-event pruning fails', async () => {
mocks.prune.mockRejectedValueOnce(new Error('statement timeout'))
await expect(runOutboxProcessor()).resolves.toEqual({
result,
recoveredDocuments: 2,
reapedBackgroundWork: 3,
prunedEvents: 0,
})
})

Expand All @@ -94,6 +109,7 @@ describe('outbox processor recovery', () => {
result,
recoveredDocuments: 0,
reapedBackgroundWork: 3,
prunedEvents: 4,
})
expect(mocks.recover).not.toHaveBeenCalled()
})
Expand Down
12 changes: 11 additions & 1 deletion apps/sim/lib/core/outbox/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import {
OUTBOX_PROCESSOR_MAX_RUNTIME_MS,
OUTBOX_PROCESSOR_RECOVERY_CUTOFF_MS,
} from '@/lib/core/outbox/constants'
import { pruneCompletedOutboxEvents } from '@/lib/core/outbox/retention'
import { type ProcessOutboxResult, processOutboxEvents } from '@/lib/core/outbox/service'
import { DeadlineExceededError } from '@/lib/core/utils/deadline'
import { directGrantOutboxHandlers } from '@/lib/invitations/direct-grant'
Expand Down Expand Up @@ -56,6 +57,7 @@ export interface OutboxProcessorResult {
result: ProcessOutboxResult
recoveredDocuments: number
reapedBackgroundWork: number
prunedEvents: number
}

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

const output = { result, reapedBackgroundWork, recoveredDocuments }
let prunedEvents = 0
try {
prunedEvents = await pruneCompletedOutboxEvents()
} catch (error) {
logger.error('Completed outbox pruning failed', { error: toError(error).message })
}

const output = { result, reapedBackgroundWork, recoveredDocuments, prunedEvents }
logger.info('Outbox processing completed', {
...result,
reapedBackgroundWork,
recoveredDocuments,
prunedEvents,
durationMs: Date.now() - startedAt,
})
return output
Expand Down
121 changes: 121 additions & 0 deletions apps/sim/lib/core/outbox/retention.integration.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
/** Real PostgreSQL retention removes only finished events that nothing reads again. */

import { outboxEvent } from '@sim/db/schema'
import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure'
import { withUtcTimestamps } from '@sim/db/timestamps'
import { generateId } from '@sim/utils/id'
import { sql } from 'drizzle-orm'
import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js'
import postgres from 'postgres'
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'

const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase | undefined }))

vi.mock('@sim/db', () => ({
get db() {
if (!database.current) throw new Error('Outbox PostgreSQL test database is not initialized')
return database.current
},
}))

import {
COMPLETED_OUTBOX_RETENTION_MS,
OUTBOX_PRUNE_BATCH_SIZE,
pruneCompletedOutboxEvents,
} from '@/lib/core/outbox/retention'

const RECOVER = 'knowledge.document.processing.recover'
const STORAGE_CLEANUP = 'knowledge.document.storage.cleanup'
const ADMIN_OPERATION = 'admin.organization-member-operation'
const OCR_CHECKPOINT_EXPIRY = 'knowledge.document.ocr-checkpoint.expire'

describe('completed outbox retention in PostgreSQL', () => {
const schemaName = `outbox_retention_${generateId().replaceAll('-', '')}`
const connection = postgres(
readTestDatabaseUrl(),
withUtcTimestamps({
max: 2,
prepare: false,
fetch_types: false,
connection: { search_path: schemaName },
onnotice: () => {},
})
)
const expired = () => new Date(Date.now() - COMPLETED_OUTBOX_RETENTION_MS - 60_000)

beforeAll(async () => {
await connection`CREATE SCHEMA ${connection(schemaName)}`
await connection`CREATE TABLE outbox_event (LIKE public.outbox_event INCLUDING ALL)`
database.current = drizzle(connection)
})

afterEach(async () => {
await connection`TRUNCATE outbox_event`
})

afterAll(async () => {
try {
await connection`DROP SCHEMA ${connection(schemaName)} CASCADE`
} finally {
await connection.end()
database.current = undefined
}
})

async function seed(eventType: string, status: string, createdAt: Date) {
const id = generateId()
await database.current!.insert(outboxEvent).values({
id,
eventType,
payload: {},
status,
createdAt,
availableAt: createdAt,
})
return id
}

async function remainingIds() {
const rows = await database.current!.select({ id: outboxEvent.id }).from(outboxEvent)
return new Set(rows.map((row) => row.id))
}

it('deletes only expired completed events of prunable types', async () => {
const pruned = [
await seed(RECOVER, 'completed', expired()),
await seed(STORAGE_CLEANUP, 'completed', expired()),
]
const kept = [
await seed(RECOVER, 'completed', new Date(Date.now() - 60_000)),
await seed(STORAGE_CLEANUP, 'pending', expired()),
await seed(RECOVER, 'processing', expired()),
await seed(STORAGE_CLEANUP, 'dead_letter', expired()),
await seed(ADMIN_OPERATION, 'completed', expired()),
await seed(OCR_CHECKPOINT_EXPIRY, 'completed', expired()),
]

expect(await pruneCompletedOutboxEvents()).toBe(pruned.length)
expect(await remainingIds()).toEqual(new Set(kept))
})

it('deletes at most one oldest batch per type per run, and the next run continues', async () => {
const createdAt = expired().toISOString()
for (const [prefix, eventType] of [
['recover', RECOVER],
['storage', STORAGE_CLEANUP],
]) {
await database.current!.execute(sql`
INSERT INTO outbox_event (id, event_type, payload, status, available_at, created_at)
SELECT ${prefix} || ':' || n, ${eventType}, '{}'::json, 'completed',
${createdAt}::timestamp, ${createdAt}::timestamp - n * interval '1 millisecond'
FROM generate_series(0, ${OUTBOX_PRUNE_BATCH_SIZE}::integer) AS n
`)
}
const pending = await seed(RECOVER, 'pending', expired())

expect(await pruneCompletedOutboxEvents()).toBe(2 * OUTBOX_PRUNE_BATCH_SIZE)
expect(await remainingIds()).toEqual(new Set(['recover:0', 'storage:0', pending]))
expect(await pruneCompletedOutboxEvents()).toBe(2)
expect(await remainingIds()).toEqual(new Set([pending]))
})
})
64 changes: 64 additions & 0 deletions apps/sim/lib/core/outbox/retention.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
import { db } from '@sim/db'
import { outboxEvent } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { and, asc, eq, inArray, lt, sql } from 'drizzle-orm'
import { KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT } from '@/lib/knowledge/documents/processing-recovery'
import { KNOWLEDGE_STORAGE_CLEANUP_EVENT } from '@/lib/knowledge/documents/storage-cleanup'

const logger = createLogger('OutboxRetention')

/** How long a completed event stays readable for operators after it was enqueued. */
export const COMPLETED_OUTBOX_RETENTION_MS = 7 * 24 * 60 * 60_000
/**
* Rows deleted per type per run. The processor runs once a minute, so each type drains at most
* 1,000 rows × 1,440 runs = 1.44M rows a day: a steady trickle whose WAL and dead tuples
* autovacuum absorbs, yet five times what recovery can enqueue (200 per run × 1,440 runs).
*/
export const OUTBOX_PRUNE_BATCH_SIZE = 1_000

/**
* Event types whose completed rows nothing reads again: each carries a fresh random id and is
* looked up only while pending or processing. Operation records, idempotency keys and
* deterministic-id dedupe gates, such as the checkpoint expiry events, must outlive completion.
*/
const PRUNABLE_OUTBOX_EVENT_TYPES = [
KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT,
KNOWLEDGE_STORAGE_CLEANUP_EVENT,
] as const

async function pruneCompletedBatch(eventType: string, cutoff: Date): Promise<number> {
return db.transaction(async (tx) => {
await tx.execute(sql`SELECT set_config('statement_timeout', '10000', true)`)
Comment thread
waleedlatif1 marked this conversation as resolved.
const expired = tx
.select({ id: outboxEvent.id })
.from(outboxEvent)
.where(
and(
eq(outboxEvent.eventType, eventType),
lt(outboxEvent.createdAt, cutoff),
eq(outboxEvent.status, 'completed')
)
)
.orderBy(asc(outboxEvent.createdAt))
.limit(OUTBOX_PRUNE_BATCH_SIZE)
.for('update', { skipLocked: true })
Comment thread
waleedlatif1 marked this conversation as resolved.
const deleted = await tx.delete(outboxEvent).where(inArray(outboxEvent.id, expired))
return deleted.count
})
}

/**
* Deletes one oldest-first batch of completed prunable events per type, enqueued before the
* retention window, through the type/creation index. One bounded batch per run keeps a large
* backlog from turning into a burst of deletes; overlapping runs skip each other's locked rows.
* Pending, processing and dead-letter rows are never deleted.
*/
export async function pruneCompletedOutboxEvents(now = new Date()): Promise<number> {
const cutoff = new Date(now.getTime() - COMPLETED_OUTBOX_RETENTION_MS)
let pruned = 0
for (const eventType of PRUNABLE_OUTBOX_EVENT_TYPES) {
pruned += await pruneCompletedBatch(eventType, cutoff)
}
if (pruned > 0) logger.info('Pruned completed outbox events', { pruned })
return pruned
}
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,7 @@ import {
DOCUMENT_RECOVERY_BATCH_SIZE,
KNOWLEDGE_DOCUMENT_RECOVERY_OUTBOX_EVENT,
recoverKnowledgeDocumentProcessing,
recoveryCandidatesQuery,
} from '@/lib/knowledge/documents/processing-recovery'
import {
processDocumentAsync,
Expand All @@ -107,6 +108,11 @@ import { MAX_PROCESSING_ATTEMPTS, QUEUED_DISPATCH_GRACE_MS } from '@/lib/knowled
import { searchScopedKnowledge } from '@/lib/sim-search/indexed/search/scoped-search'
import type { SyncResult } from '@/connectors/types'

interface QueryPlan {
'Shared Hit Blocks': number
'Shared Read Blocks': number
}

const fixtures: ReturnType<typeof createKnowledgeAclFixtureIds>[] = []
const old = () => new Date(Date.now() - QUEUED_DISPATCH_GRACE_MS - 60_000)
async function seed() {
Expand Down Expand Up @@ -956,6 +962,80 @@ describe('independent recovery of retained connector documents', () => {
}
})

it('reads a bounded page however many older documents belong to paused sources', async () => {
const paused = await seed()
const file = await failedFile(paused)
const [original] = await db.select().from(document).where(eq(document.id, file.documentId))
const backlogStart = original.uploadedAt.getTime() - 60 * 60_000
for (let offset = 0; offset < 5_000; offset += 1_000) {
await db.insert(document).values(
Array.from({ length: 1_000 }, (_, index) => ({
...original,
id: generateId(),
externalId: generateId(),
secretProvenanceVersion: null,
uploadedAt: new Date(backlogStart + offset + index),
}))
)
}
await db
.update(knowledgeConnector)
.set({ status: 'paused' })
.where(eq(knowledgeConnector.id, paused.connectorId))
const healthy = await seed()
const recoverable = await failedFile(healthy)
await db.execute(sql`ANALYZE ${document}`)

const plans = await db.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(sql`
EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${recoveryCandidatesQuery(db, new Date(), {
limit: DOCUMENT_RECOVERY_BATCH_SIZE,
attemptedConnectors: new Set(),
blockedKnowledgeBases: new Set(),
})}
`)
const plan = plans[0]['QUERY PLAN'][0].Plan
/** Buffer work, unlike wall-clock time, catches a walk through the paused backlog. */
expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(1_000)

expect(await recoverKnowledgeDocumentProcessing()).toBe(1)
const [event] = await eventsFor(healthy)
expect(event.payload).toMatchObject({ documentId: recoverable.documentId })
await db
.update(knowledgeConnector)
.set({ status: 'paused' })
.where(eq(knowledgeConnector.id, healthy.connectorId))
})

it('recovers a sibling source in the same call when a full batch of older work is live', async () => {
const live = await seed()
const file = await failedFile(live)
const [original] = await db.select().from(document).where(eq(document.id, file.documentId))
await db.insert(document).values(
Array.from({ length: DOCUMENT_RECOVERY_BATCH_SIZE - 1 }, () => ({
...original,
id: generateId(),
externalId: generateId(),
secretProvenanceVersion: null,
}))
)
const sibling = await seed()
const siblingFile = await failedFile(sibling)
fixture.useTrigger = true
fixture.listRuns.mockImplementation(async ({ tag }: { tag: string }) => ({
data: tag === `documentId:${siblingFile.documentId}` ? [] : [{ status: 'QUEUED' }],
hasNextPage: () => false,
}))

expect(await recoverKnowledgeDocumentProcessing()).toBe(1)
expect(await eventsFor(live)).toHaveLength(0)
const [event] = await eventsFor(sibling)
expect(event.payload).toMatchObject({ documentId: siblingFile.documentId })
await db
.update(knowledgeConnector)
.set({ status: 'paused' })
.where(inArray(knowledgeConnector.id, [live.connectorId, sibling.connectorId]))
})

it('recovers another document while an index transaction holds the same KB foreign-key lock', async () => {
const ids = await seed()
const file = await failedFile(ids)
Expand Down
Loading
Loading