Skip to content

Commit 0bc8747

Browse files
committed
improvement(outbox): prune one bounded batch per type per run
1 parent 3af0700 commit 0bc8747

2 files changed

Lines changed: 27 additions & 24 deletions

File tree

‎apps/sim/lib/core/outbox/retention.integration.ts‎

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -98,17 +98,24 @@ describe('completed outbox retention in PostgreSQL', () => {
9898
expect(await remainingIds()).toEqual(new Set(kept))
9999
})
100100

101-
it('drains a backlog larger than one batch within a call', async () => {
101+
it('deletes at most one oldest batch per type per run, and the next run continues', async () => {
102102
const createdAt = expired().toISOString()
103-
await database.current!.execute(sql`
104-
INSERT INTO outbox_event (id, event_type, payload, status, available_at, created_at)
105-
SELECT 'retention:' || n, ${RECOVER}, '{}'::json, 'completed',
106-
${createdAt}::timestamp, ${createdAt}::timestamp - n * interval '1 millisecond'
107-
FROM generate_series(0, ${OUTBOX_PRUNE_BATCH_SIZE}::integer) AS n
108-
`)
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+
}
109114
const pending = await seed(RECOVER, 'pending', expired())
110115

111-
expect(await pruneCompletedOutboxEvents()).toBe(OUTBOX_PRUNE_BATCH_SIZE + 1)
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)
112119
expect(await remainingIds()).toEqual(new Set([pending]))
113120
})
114121
})

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

Lines changed: 12 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,12 @@ const logger = createLogger('OutboxRetention')
99

1010
/** How long a completed event stays readable for operators after it was enqueued. */
1111
export const COMPLETED_OUTBOX_RETENTION_MS = 7 * 24 * 60 * 60_000
12-
export const OUTBOX_PRUNE_BATCH_SIZE = 5_000
13-
const OUTBOX_PRUNE_BUDGET_MS = 10_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
1418

1519
/**
1620
* Event types whose completed rows nothing reads again: each carries a fresh random id and is
@@ -44,24 +48,16 @@ async function pruneCompletedBatch(eventType: string, cutoff: Date): Promise<num
4448
}
4549

4650
/**
47-
* Deletes completed prunable events enqueued before the retention window, one oldest-first
48-
* batch per type in turn through the type/creation index, until the budget is spent. Pending,
49-
* processing and dead-letter rows are never deleted.
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.
5055
*/
5156
export async function pruneCompletedOutboxEvents(now = new Date()): Promise<number> {
52-
const deadlineAt = Date.now() + OUTBOX_PRUNE_BUDGET_MS
5357
const cutoff = new Date(now.getTime() - COMPLETED_OUTBOX_RETENTION_MS)
5458
let pruned = 0
55-
let backlogged: string[] = [...PRUNABLE_OUTBOX_EVENT_TYPES]
56-
while (backlogged.length > 0 && Date.now() < deadlineAt) {
57-
const remaining: string[] = []
58-
for (const eventType of backlogged) {
59-
if (Date.now() >= deadlineAt) break
60-
const deleted = await pruneCompletedBatch(eventType, cutoff)
61-
pruned += deleted
62-
if (deleted === OUTBOX_PRUNE_BATCH_SIZE) remaining.push(eventType)
63-
}
64-
backlogged = remaining
59+
for (const eventType of PRUNABLE_OUTBOX_EVENT_TYPES) {
60+
pruned += await pruneCompletedBatch(eventType, cutoff)
6561
}
6662
if (pruned > 0) logger.info('Pruned completed outbox events', { pruned })
6763
return pruned

0 commit comments

Comments
 (0)