-
Notifications
You must be signed in to change notification settings - Fork 3.8k
improvement(knowledge): scan processing recovery per connector and prune finished outbox events #8330
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
waleedlatif1
merged 2 commits into
staging
from
improvement/processing-recovery-per-connector
Sep 26, 2026
Merged
improvement(knowledge): scan processing recovery per connector and prune finished outbox events #8330
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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])) | ||
| }) | ||
| }) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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)`) | ||
| 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 }) | ||
|
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 | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.