Skip to content

Commit dfbfb4b

Browse files
committed
Merge staging into fix/workspace-file-content-revision-precision
2 parents 4ad1df1 + 2112998 commit dfbfb4b

7 files changed

Lines changed: 27811 additions & 60 deletions

File tree

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
import { outboxEvent } from '@sim/db/schema'
2+
import { sql } from 'drizzle-orm'
3+
4+
const MAX_READY_EVENT_TYPES = 128
5+
6+
/**
7+
* Walks the pending index one type at a time, reading only its earliest availability.
8+
* The time filter belongs AFTER the walk: filtering inside each seek would scan
9+
* all future rows of a type with no ready events. A strictly increasing type ends
10+
* the recursion, including unknown types left by rolling deployments.
11+
*
12+
* Sort and cap ready heads after discovery so future types cannot hide ready ones,
13+
* and preserve the scheduler's oldest-first ordering. Database work scales with
14+
* distinct pending types (plus MVCC visibility checks), not their event counts;
15+
* at most 128 metadata rows cross into the worker, with no payloads or fan-out.
16+
*/
17+
export function readyEventTypesQuery(now: Date) {
18+
return sql`
19+
WITH RECURSIVE pending_heads AS (
20+
(
21+
SELECT event_type, available_at
22+
FROM ${outboxEvent}
23+
WHERE status = 'pending'
24+
ORDER BY event_type, available_at
25+
LIMIT 1
26+
)
27+
UNION ALL
28+
SELECT next_head.event_type, next_head.available_at
29+
FROM pending_heads
30+
CROSS JOIN LATERAL (
31+
SELECT event_type, available_at
32+
FROM ${outboxEvent}
33+
WHERE status = 'pending' AND event_type > pending_heads.event_type
34+
ORDER BY event_type, available_at
35+
LIMIT 1
36+
) next_head
37+
)
38+
SELECT event_type AS "eventType"
39+
FROM pending_heads
40+
WHERE available_at <= ${now.toISOString()}::timestamp
41+
ORDER BY available_at, event_type
42+
LIMIT ${MAX_READY_EVENT_TYPES}
43+
`
44+
}

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

Lines changed: 115 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ vi.mock('@sim/db', () => ({
1717
},
1818
}))
1919

20+
import { readyEventTypesQuery } from '@/lib/core/outbox/queries'
2021
import {
2122
type OutboxHandler,
2223
processOutboxEvents,
@@ -26,6 +27,8 @@ import {
2627
interface QueryPlan {
2728
'Node Type': string
2829
'Index Name'?: string
30+
'Shared Hit Blocks': number
31+
'Shared Read Blocks': number
2932
Plans?: QueryPlan[]
3033
}
3134

@@ -90,10 +93,10 @@ describe('outbox scheduling in PostgreSQL', () => {
9093
async function seedBacklog(
9194
eventType: string,
9295
count: number,
93-
status: 'pending' | 'completed' | 'processing'
96+
status: 'pending' | 'completed' | 'processing',
97+
availableAt = new Date(Date.now() - 60_000)
9498
) {
9599
const prefix = generateId()
96-
const availableAt = new Date(Date.now() - 60_000)
97100
const createdAt = new Date(Date.now() - 24 * 60 * 60_000)
98101
const lockedAt = status === 'processing' ? new Date(Date.now() - 11 * 60_000) : null
99102
eventTypes.add(eventType)
@@ -110,6 +113,114 @@ describe('outbox scheduling in PostgreSQL', () => {
110113
}
111114
}
112115

116+
async function expectBoundedDiscovery(now: Date) {
117+
const plans = await db.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(sql`
118+
EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${readyEventTypesQuery(now)}
119+
`)
120+
const plan = plans[0]['QUERY PLAN'][0].Plan
121+
const nodes = planNodes(plan)
122+
expect(nodes.some((node) => node['Node Type'] === 'Recursive Union')).toBe(true)
123+
expect(nodes.some((node) => node['Node Type'] === 'Seq Scan')).toBe(false)
124+
/** Buffer work, unlike wall-clock time, catches a backlog scan even on a warm local database. */
125+
expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(1_000)
126+
}
127+
128+
it('discovers an empty queue without returning a null type', async () => {
129+
expect(await db.execute(readyEventTypesQuery(new Date()))).toEqual([])
130+
})
131+
132+
it('uses earliest availability, inclusive deadlines, and deterministic type ties', async () => {
133+
const now = new Date('2026-01-01T12:00:00.000Z')
134+
const fixtures = [
135+
{ eventType: 'test.outbox.z-first', availableAt: new Date(now.getTime() - 1) },
136+
{ eventType: 'test.outbox.z-first', availableAt: new Date(now.getTime() + 60_000) },
137+
{ eventType: 'test.outbox.b-tie', availableAt: now },
138+
{ eventType: 'test.outbox.a-tie', availableAt: now },
139+
{ eventType: 'test.outbox.future', availableAt: new Date(now.getTime() + 1) },
140+
{ eventType: 'test.outbox.completed', availableAt: now, status: 'completed' },
141+
{ eventType: 'test.outbox.processing', availableAt: now, status: 'processing' },
142+
{ eventType: 'test.outbox.dead', availableAt: now, status: 'dead_letter' },
143+
]
144+
for (const fixture of fixtures) eventTypes.add(fixture.eventType)
145+
await db
146+
.insert(outboxEvent)
147+
.values(fixtures.map((row) => ({ id: generateId(), payload: {}, ...row })))
148+
149+
expect(await db.execute(readyEventTypesQuery(now))).toEqual([
150+
{ eventType: 'test.outbox.z-first' },
151+
{ eventType: 'test.outbox.a-tie' },
152+
{ eventType: 'test.outbox.b-tie' },
153+
])
154+
})
155+
156+
it('caps ready types after ordering all heads, including types unknown to this worker', async () => {
157+
const now = new Date()
158+
const rows = Array.from({ length: 140 }, (_, index) => ({
159+
id: generateId(),
160+
eventType: `test.outbox.type-${String(index).padStart(3, '0')}`,
161+
payload: {},
162+
availableAt: new Date(now.getTime() - index - 1),
163+
}))
164+
for (const row of rows) eventTypes.add(row.eventType)
165+
await db.insert(outboxEvent).values(rows)
166+
167+
expect(await db.execute(readyEventTypesQuery(now))).toEqual(
168+
[...rows]
169+
.reverse()
170+
.slice(0, 128)
171+
.map(({ eventType }) => ({ eventType }))
172+
)
173+
})
174+
175+
it('skips large future backlogs and more than 128 future types without hiding ready work', async () => {
176+
const future = new Date(Date.now() + 48 * 60 * 60_000)
177+
await seedBacklog('test.outbox.a-expiry', 100_000, 'pending', future)
178+
const futureTypes = Array.from({ length: 130 }, (_, index) => ({
179+
id: generateId(),
180+
eventType: `test.outbox.future-${index}`,
181+
payload: {},
182+
availableAt: future,
183+
}))
184+
for (const row of futureTypes) eventTypes.add(row.eventType)
185+
await db.insert(outboxEvent).values(futureTypes)
186+
await enqueue('test.outbox.z-ready', 1)
187+
await connection`VACUUM (ANALYZE) outbox_event`
188+
189+
const now = new Date()
190+
expect(await db.execute(readyEventTypesQuery(now))).toEqual([
191+
{ eventType: 'test.outbox.z-ready' },
192+
])
193+
await expectBoundedDiscovery(now)
194+
expect(await processOutboxEvents({ 'test.outbox.z-ready': async () => {} })).toMatchObject({
195+
processed: 1,
196+
})
197+
expect(await db.execute(readyEventTypesQuery(new Date()))).toEqual([])
198+
}, 60_000)
199+
200+
it('retains bounded retries for a type missing during a rolling deployment', async () => {
201+
const [event] = await enqueue('test.outbox.unknown', 1)
202+
203+
expect(await processOutboxEvents({})).toMatchObject({ retried: 1 })
204+
const [pending] = await db.select().from(outboxEvent).where(eq(outboxEvent.id, event.id))
205+
expect(pending).toMatchObject({ status: 'pending', attempts: 1 })
206+
expect(pending.availableAt.getTime()).toBeGreaterThan(Date.now())
207+
})
208+
209+
it('lets claims skip a locked head without hiding other rows of that type', async () => {
210+
const [locked, available] = await enqueue('test.outbox.locked', 2)
211+
const delivered: string[] = []
212+
await connection.begin(async (transaction) => {
213+
await transaction`SELECT id FROM outbox_event WHERE id = ${locked.id} FOR UPDATE`
214+
const result = await processOutboxEvents({
215+
'test.outbox.locked': async (_payload, context) => {
216+
delivered.push(context.eventId)
217+
},
218+
})
219+
expect(result.processed).toBe(1)
220+
})
221+
expect(delivered).toEqual([available.id])
222+
})
223+
113224
it('serves newer event types before exhausting an older cleanup backlog', async () => {
114225
await enqueue('test.outbox.cleanup', 1_000)
115226
const [dispatch] = await enqueue('test.outbox.dispatch', 1, 2_000)
@@ -191,7 +302,8 @@ describe('outbox scheduling in PostgreSQL', () => {
191302
await seedBacklog('test.outbox.cleanup', 100_000, 'pending')
192303
const [dispatch] = await enqueue('test.outbox.dispatch', 1)
193304
const [billing] = await enqueue('test.outbox.billing', 1)
194-
await db.execute(sql`ANALYZE outbox_event`)
305+
await connection`VACUUM (ANALYZE) outbox_event`
306+
await expectBoundedDiscovery(new Date())
195307

196308
const plans = await db.execute<{ 'QUERY PLAN': { Plan: QueryPlan }[] }>(sql`
197309
EXPLAIN (FORMAT JSON)

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

Lines changed: 73 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
*/
44

55
import { outboxEvent } from '@sim/db/schema'
6+
import { createLogger } from '@sim/logger'
67
import { dbChainMock, dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing'
78
import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
89

@@ -36,6 +37,11 @@ import {
3637
withOutboxHandlerTimeout,
3738
} from '@/lib/core/outbox/service'
3839

40+
const logger =
41+
vi.mocked(createLogger).mock.results[
42+
vi.mocked(createLogger).mock.calls.findIndex(([name]) => name === 'OutboxService')
43+
].value
44+
3945
function makePendingRow(overrides: Partial<OutboxRow> = {}): OutboxRow {
4046
return {
4147
id: 'evt-1',
@@ -69,8 +75,7 @@ function holdLease() {
6975

7076
/** Queue metadata discovery followed by individually claimed rows. */
7177
function queuePendingEvents(rows: OutboxRow[]) {
72-
queueTableRows(
73-
outboxEvent,
78+
dbChainMockFns.execute.mockResolvedValueOnce(
7479
[...new Set(rows.map(({ eventType }) => eventType))].map((eventType) => ({ eventType }))
7580
)
7681
for (const row of rows) queueTableRows(outboxEvent, [row])
@@ -291,6 +296,68 @@ describe('processOutboxEvents — empty / no handler', () => {
291296
})
292297
})
293298

299+
describe('processOutboxEvents — infrastructure diagnostics', () => {
300+
beforeEach(() => {
301+
vi.clearAllMocks()
302+
resetDbChainMock()
303+
})
304+
305+
it('logs the nested database cause and rethrows discovery failures without exposing parameters', async () => {
306+
const cause = Object.assign(new Error('canceling statement due to statement timeout'), {
307+
code: '57014',
308+
})
309+
const error = new Error('Failed query: select event_type\nparams: private-token', { cause })
310+
dbChainMockFns.execute.mockRejectedValueOnce(error)
311+
312+
await expect(processOutboxEvents({})).rejects.toBe(error)
313+
314+
expect(logger.error).toHaveBeenCalledWith(
315+
'Outbox processing failed',
316+
expect.objectContaining({
317+
phase: 'discover',
318+
processed: 0,
319+
error: expect.objectContaining({
320+
code: '57014',
321+
message: 'canceling statement due to statement timeout',
322+
}),
323+
})
324+
)
325+
expect(JSON.stringify(vi.mocked(logger.error).mock.calls)).not.toContain('private-token')
326+
expect(dbChainMockFns.transaction).not.toHaveBeenCalled()
327+
})
328+
329+
it('distinguishes a reaper failure from discovery and stops the poll', async () => {
330+
const error = new Error('connection closed')
331+
dbChainMockFns.returning.mockRejectedValueOnce(error)
332+
333+
await expect(processOutboxEvents({})).rejects.toBe(error)
334+
335+
expect(logger.error).toHaveBeenCalledWith(
336+
'Outbox processing failed',
337+
expect.objectContaining({ phase: 'reap' })
338+
)
339+
expect(dbChainMockFns.execute).not.toHaveBeenCalled()
340+
})
341+
342+
it('reports completed work when a later claim fails without rerunning handlers', async () => {
343+
const handler = vi.fn(async () => {})
344+
const error = new Error('connection closed')
345+
queuePendingEvents([makePendingRow()])
346+
holdLease()
347+
dbChainMockFns.transaction
348+
.mockImplementationOnce(async (callback) => callback(dbChainMock.db))
349+
.mockRejectedValueOnce(error)
350+
351+
await expect(processOutboxEvents({ 'test.event': handler })).rejects.toBe(error)
352+
353+
expect(logger.error).toHaveBeenCalledWith(
354+
'Outbox processing failed',
355+
expect.objectContaining({ phase: 'claim', processed: 1 })
356+
)
357+
expect(handler).toHaveBeenCalledOnce()
358+
})
359+
})
360+
294361
describe('processOutboxEvents — handler success and retry', () => {
295362
beforeEach(() => {
296363
vi.clearAllMocks()
@@ -542,7 +609,10 @@ describe('processOutboxEvents — handler timeout', () => {
542609
550_000
543610
)
544611
const shortHandler = vi.fn(async () => {})
545-
queueTableRows(outboxEvent, [{ eventType: 'test.long' }, { eventType: 'test.short' }])
612+
dbChainMockFns.execute.mockResolvedValueOnce([
613+
{ eventType: 'test.long' },
614+
{ eventType: 'test.short' },
615+
])
546616
queueTableRows(outboxEvent, [makePendingRow({ eventType: 'test.short' })])
547617
holdLease()
548618

0 commit comments

Comments
 (0)