Skip to content

Commit a9a825d

Browse files
committed
improvement(outbox): only start the processor when work is due or maintenance is scheduled
1 parent 8fc4c1a commit a9a825d

9 files changed

Lines changed: 204 additions & 12 deletions

File tree

‎apps/sim/app/api/webhooks/outbox/process/route.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,9 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
2020
const requestId = generateRequestId()
2121
try {
2222
const accepted = await enqueueOutboxProcessor()
23+
if (accepted.backend === null) {
24+
return NextResponse.json({ success: true, requestId, triggered: false })
25+
}
2326
if (accepted.backend === 'trigger-dev') {
2427
logger.info('Outbox processor accepted', { jobId: accepted.jobId })
2528
return NextResponse.json(

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,3 +6,5 @@ export const OUTBOX_PROCESSOR_INTERVAL_MS = 60_000
66
export const OUTBOX_PROCESSOR_CONCURRENCY = Math.ceil(
77
(OUTBOX_PROCESSOR_MAX_DURATION_SECONDS * 1000) / OUTBOX_PROCESSOR_INTERVAL_MS
88
)
9+
/** Every Nth scheduled tick runs the processor's recovery, reaping and pruning even with nothing due. */
10+
export const OUTBOX_MAINTENANCE_EVERY_INTERVALS = 5

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

Lines changed: 54 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,45 @@
11
import { asyncJobsRegionMock } from '@sim/testing/mocks/async-jobs-region.mock'
22
import { setEnvFlags } from '@sim/testing/mocks/env-flags.mock'
3+
import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock'
34
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
45

56
const hoisted = vi.hoisted(() => ({
67
processor: vi.fn(),
78
}))
89
vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock)
910
vi.mock('@/lib/core/outbox/processor', () => ({ runOutboxProcessor: hoisted.processor }))
11+
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)
1012

1113
import { tasks } from '@trigger.dev/sdk'
1214
import { enqueueOutboxProcessor } from '@/lib/core/outbox/enqueue'
1315

14-
const mocks = { ...hoisted, trigger: vi.mocked(tasks.trigger) }
16+
const mocks = {
17+
...hoisted,
18+
trigger: vi.mocked(tasks.trigger),
19+
hasDueWork: outboxServiceMockFns.mockHasDueOutboxWork,
20+
}
21+
22+
/** 12:34 is not a maintenance minute (754 % 5 = 4); 12:35 is. */
23+
const IDLE_MINUTE = new Date('2026-09-16T12:34:45Z')
24+
const MAINTENANCE_MINUTE = new Date('2026-09-16T12:35:10Z')
25+
26+
const BACKENDS = [
27+
{ name: 'Trigger.dev', isTriggerDevEnabled: true },
28+
{ name: 'self-hosted inline', isTriggerDevEnabled: false },
29+
] as const
30+
31+
function processorStarted() {
32+
return mocks.trigger.mock.calls.length + mocks.processor.mock.calls.length
33+
}
1534

1635
describe('outbox processor enqueue', () => {
1736
beforeEach(() => {
1837
vi.useFakeTimers()
19-
vi.setSystemTime(new Date('2026-09-16T12:34:45Z'))
38+
vi.setSystemTime(IDLE_MINUTE)
2039
setEnvFlags({ isTriggerDevEnabled: true })
2140
mocks.trigger.mockResolvedValue({ id: 'run-1' })
41+
mocks.hasDueWork.mockReset()
42+
mocks.hasDueWork.mockResolvedValue(true)
2243
})
2344
afterEach(() => vi.useRealTimers())
2445

@@ -49,4 +70,35 @@ describe('outbox processor enqueue', () => {
4970
await expect(enqueueOutboxProcessor()).resolves.toEqual({ backend: 'inline', output })
5071
expect(mocks.trigger).not.toHaveBeenCalled()
5172
})
73+
74+
it.each(BACKENDS)(
75+
'starts no $name processor on an idle queue outside the maintenance window',
76+
async ({ isTriggerDevEnabled }) => {
77+
setEnvFlags({ isTriggerDevEnabled })
78+
mocks.hasDueWork.mockResolvedValue(false)
79+
await expect(enqueueOutboxProcessor()).resolves.toEqual({ backend: null, triggered: false })
80+
expect(processorStarted()).toBe(0)
81+
}
82+
)
83+
84+
it.each(BACKENDS)(
85+
'starts the $name processor on an idle queue in the maintenance window',
86+
async ({ isTriggerDevEnabled }) => {
87+
setEnvFlags({ isTriggerDevEnabled })
88+
vi.setSystemTime(MAINTENANCE_MINUTE)
89+
mocks.hasDueWork.mockResolvedValue(false)
90+
await enqueueOutboxProcessor()
91+
expect(processorStarted()).toBe(1)
92+
}
93+
)
94+
95+
it.each(BACKENDS)(
96+
'starts the $name processor when the work check fails',
97+
async ({ isTriggerDevEnabled }) => {
98+
setEnvFlags({ isTriggerDevEnabled })
99+
mocks.hasDueWork.mockRejectedValue(new Error('connection refused'))
100+
await enqueueOutboxProcessor()
101+
expect(processorStarted()).toBe(1)
102+
}
103+
)
52104
})

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

Lines changed: 31 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,48 @@
1+
import { createLogger } from '@sim/logger'
2+
import { getErrorMessage } from '@sim/utils/errors'
13
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
24
import {
5+
OUTBOX_MAINTENANCE_EVERY_INTERVALS,
36
OUTBOX_PROCESSOR_INTERVAL_MS,
47
OUTBOX_PROCESSOR_MAX_DURATION_SECONDS,
58
} from '@/lib/core/outbox/constants'
69
import type { OutboxProcessorResult } from '@/lib/core/outbox/processor'
10+
import { hasDueOutboxWork } from '@/lib/core/outbox/service'
711
import type { processOutboxTask } from '@/background/process-outbox'
812

13+
const logger = createLogger('OutboxProcessorEnqueue')
14+
915
type OutboxProcessorEnqueueResult =
1016
| { backend: 'trigger-dev'; jobId: string }
1117
| { backend: 'inline'; output: OutboxProcessorResult }
18+
| { backend: null; triggered: false }
19+
20+
/**
21+
* Every tick in the maintenance window runs, so document recovery, background-work reaping and
22+
* pruning keep their cadence on an idle queue. Other ticks run only when events are due or a lease
23+
* is stale; a gate that cannot answer runs anyway, since billing, seat sync, invitations and
24+
* document dispatch all ride the outbox.
25+
*/
26+
async function shouldRunOutboxProcessor(now: Date, scheduleWindow: number): Promise<boolean> {
27+
if (scheduleWindow % OUTBOX_MAINTENANCE_EVERY_INTERVALS === 0) return true
28+
try {
29+
return await hasDueOutboxWork(now)
30+
} catch (error) {
31+
logger.warn('Outbox work check failed; running the processor anyway', {
32+
error: getErrorMessage(error),
33+
})
34+
return true
35+
}
36+
}
1237

1338
/** The database owns delivery state; the cron request waits only for durable worker acceptance. */
1439
export async function enqueueOutboxProcessor(): Promise<OutboxProcessorEnqueueResult> {
40+
const now = new Date()
41+
const scheduleWindow = Math.floor(now.getTime() / OUTBOX_PROCESSOR_INTERVAL_MS)
42+
if (!(await shouldRunOutboxProcessor(now, scheduleWindow))) {
43+
return { backend: null, triggered: false }
44+
}
45+
1546
if (!isTriggerDevEnabled) {
1647
const { runOutboxProcessor } = await import('@/lib/core/outbox/processor')
1748
return { backend: 'inline', output: await runOutboxProcessor() }
@@ -21,7 +52,6 @@ export async function enqueueOutboxProcessor(): Promise<OutboxProcessorEnqueueRe
2152
import('@trigger.dev/sdk'),
2253
import('@/lib/core/async-jobs/region'),
2354
])
24-
const scheduleWindow = Math.floor(Date.now() / OUTBOX_PROCESSOR_INTERVAL_MS)
2555
const handle = await tasks.trigger<typeof processOutboxTask>('process-outbox', undefined, {
2656
idempotencyKey: `process-outbox:${scheduleWindow}`,
2757
idempotencyKeyTTL: '5m',

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

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,18 @@
1+
import { db } from '@sim/db'
12
import { outboxEvent } from '@sim/db/schema'
2-
import { sql } from 'drizzle-orm'
3+
import { and, eq, exists, lte, or, type SQL, sql } from 'drizzle-orm'
34

45
const MAX_READY_EVENT_TYPES = 128
6+
/** How long a `processing` lease may go without a terminal write before the reaper reclaims it. */
7+
export const STUCK_PROCESSING_THRESHOLD_MS = 10 * 60 * 1000
8+
9+
/** A `processing` row whose lease the reaper may reclaim at `now`. */
10+
export function isStuckProcessing(now: Date): SQL | undefined {
11+
return and(
12+
eq(outboxEvent.status, 'processing'),
13+
lte(outboxEvent.lockedAt, new Date(now.getTime() - STUCK_PROCESSING_THRESHOLD_MS))
14+
)
15+
}
516

617
/**
718
* Walks the pending index one type at a time, reading only its earliest availability.
@@ -42,3 +53,18 @@ export function readyEventTypesQuery(now: Date) {
4253
LIMIT ${MAX_READY_EVENT_TYPES}
4354
`
4455
}
56+
57+
/**
58+
* One row, `due`: whether `processOutboxEvents` would act at `now`. The pending leg is the
59+
* claim phase's own discovery walk, so it reads only each type's head and cannot be planned as a
60+
* scan of future or completed rows; the lease leg is the reaper's predicate, over the few rows in
61+
* `processing`.
62+
*/
63+
export function dueOutboxWorkQuery(now: Date) {
64+
const staleLease = db
65+
.select({ id: outboxEvent.id })
66+
.from(outboxEvent)
67+
.where(isStuckProcessing(now))
68+
.limit(1)
69+
return sql`SELECT ${or(sql`EXISTS (${readyEventTypesQuery(now)})`, exists(staleLease))} AS due`
70+
}

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

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,11 @@ const logger = createLogger('OutboxRetention')
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
1212
/**
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).
13+
* Rows deleted per type per run. The processor runs on every fifth minute's maintenance tick and
14+
* on any other minute with due work, so each type drains at least 1,000 rows × 288 runs = 288k
15+
* rows a day and at most 1.44M: a steady trickle whose WAL and dead tuples autovacuum absorbs.
16+
* Recovery enqueues only inside a run (at most 200), so pruning keeps five times its pace at any
17+
* run count.
1618
*/
1719
export const OUTBOX_PRUNE_BATCH_SIZE = 1_000
1820

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

Lines changed: 64 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,13 @@ vi.mock('@sim/db', () => ({
1919
},
2020
}))
2121

22-
import { readyEventTypesQuery } from '@/lib/core/outbox/queries'
2322
import {
23+
dueOutboxWorkQuery,
24+
readyEventTypesQuery,
25+
STUCK_PROCESSING_THRESHOLD_MS,
26+
} from '@/lib/core/outbox/queries'
27+
import {
28+
hasDueOutboxWork,
2429
type OutboxHandler,
2530
processOutboxEvents,
2631
withOutboxHandlerTimeout,
@@ -386,4 +391,62 @@ describe('outbox scheduling in PostgreSQL', () => {
386391
const second = await processOutboxEvents({}, { batchSize: 0 })
387392
expect(second.reaped).toBe(5)
388393
})
394+
395+
describe('due-work gate', () => {
396+
const now = new Date('2026-09-16T12:34:00.000Z')
397+
398+
async function insertRow(
399+
status: 'pending' | 'processing' | 'completed' | 'dead_letter',
400+
times: { availableAt?: Date; lockedAt?: Date | null }
401+
) {
402+
eventTypes.add('test.outbox.gate')
403+
await db.insert(outboxEvent).values({
404+
id: generateId(),
405+
eventType: 'test.outbox.gate',
406+
payload: {},
407+
status,
408+
createdAt: new Date(now.getTime() - 60 * 60_000),
409+
availableAt: times.availableAt ?? new Date(now.getTime() - 60 * 60_000),
410+
lockedAt: times.lockedAt ?? null,
411+
})
412+
}
413+
414+
it('reports no work for future, freshly leased, completed and dead-lettered rows', async () => {
415+
expect(await hasDueOutboxWork(now)).toBe(false)
416+
await insertRow('pending', { availableAt: new Date(now.getTime() + 1) })
417+
await insertRow('processing', {
418+
lockedAt: new Date(now.getTime() - STUCK_PROCESSING_THRESHOLD_MS + 1),
419+
})
420+
await insertRow('completed', {})
421+
await insertRow('dead_letter', {})
422+
expect(await hasDueOutboxWork(now)).toBe(false)
423+
})
424+
425+
it('reports a pending event due exactly now, as the claim phase would take it', async () => {
426+
await insertRow('pending', { availableAt: now })
427+
expect(await hasDueOutboxWork(now)).toBe(true)
428+
})
429+
430+
it('reports a lease exactly at the stale threshold, as the reaper would reclaim it', async () => {
431+
await insertRow('processing', {
432+
lockedAt: new Date(now.getTime() - STUCK_PROCESSING_THRESHOLD_MS),
433+
})
434+
expect(await hasDueOutboxWork(now)).toBe(true)
435+
})
436+
437+
it('probes both legs through indexes beside a large completed and future backlog', async () => {
438+
await seedBacklog('test.outbox.gate', 50_000, 'completed')
439+
await seedBacklog('test.outbox.gate', 50_000, 'pending', new Date(Date.now() + 60 * 60_000))
440+
await connection`VACUUM (ANALYZE) outbox_event`
441+
const gateNow = new Date()
442+
expect(await hasDueOutboxWork(gateNow)).toBe(false)
443+
444+
const plans = await db.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(sql`
445+
EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${dueOutboxWorkQuery(gateNow)}
446+
`)
447+
const plan = plans[0]['QUERY PLAN'][0].Plan
448+
expect(planNodes(plan).some((node) => node['Node Type'] === 'Seq Scan')).toBe(false)
449+
expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(100)
450+
}, 60_000)
451+
})
389452
})

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

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,12 @@ import { describeError, toError } from '@sim/utils/errors'
55
import { generateId } from '@sim/utils/id'
66
import { truncate } from '@sim/utils/string'
77
import { and, asc, desc, eq, inArray, lte, sql } from 'drizzle-orm'
8-
import { readyEventTypesQuery } from '@/lib/core/outbox/queries'
8+
import {
9+
dueOutboxWorkQuery,
10+
isStuckProcessing,
11+
readyEventTypesQuery,
12+
STUCK_PROCESSING_THRESHOLD_MS,
13+
} from '@/lib/core/outbox/queries'
914

1015
const logger = createLogger('OutboxService')
1116

@@ -24,7 +29,6 @@ const MAX_REAPED_EVENTS = 1_000
2429
function toPersistedHandlerError(error: unknown): string {
2530
return truncate(toError(error).message.split(/\nparams: /)[0], MAX_PERSISTED_ERROR_LENGTH)
2631
}
27-
const STUCK_PROCESSING_THRESHOLD_MS = 10 * 60 * 1000 // 10 minutes
2832
const MAX_BACKOFF_MS = 60 * 60 * 1000 // 1 hour
2933
const BASE_BACKOFF_MS = 1000 // 1 second, doubled per attempt
3034
/** Ordinary handlers keep a short window; longer handlers explicitly opt in below the stale-lease limit. */
@@ -542,18 +546,26 @@ export async function processOutboxEventById(
542546
return runHandler(event, handlers)
543547
}
544548

549+
/**
550+
* Whether {@link processOutboxEvents} would find anything at `now`: a pending event the claim
551+
* phase may take, or a stale lease the reaper may reclaim.
552+
*/
553+
export async function hasDueOutboxWork(now: Date): Promise<boolean> {
554+
const [row] = await db.execute<{ due: boolean }>(dueOutboxWorkQuery(now))
555+
return row?.due === true
556+
}
557+
545558
/**
546559
* Reaper: move `processing` rows whose worker died (stale `lockedAt`)
547560
* back to `pending` so another worker can pick them up. Without this,
548561
* a SIGKILL between claim and result-write would permanently strand
549562
* the row in `processing`.
550563
*/
551564
async function reapStuckProcessingRows(): Promise<number> {
552-
const stuckBefore = new Date(Date.now() - STUCK_PROCESSING_THRESHOLD_MS)
553565
const stuckRows = db
554566
.select({ id: outboxEvent.id })
555567
.from(outboxEvent)
556-
.where(and(eq(outboxEvent.status, 'processing'), lte(outboxEvent.lockedAt, stuckBefore)))
568+
.where(isStuckProcessing(new Date()))
557569
.orderBy(asc(outboxEvent.lockedAt), asc(outboxEvent.id))
558570
.limit(MAX_REAPED_EVENTS)
559571
.for('update', { skipLocked: true })

‎packages/testing/src/mocks/outbox-service.mock.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ export const outboxServiceMockFns = {
6969
}),
7070
mockFindDeadLetteredEvents: vi.fn(),
7171
mockHasInflightOutboxEvent: vi.fn(),
72+
mockHasDueOutboxWork: vi.fn(),
7273
mockProcessOutboxEvents: vi.fn(),
7374
mockProcessOutboxEventById: vi.fn(),
7475
}
@@ -97,6 +98,7 @@ export const outboxServiceMock = {
9798
outboxPayloadHasSourceOperationId: outboxServiceMockFns.mockOutboxPayloadHasSourceOperationId,
9899
findDeadLetteredEvents: outboxServiceMockFns.mockFindDeadLetteredEvents,
99100
hasInflightOutboxEvent: outboxServiceMockFns.mockHasInflightOutboxEvent,
101+
hasDueOutboxWork: outboxServiceMockFns.mockHasDueOutboxWork,
100102
processOutboxEvents: outboxServiceMockFns.mockProcessOutboxEvents,
101103
processOutboxEventById: outboxServiceMockFns.mockProcessOutboxEventById,
102104
}

0 commit comments

Comments
 (0)