Skip to content
4 changes: 2 additions & 2 deletions apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ export const GET = withRouteHandler(async (request: NextRequest) => {

try {
const result = await enqueueWorkspaceFileSearchDispatch()
logger.info('Workspace file search dispatcher accepted', result)
return NextResponse.json({ success: true, triggered: true, ...result }, { status: 202 })
if (result.triggered) logger.info('Workspace file search dispatcher accepted', result)
return NextResponse.json({ success: true, ...result }, { status: result.triggered ? 202 : 200 })
} catch (error) {
logger.error('Workspace file search dispatcher enqueue failed', {
error: toError(error).message,
Expand Down
8 changes: 4 additions & 4 deletions apps/sim/app/api/webhooks/outbox/process/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,12 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
const requestId = generateRequestId()
try {
const accepted = await enqueueOutboxProcessor()
if (!accepted.triggered) {
return NextResponse.json({ success: true, requestId, triggered: false })
}
if (accepted.backend === 'trigger-dev') {
logger.info('Outbox processor accepted', { jobId: accepted.jobId })
return NextResponse.json(
{ success: true, requestId, triggered: true, ...accepted },
{ status: 202 }
)
return NextResponse.json({ success: true, requestId, ...accepted }, { status: 202 })
}
return NextResponse.json({ success: true, requestId, ...accepted.output })
} catch (error) {
Expand Down
73 changes: 73 additions & 0 deletions apps/sim/lib/core/async-jobs/scheduled-pass.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import type { TriggerOptions } from '@trigger.dev/sdk'

/**
* What one cron tick did with a scheduled pass: started none, handed one to Trigger.dev, or
* started one in this process.
*/
export type ScheduledPassResult =
| { triggered: false; backend: null; jobId: null }
| { triggered: true; backend: 'trigger-dev'; jobId: string }
| { triggered: true; backend: 'inline'; jobId: null }

/** The Trigger.dev run a scheduled pass starts, deduplicated to one per schedule window. */
interface ScheduledPassTrigger {
taskId: string
/** The idempotency key is `${keyPrefix}:${window}`; defaults to the task id. */
keyPrefix?: string
/** The schedule window's length: every tick inside one window starts the same run. */
intervalMs: number
/** The instant whose window keys the run; defaults to the moment the run is triggered. */
at?: Date
options?: Pick<TriggerOptions, 'maxDuration' | 'ttl'>
}

interface ScheduledPass {
/** Whether this tick owes a pass. The call site owns that policy, including a failed check. */
due: boolean
/**
* Whether Trigger.dev runs the pass. Read only once the pass is due, so an availability check
* with side effects of its own never runs on a tick that starts nothing.
*/
triggerAvailable: () => boolean
/** Starts the pass in this process without waiting for it. */
startInline: () => void
trigger: ScheduledPassTrigger
}

/**
* Triggers the window's run and resolves its id once Trigger.dev durably accepts it. Duplicate
* ticks in one window share the idempotency key, and the key outlives its window so a late
* duplicate is still folded into the run its window started.
*/
export async function triggerScheduledPass(trigger: ScheduledPassTrigger): Promise<string> {
const [{ tasks }, { resolveTriggerRegion }] = await Promise.all([
import('@trigger.dev/sdk'),
import('@/lib/core/async-jobs/region'),
])
const window = Math.floor((trigger.at?.getTime() ?? Date.now()) / trigger.intervalMs)
const handle = await tasks.trigger(trigger.taskId, undefined, {
idempotencyKey: `${trigger.keyPrefix ?? trigger.taskId}:${window}`,
idempotencyKeyTTL: '5m',
...trigger.options,
region: await resolveTriggerRegion(),
})
return handle.id
}

/**
* Runs one cron tick of a scheduled pass: nothing when no pass is due, the window's Trigger.dev
* run when Trigger.dev is available, and otherwise a pass started in this process. The result
* returns once the pass is accepted, never after it finishes.
*/
export async function startScheduledPass(pass: ScheduledPass): Promise<ScheduledPassResult> {
if (!pass.due) return { triggered: false, backend: null, jobId: null }
if (!pass.triggerAvailable()) {
pass.startInline()
return { triggered: true, backend: 'inline', jobId: null }
}
return {
triggered: true,
backend: 'trigger-dev',
jobId: await triggerScheduledPass(pass.trigger),
}
}
2 changes: 2 additions & 0 deletions apps/sim/lib/core/outbox/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,5 @@ export const OUTBOX_PROCESSOR_INTERVAL_MS = 60_000
export const OUTBOX_PROCESSOR_CONCURRENCY = Math.ceil(
(OUTBOX_PROCESSOR_MAX_DURATION_SECONDS * 1000) / OUTBOX_PROCESSOR_INTERVAL_MS
)
/** Every Nth scheduled tick runs the processor's recovery, reaping and pruning even with nothing due. */
export const OUTBOX_MAINTENANCE_EVERY_INTERVALS = 5
106 changes: 93 additions & 13 deletions apps/sim/lib/core/outbox/enqueue.test.ts
Original file line number Diff line number Diff line change
@@ -1,35 +1,78 @@
import { asyncJobsRegionMock } from '@sim/testing/mocks/async-jobs-region.mock'
import { setEnvFlags } from '@sim/testing/mocks/env-flags.mock'
import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock'
import {
createIdempotentTasksTrigger,
triggerSdkMockFns,
} from '@sim/testing/mocks/trigger-sdk.mock'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'

const hoisted = vi.hoisted(() => ({
processor: vi.fn(),
}))
vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock)
vi.mock('@/lib/core/outbox/processor', () => ({ runOutboxProcessor: hoisted.processor }))
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)

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

const mocks = { ...hoisted, trigger: vi.mocked(tasks.trigger) }
const mocks = {
...hoisted,
trigger: triggerSdkMockFns.mockTasksTrigger,
hasDueWork: outboxServiceMockFns.mockHasDueOutboxWork,
}

/** 12:34 is not a maintenance minute (754 % 5 = 4); 12:35 is. */
const IDLE_MINUTE = new Date('2026-09-16T12:34:45Z')
const MAINTENANCE_MINUTE = new Date('2026-09-16T12:35:10Z')

const INLINE_OUTPUT = {
result: { processed: 0, retried: 0, deadLettered: 0, leaseLost: 0, reaped: 0 },
recoveredDocuments: 0,
reapedBackgroundWork: 0,
}

/** Each backend's own proof that it, and only it, ran the processor. */
const BACKENDS = [
{
name: 'Trigger.dev',
isTriggerDevEnabled: true,
started: { triggered: true, backend: 'trigger-dev', jobId: 'run-1' },
},
{
name: 'self-hosted inline',
isTriggerDevEnabled: false,
started: { triggered: true, backend: 'inline', jobId: null, output: INLINE_OUTPUT },
},
] as const

describe('outbox processor enqueue', () => {
beforeEach(() => {
vi.useFakeTimers()
vi.setSystemTime(new Date('2026-09-16T12:34:45Z'))
vi.setSystemTime(IDLE_MINUTE)
setEnvFlags({ isTriggerDevEnabled: true })
mocks.trigger.mockResolvedValue({ id: 'run-1' })
mocks.trigger.mockImplementation(createIdempotentTasksTrigger())
mocks.processor.mockReset()
mocks.processor.mockResolvedValue(INLINE_OUTPUT)
mocks.hasDueWork.mockReset()
mocks.hasDueWork.mockResolvedValue(true)
})
afterEach(() => vi.useRealTimers())

it('deduplicates duplicate ticks while allowing the next minute to drain more work', async () => {
await enqueueOutboxProcessor()
await enqueueOutboxProcessor()
vi.advanceTimersByTime(60_000)
await enqueueOutboxProcessor()
const keys = mocks.trigger.mock.calls.map((call) => call[2].idempotencyKey)
expect(keys[0]).toBe(keys[1])
expect(keys[2]).not.toBe(keys[0])
it('folds ticks into the run of the minute their work check was made in', async () => {
const first = await enqueueOutboxProcessor()
const duplicate = await enqueueOutboxProcessor()
mocks.hasDueWork.mockImplementationOnce(async () => {
vi.advanceTimersByTime(60_000)
return true
})
const checkedAcrossTheBoundary = await enqueueOutboxProcessor()
const nextMinute = await enqueueOutboxProcessor()

expect(first.jobId).toBe('run-1')
expect(duplicate.jobId).toBe('run-1')
expect(checkedAcrossTheBoundary.jobId).toBe('run-1')
expect(nextMinute.jobId).toBe('run-2')
})

it('fails closed on an enqueue error without starting concurrent inline work', async () => {
Expand All @@ -46,7 +89,44 @@ describe('outbox processor enqueue', () => {
reapedBackgroundWork: 1,
}
mocks.processor.mockResolvedValueOnce(output)
await expect(enqueueOutboxProcessor()).resolves.toEqual({ backend: 'inline', output })
await expect(enqueueOutboxProcessor()).resolves.toEqual({
triggered: true,
backend: 'inline',
jobId: null,
output,
})
expect(mocks.trigger).not.toHaveBeenCalled()
})

it.each(BACKENDS)(
'starts no $name processor on an idle queue outside the maintenance window',
async ({ isTriggerDevEnabled }) => {
setEnvFlags({ isTriggerDevEnabled })
mocks.hasDueWork.mockResolvedValue(false)
await expect(enqueueOutboxProcessor()).resolves.toEqual({
triggered: false,
backend: null,
jobId: null,
})
}
)

it.each(BACKENDS)(
'starts the $name processor on an idle queue in the maintenance window',
async ({ isTriggerDevEnabled, started }) => {
setEnvFlags({ isTriggerDevEnabled })
vi.setSystemTime(MAINTENANCE_MINUTE)
mocks.hasDueWork.mockResolvedValue(false)
await expect(enqueueOutboxProcessor()).resolves.toEqual(started)
}
)

it.each(BACKENDS)(
'starts the $name processor when the work check fails',
async ({ isTriggerDevEnabled, started }) => {
setEnvFlags({ isTriggerDevEnabled })
mocks.hasDueWork.mockRejectedValue(new Error('connection refused'))
await expect(enqueueOutboxProcessor()).resolves.toEqual(started)
}
)
})
65 changes: 49 additions & 16 deletions apps/sim/lib/core/outbox/enqueue.ts
Original file line number Diff line number Diff line change
@@ -1,32 +1,65 @@
import { createLogger } from '@sim/logger'
import { getErrorMessage } from '@sim/utils/errors'
import {
type ScheduledPassResult,
triggerScheduledPass,
} from '@/lib/core/async-jobs/scheduled-pass'
import { isTriggerDevEnabled } from '@/lib/core/config/env-flags'
import {
OUTBOX_MAINTENANCE_EVERY_INTERVALS,
OUTBOX_PROCESSOR_INTERVAL_MS,
OUTBOX_PROCESSOR_MAX_DURATION_SECONDS,
} from '@/lib/core/outbox/constants'
import type { OutboxProcessorResult } from '@/lib/core/outbox/processor'
import type { processOutboxTask } from '@/background/process-outbox'
import { hasDueOutboxWork } from '@/lib/core/outbox/service'

const logger = createLogger('OutboxProcessorEnqueue')

/** A self-hosted inline run finishes before the cron request returns, so it carries its output. */
type OutboxProcessorEnqueueResult =
| { backend: 'trigger-dev'; jobId: string }
| { backend: 'inline'; output: OutboxProcessorResult }
| Exclude<ScheduledPassResult, { backend: 'inline' }>
| { triggered: true; backend: 'inline'; jobId: null; output: OutboxProcessorResult }

/** The database owns delivery state; the cron request waits only for durable worker acceptance. */
/**
* Every tick in the maintenance window runs, so document recovery, background-work reaping and
* pruning keep their cadence on an idle queue. Other ticks run only when events are due or a lease
* is stale; a gate that cannot answer runs anyway, since billing, seat sync, invitations and
* document dispatch all ride the outbox.
*/
async function shouldRunOutboxProcessor(now: Date, scheduleWindow: number): Promise<boolean> {
if (scheduleWindow % OUTBOX_MAINTENANCE_EVERY_INTERVALS === 0) return true
try {
return await hasDueOutboxWork(now)
} catch (error) {
logger.warn('Outbox work check failed; running the processor anyway', {
error: getErrorMessage(error),
})
return true
}
}

/**
* The database owns delivery state; the cron request waits only for durable worker acceptance.
* The self-hosted branch processes synchronously and returns the output the route reports, which
* is why it stays here instead of in the scheduled pass's detached inline start.
*/
export async function enqueueOutboxProcessor(): Promise<OutboxProcessorEnqueueResult> {
const now = new Date()
const scheduleWindow = Math.floor(now.getTime() / OUTBOX_PROCESSOR_INTERVAL_MS)
if (!(await shouldRunOutboxProcessor(now, scheduleWindow))) {
return { triggered: false, backend: null, jobId: null }
}

if (!isTriggerDevEnabled) {
const { runOutboxProcessor } = await import('@/lib/core/outbox/processor')
return { backend: 'inline', output: await runOutboxProcessor() }
return { triggered: true, backend: 'inline', jobId: null, output: await runOutboxProcessor() }
}

const [{ tasks }, { resolveTriggerRegion }] = await Promise.all([
import('@trigger.dev/sdk'),
import('@/lib/core/async-jobs/region'),
])
const scheduleWindow = Math.floor(Date.now() / OUTBOX_PROCESSOR_INTERVAL_MS)
const handle = await tasks.trigger<typeof processOutboxTask>('process-outbox', undefined, {
idempotencyKey: `process-outbox:${scheduleWindow}`,
idempotencyKeyTTL: '5m',
maxDuration: OUTBOX_PROCESSOR_MAX_DURATION_SECONDS,
region: await resolveTriggerRegion(),
const jobId = await triggerScheduledPass({
taskId: 'process-outbox',
intervalMs: OUTBOX_PROCESSOR_INTERVAL_MS,
at: now,
options: { maxDuration: OUTBOX_PROCESSOR_MAX_DURATION_SECONDS },
})
return { backend: 'trigger-dev', jobId: handle.id }
return { triggered: true, backend: 'trigger-dev', jobId }
}
28 changes: 27 additions & 1 deletion apps/sim/lib/core/outbox/queries.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,18 @@
import { db } from '@sim/db'
import { outboxEvent } from '@sim/db/schema'
import { sql } from 'drizzle-orm'
import { and, eq, exists, lte, or, type SQL, sql } from 'drizzle-orm'

const MAX_READY_EVENT_TYPES = 128
/** How long a `processing` lease may go without a terminal write before the reaper reclaims it. */
export const STUCK_PROCESSING_THRESHOLD_MS = 10 * 60 * 1000

/** A `processing` row whose lease the reaper may reclaim at `now`. */
export function isStuckProcessing(now: Date): SQL | undefined {
return and(
eq(outboxEvent.status, 'processing'),
lte(outboxEvent.lockedAt, new Date(now.getTime() - STUCK_PROCESSING_THRESHOLD_MS))
)
}

/**
* Walks the pending index one type at a time, reading only its earliest availability.
Expand Down Expand Up @@ -42,3 +53,18 @@ export function readyEventTypesQuery(now: Date) {
LIMIT ${MAX_READY_EVENT_TYPES}
`
}

/**
* One row, `due`: whether `processOutboxEvents` would act at `now`. The pending leg is the
* claim phase's own discovery walk, so it reads only each type's head and cannot be planned as a
* scan of future or completed rows; the lease leg is the reaper's predicate, over the few rows in
* `processing`.
*/
export function dueOutboxWorkQuery(now: Date) {
const staleLease = db
.select({ id: outboxEvent.id })
.from(outboxEvent)
.where(isStuckProcessing(now))
.limit(1)
return sql`SELECT ${or(sql`EXISTS (${readyEventTypesQuery(now)})`, exists(staleLease))} AS due`
}
5 changes: 2 additions & 3 deletions apps/sim/lib/core/outbox/retention.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,8 @@ 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).
* Most rows deleted per type per run. Recovery enqueues only inside a run, at most 200 per run, so
* a 1,000-row prune keeps five times its pace however often the processor runs.
*/
export const OUTBOX_PRUNE_BATCH_SIZE = 1_000

Expand Down
Loading
Loading