Skip to content

Commit 40b7cbc

Browse files
committed
fix(outbox): fail the run after maintenance when a handler module could not load
1 parent cfa860d commit 40b7cbc

5 files changed

Lines changed: 61 additions & 9 deletions

File tree

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

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,14 @@ const IDLE_MINUTE = new Date('2026-09-16T12:34:45Z')
2727
const MAINTENANCE_MINUTE = new Date('2026-09-16T12:35:10Z')
2828

2929
const INLINE_OUTPUT = {
30-
result: { processed: 0, retried: 0, deadLettered: 0, leaseLost: 0, reaped: 0 },
30+
result: {
31+
processed: 0,
32+
retried: 0,
33+
deadLettered: 0,
34+
leaseLost: 0,
35+
reaped: 0,
36+
unloadedEventTypes: [],
37+
},
3138
recoveredDocuments: 0,
3239
reapedBackgroundWork: 0,
3340
}
@@ -84,7 +91,14 @@ describe('outbox processor enqueue', () => {
8491
it('preserves synchronous processing for self-hosted deployments without Trigger', async () => {
8592
setEnvFlags({ isTriggerDevEnabled: false })
8693
const output = {
87-
result: { processed: 4, retried: 0, deadLettered: 0, leaseLost: 0, reaped: 0 },
94+
result: {
95+
processed: 4,
96+
retried: 0,
97+
deadLettered: 0,
98+
leaseLost: 0,
99+
reaped: 0,
100+
unloadedEventTypes: [],
101+
},
88102
recoveredDocuments: 2,
89103
reapedBackgroundWork: 1,
90104
}

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

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,14 @@ import { runOutboxProcessor } from '@/lib/core/outbox/processor'
2424
const mockProcessOutboxEvents = outboxServiceMockFns.mockProcessOutboxEvents
2525

2626
describe('outbox processor recovery', () => {
27-
const result = { processed: 5, retried: 1, deadLettered: 0, leaseLost: 0, reaped: 0 }
27+
const result = {
28+
processed: 5,
29+
retried: 1,
30+
deadLettered: 0,
31+
leaseLost: 0,
32+
reaped: 0,
33+
unloadedEventTypes: [],
34+
}
2835

2936
beforeEach(() => {
3037
vi.resetAllMocks()
@@ -86,4 +93,18 @@ describe('outbox processor recovery', () => {
8693
expect(mocks.recover).not.toHaveBeenCalled()
8794
expect(mocks.reap).not.toHaveBeenCalled()
8895
})
96+
97+
it('finishes maintenance, then fails the run naming event types whose handler module failed to load', async () => {
98+
mockProcessOutboxEvents.mockResolvedValueOnce({
99+
...result,
100+
unloadedEventTypes: ['test.broken', 'test.broken-too'],
101+
})
102+
103+
await expect(runOutboxProcessor()).rejects.toThrow(
104+
'Outbox handler modules failed to load; left pending: test.broken, test.broken-too'
105+
)
106+
expect(mocks.recover).toHaveBeenCalledOnce()
107+
expect(mocks.reap).toHaveBeenCalledOnce()
108+
expect(mocks.prune).toHaveBeenCalledOnce()
109+
})
89110
})

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,5 +71,11 @@ export async function runOutboxProcessor(): Promise<OutboxProcessorResult> {
7171
prunedEvents,
7272
durationMs: Date.now() - startedAt,
7373
})
74+
/** Fail the run so a broken handler module stays as visible as the crash its static import caused. */
75+
if (result.unloadedEventTypes.length > 0) {
76+
throw new Error(
77+
`Outbox handler modules failed to load; left pending: ${result.unloadedEventTypes.join(', ')}`
78+
)
79+
}
7480
return output
7581
}

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

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -637,6 +637,7 @@ describe('processOutboxEvents — lazy handler groups', () => {
637637
processed: 1,
638638
retried: 0,
639639
deadLettered: 0,
640+
unloadedEventTypes: ['test.flaky'],
640641
})
641642
expect(healthyHandler).toHaveBeenCalledOnce()
642643
expect(claimedEventTypes()).not.toContain('test.flaky')
@@ -646,7 +647,10 @@ describe('processOutboxEvents — lazy handler groups', () => {
646647
queuePendingEvents([makePendingRow({ eventType: 'test.flaky' })])
647648
holdLease()
648649

649-
expect(await processOutboxEvents(groups)).toMatchObject({ processed: 1 })
650+
expect(await processOutboxEvents(groups)).toMatchObject({
651+
processed: 1,
652+
unloadedEventTypes: [],
653+
})
650654
expect(recoveredHandler).toHaveBeenCalledOnce()
651655
})
652656

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

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -166,6 +166,8 @@ export interface ProcessOutboxResult {
166166
deadLettered: number
167167
leaseLost: number
168168
reaped: number
169+
/** Ready event types left pending because their handler module failed to import. */
170+
unloadedEventTypes: string[]
169171
}
170172

171173
export type ProcessSingleOutboxResult =
@@ -457,7 +459,7 @@ export async function processOutboxEvents(
457459
reaped = await reapStuckProcessingRows()
458460
phase = 'discover'
459461
const readyTypes = await db.execute<{ eventType: string }>(readyEventTypesQuery(new Date()))
460-
const { handlers, eligibleTypes } = await resolveOutboxHandlers(
462+
const { handlers, eligibleTypes, unloadedEventTypes } = await resolveOutboxHandlers(
461463
handlerGroups,
462464
readyTypes.map(({ eventType }) => eventType)
463465
)
@@ -497,7 +499,7 @@ export async function processOutboxEvents(
497499
else retried++
498500
}
499501

500-
return { processed, retried, deadLettered, leaseLost, reaped }
502+
return { processed, retried, deadLettered, leaseLost, reaped, unloadedEventTypes }
501503
} catch (error) {
502504
logger.error('Outbox processing failed', {
503505
phase,
@@ -517,13 +519,17 @@ export async function processOutboxEvents(
517519
* Imports the groups that serve any ready event type and returns the types to claim this run. A
518520
* group whose import fails leaves its event types unclaimed: unlike a missing handler, which
519521
* spends an attempt and eventually dead-letters, a failed import says nothing about the events,
520-
* so they stay pending for a later run. Event types outside every group stay eligible and reach
521-
* the missing-handler path.
522+
* so they stay pending for a later run and are reported as `unloadedEventTypes`. Event types
523+
* outside every group stay eligible and reach the missing-handler path.
522524
*/
523525
async function resolveOutboxHandlers(
524526
groups: readonly LazyOutboxHandlerGroup[],
525527
readyEventTypes: string[]
526-
): Promise<{ handlers: OutboxHandlerRegistry; eligibleTypes: string[] }> {
528+
): Promise<{
529+
handlers: OutboxHandlerRegistry
530+
eligibleTypes: string[]
531+
unloadedEventTypes: string[]
532+
}> {
527533
const unavailableEventTypes = new Set<string>()
528534
const ready = new Set(readyEventTypes)
529535
const dueGroups = groups.filter((group) => group.events.some((eventType) => ready.has(eventType)))
@@ -547,6 +553,7 @@ async function resolveOutboxHandlers(
547553
return {
548554
handlers,
549555
eligibleTypes: readyEventTypes.filter((eventType) => !unavailableEventTypes.has(eventType)),
556+
unloadedEventTypes: readyEventTypes.filter((eventType) => unavailableEventTypes.has(eventType)),
550557
}
551558
}
552559

0 commit comments

Comments
 (0)