Skip to content

Commit dda4edc

Browse files
priosshrsthclaude
andcommitted
OUT-4105 | Stop dropping buffered grouped emails when the DB is briefly unreachable
The flush job's retry budget was ~3 seconds (3 attempts, 1s/2s backoff). When the Supabase pooler went unreachable for 47 minutes on 8/24 every window exhausted it on the first $queryRaw, and onFailure then DELETEd the whole window — unsent rows included — so those grouped emails were destroyed, or orphaned as sentAt IS NULL when the delete failed too. bufferGroupedEmailEvent only reuses windows younger than 5 minutes, so orphans are never re-flushed. Widen the in-run backoff to ~75s, and on exhaustion re-enqueue the window with a 15m delay for up to two rounds before cleaning up. The run is already idempotent (it reads sentAt IS NULL and marks each recipient as it goes), so a re-enqueue only sends what is still outstanding. Sentry now fires once the rounds are spent rather than on every window of a transient outage. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y4tF84tzJW7Bo1DehsDX22
1 parent 9806fe8 commit dda4edc

2 files changed

Lines changed: 47 additions & 12 deletions

File tree

‎src/jobs/notifications/flush-grouped-email.test.ts‎

Lines changed: 23 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -318,25 +318,41 @@ describe('flushGroupedEmailRun', () => {
318318
})
319319

320320
describe('enqueueGroupedEmailFlush', () => {
321-
it('triggers with a 5-minute delay and a workspace-scoped idempotency key', () => {
321+
it('triggers with a 5-minute delay and a round-scoped idempotency key', () => {
322322
enqueueGroupedEmailFlush(payload)
323323

324324
expect(flushGroupedEmail.trigger).toHaveBeenCalledTimes(1)
325325
const [triggeredPayload, opts] = (flushGroupedEmail.trigger as jest.Mock).mock.calls[0]
326326
expect(triggeredPayload).toEqual(payload)
327327
expect(opts).toMatchObject({
328328
delay: '5m',
329-
idempotencyKey: 'ws_1:client_1:win_1',
329+
idempotencyKey: 'ws_1:client_1:win_1:0',
330330
idempotencyKeyTTL: '10m',
331331
})
332332
})
333+
334+
it('waits longer and uses a distinct idempotency key on a retry round', () => {
335+
enqueueGroupedEmailFlush({ ...payload, retryRound: 1 })
336+
337+
const [, opts] = (flushGroupedEmail.trigger as jest.Mock).mock.calls[0]
338+
expect(opts).toMatchObject({ delay: '15m', idempotencyKey: 'ws_1:client_1:win_1:1' })
339+
})
333340
})
334341

335342
describe('flushGroupedEmailOnFailure', () => {
336-
it('captures to Sentry with workspace and window tags', async () => {
343+
it('re-enqueues the window instead of dropping buffered events', async () => {
344+
await flushGroupedEmailOnFailure({ payload, error: new Error('db unreachable') })
345+
346+
expect(mockExecuteRaw).not.toHaveBeenCalled()
347+
expect(mockCaptureException).not.toHaveBeenCalled()
348+
const [triggeredPayload] = (flushGroupedEmail.trigger as jest.Mock).mock.calls[0]
349+
expect(triggeredPayload).toEqual({ ...payload, retryRound: 1 })
350+
})
351+
352+
it('captures to Sentry with workspace and window tags once rounds are exhausted', async () => {
337353
const error = new Error('terminal')
338354

339-
await flushGroupedEmailOnFailure({ payload, error })
355+
await flushGroupedEmailOnFailure({ payload: { ...payload, retryRound: 2 }, error })
340356

341357
expect(mockCaptureException).toHaveBeenCalledTimes(1)
342358
const [captured, opts] = mockCaptureException.mock.calls[0]
@@ -348,9 +364,10 @@ describe('flushGroupedEmailOnFailure', () => {
348364
})
349365
})
350366

351-
it('deletes all rows for the window so orphaned events do not accumulate', async () => {
352-
await flushGroupedEmailOnFailure({ payload, error: new Error('terminal') })
367+
it('deletes all rows for the window only after the last retry round', async () => {
368+
await flushGroupedEmailOnFailure({ payload: { ...payload, retryRound: 2 }, error: new Error('terminal') })
353369

354370
expect(mockExecuteRaw).toHaveBeenCalledTimes(1)
371+
expect(flushGroupedEmail.trigger).not.toHaveBeenCalled()
355372
})
356373
})

‎src/jobs/notifications/flush-grouped-email.ts‎

Lines changed: 24 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import { sendGroupedEmail } from './send-grouped-email'
1818
export type FlushGroupedEmailPayload = {
1919
workspaceId: string
2020
windowKey: string
21+
retryRound?: number
2122
}
2223

2324
type WindowEvent = GroupedEmailEventInput & { individualEmail: NotificationRequestBody | null }
@@ -40,6 +41,10 @@ type IuRecipientGroup = {
4041
}
4142

4243
const TASK_ID = 'flush-grouped-email'
44+
const FLUSH_DELAY = '5m'
45+
// A pooler blip outlives the in-run backoff, so a dead window gets re-enqueued instead of dropped.
46+
const RETRY_ROUND_DELAY = '15m'
47+
const MAX_RETRY_ROUNDS = 2
4348

4449
const readUnsentWindowEvents = (db: ReturnType<typeof DBClient.getInstance>, windowKey: string) =>
4550
db.$queryRaw<BufferedRow[]>`
@@ -319,11 +324,24 @@ export const flushGroupedEmailRun = async (payload: FlushGroupedEmailPayload) =>
319324
}
320325

321326
export const flushGroupedEmailOnFailure = async ({ payload, error }: { payload: unknown; error: unknown }) => {
322-
const { workspaceId, windowKey } = payload as FlushGroupedEmailPayload
323-
Sentry.captureException(error, { tags: { job: TASK_ID, workspaceId, windowKey } })
324-
logger.error('flush-grouped-email: retries exhausted, cleaning up window', {
327+
const { workspaceId, windowKey, retryRound = 0 } = payload as FlushGroupedEmailPayload
328+
329+
if (retryRound < MAX_RETRY_ROUNDS) {
330+
logger.warn('flush-grouped-email: attempts exhausted, re-enqueueing window', {
331+
workspaceId,
332+
windowKey,
333+
retryRound,
334+
error: serializeError(error),
335+
})
336+
await enqueueGroupedEmailFlush({ workspaceId, windowKey, retryRound: retryRound + 1 })
337+
return
338+
}
339+
340+
Sentry.captureException(error, { tags: { job: TASK_ID, workspaceId, windowKey, retryRound: String(retryRound) } })
341+
logger.error('flush-grouped-email: retry rounds exhausted, cleaning up window', {
325342
workspaceId,
326343
windowKey,
344+
retryRound,
327345
error: serializeError(error),
328346
})
329347
const db = DBClient.getInstance()
@@ -341,7 +359,7 @@ export const flushGroupedEmailOnFailure = async ({ payload, error }: { payload:
341359
export const flushGroupedEmail = task({
342360
id: TASK_ID,
343361
queue: { concurrencyLimit: 5 },
344-
retry: { maxAttempts: 3, factor: 2, minTimeoutInMs: 1_000, maxTimeoutInMs: 15_000, randomize: true },
362+
retry: { maxAttempts: 5, factor: 2, minTimeoutInMs: 5_000, maxTimeoutInMs: 60_000, randomize: true },
345363
maxDuration: 60,
346364
run: flushGroupedEmailRun,
347365
})
@@ -350,7 +368,7 @@ tasks.onFailure(TASK_ID, flushGroupedEmailOnFailure)
350368

351369
export const enqueueGroupedEmailFlush = (payload: FlushGroupedEmailPayload) =>
352370
flushGroupedEmail.trigger(payload, {
353-
delay: '5m',
354-
idempotencyKey: `${payload.workspaceId}:${payload.windowKey}`,
371+
delay: payload.retryRound ? RETRY_ROUND_DELAY : FLUSH_DELAY,
372+
idempotencyKey: `${payload.workspaceId}:${payload.windowKey}:${payload.retryRound ?? 0}`,
355373
idempotencyKeyTTL: '10m',
356374
})

0 commit comments

Comments
 (0)