Skip to content

Commit 493375b

Browse files
Bill LeoutsakosBill Leoutsakos
authored andcommitted
fix(newsletters): fence resend recovery state
1 parent 07ff1c5 commit 493375b

10 files changed

Lines changed: 499 additions & 66 deletions

File tree

apps/sim/app/api/superuser/newsletters/runs/[id]/export.csv/route.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { exportNewsletterRunCsvContract } from '@/lib/api/contracts/newsletters'
55
import { parseRequest } from '@/lib/api/server'
66
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
77
import { validateNewsletterSuperuser } from '@/lib/newsletters/auth'
8+
import { isNewsletterResendError } from '@/lib/newsletters/resend'
89
import { createNewsletterCsvExport } from '@/lib/newsletters/runs'
910

1011
const logger = createLogger('NewsletterCsvExportAPI')
@@ -58,7 +59,7 @@ export const GET = withRouteHandler(
5859
if (/Finalize/i.test(message)) {
5960
return NextResponse.json({ error: message }, { status: 400 })
6061
}
61-
if (/RESEND_API_KEY|Resend .*list|Resend request/i.test(message)) {
62+
if (isNewsletterResendError(error)) {
6263
return NextResponse.json({ error: message }, { status: 503 })
6364
}
6465
logger.error('Failed to export newsletter CSV', { error: message })
Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,63 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { createMockRequest } from '@sim/testing'
5+
import { beforeEach, describe, expect, it, vi } from 'vitest'
6+
7+
const { mockFinalizeNewsletterRun, mockValidateNewsletterSuperuser } = vi.hoisted(() => ({
8+
mockFinalizeNewsletterRun: vi.fn(),
9+
mockValidateNewsletterSuperuser: vi.fn(),
10+
}))
11+
12+
vi.mock('@/lib/newsletters/auth', () => ({
13+
validateNewsletterSuperuser: mockValidateNewsletterSuperuser,
14+
}))
15+
16+
vi.mock('@/lib/newsletters/runs', () => ({
17+
finalizeNewsletterRun: mockFinalizeNewsletterRun,
18+
}))
19+
20+
import { NewsletterResendError } from '@/lib/newsletters/resend'
21+
import { POST } from '@/app/api/superuser/newsletters/runs/[id]/finalize/route'
22+
23+
function callRoute() {
24+
const request = createMockRequest(
25+
'POST',
26+
undefined,
27+
{},
28+
'http://localhost:3000/api/superuser/newsletters/runs/run-1/finalize'
29+
)
30+
return POST(request, { params: Promise.resolve({ id: 'run-1' }) })
31+
}
32+
33+
describe('newsletter run finalization', () => {
34+
beforeEach(() => {
35+
vi.clearAllMocks()
36+
mockValidateNewsletterSuperuser.mockResolvedValue({
37+
success: true,
38+
userId: 'admin-1',
39+
})
40+
})
41+
42+
it.each([
43+
'Resend suppression pagination returned no cursor',
44+
'Resend contact pagination returned no cursor',
45+
'Resend contact property pagination returned no cursor',
46+
])('maps a Resend service failure to 503: %s', async (message) => {
47+
mockFinalizeNewsletterRun.mockRejectedValueOnce(new NewsletterResendError(message))
48+
49+
const response = await callRoute()
50+
51+
expect(response.status).toBe(503)
52+
await expect(response.json()).resolves.toEqual({ error: message })
53+
})
54+
55+
it('does not classify an unrelated error by message text', async () => {
56+
mockFinalizeNewsletterRun.mockRejectedValueOnce(new Error('Resend text from unrelated code'))
57+
58+
const response = await callRoute()
59+
60+
expect(response.status).toBe(500)
61+
await expect(response.json()).resolves.toEqual({ error: 'Internal server error' })
62+
})
63+
})

apps/sim/app/api/superuser/newsletters/runs/[id]/finalize/route.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { finalizeNewsletterRunContract } from '@/lib/api/contracts/newsletters'
55
import { parseRequest } from '@/lib/api/server'
66
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
77
import { validateNewsletterSuperuser } from '@/lib/newsletters/auth'
8+
import { isNewsletterResendError } from '@/lib/newsletters/resend'
89
import { finalizeNewsletterRun } from '@/lib/newsletters/runs'
910

1011
const logger = createLogger('NewsletterFinalizeAPI')
@@ -35,7 +36,7 @@ export const POST = withRouteHandler(
3536
if (/already in progress/i.test(message)) {
3637
return NextResponse.json({ error: message }, { status: 409 })
3738
}
38-
if (/RESEND_API_KEY|Resend .*list|Resend request/i.test(message)) {
39+
if (isNewsletterResendError(error)) {
3940
return NextResponse.json({ error: message }, { status: 503 })
4041
}
4142
logger.error('Failed to finalize newsletter run', { error: message })

apps/sim/lib/api/contracts/newsletters.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -114,7 +114,7 @@ export const newsletterRunResponseSchema = z.object({
114114

115115
export const pushNewsletterRunResponseSchema = z.object({
116116
run: newsletterRunSchema,
117-
jobId: z.string(),
117+
jobId: z.string().nullable(),
118118
})
119119

120120
export const newsletterJobResponseSchema = z.object({

apps/sim/lib/newsletters/push-resend.test.ts

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@ vi.mock('@/lib/newsletters/runs', () => ({
5858
updateRecipientSyncStatus: mocks.updateRecipient,
5959
}))
6060

61+
import { pushNewsletterRunResponseSchema } from '@/lib/api/contracts/newsletters'
6162
import { enqueueNewsletterResendSync, runNewsletterResendSync } from '@/lib/newsletters/push-resend'
6263

6364
const run = {
@@ -86,6 +87,7 @@ describe('newsletter Resend queueing', () => {
8687
})
8788
mocks.queueEnqueue.mockResolvedValue('trigger-run-123')
8889
mocks.setJob.mockResolvedValue({ ...run, resendSyncJobId: 'trigger-run-123' })
90+
mocks.setSegment.mockResolvedValue(run)
8991
mocks.isAsyncJobEnqueueError.mockReturnValue(false)
9092
})
9193

@@ -238,6 +240,55 @@ describe('newsletter Resend queueing', () => {
238240
}
239241
)
240242

243+
it('does not enqueue a pushed run without a stored job id', async () => {
244+
const pushedRun = { ...run, status: 'pushed' }
245+
mocks.claimAttempt.mockResolvedValue({
246+
attempt: 2,
247+
jobId: null,
248+
run: pushedRun,
249+
shouldEnqueue: false,
250+
})
251+
252+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
253+
254+
expect(mocks.getJobQueue).not.toHaveBeenCalled()
255+
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
256+
expect(result).toEqual({ run: pushedRun, jobId: null })
257+
expect(pushNewsletterRunResponseSchema.shape.jobId.parse(result.jobId)).toBeNull()
258+
})
259+
260+
it.each([
261+
['completed', { id: 'trigger-run-existing', status: 'completed' }],
262+
['failed', { id: 'trigger-run-existing', status: 'failed', error: 'worker failed' }],
263+
['processing', { id: 'trigger-run-existing', status: 'processing' }],
264+
])(
265+
'does not enqueue when reconciliation of a %s job refreshes to pushed',
266+
async (_providerState, providerJob) => {
267+
const pushedRun = { ...run, status: 'pushed' }
268+
mocks.claimAttempt
269+
.mockResolvedValueOnce({
270+
attempt: 2,
271+
jobId: 'trigger-run-existing',
272+
run,
273+
shouldEnqueue: false,
274+
})
275+
.mockResolvedValueOnce({
276+
attempt: 2,
277+
jobId: null,
278+
run: pushedRun,
279+
shouldEnqueue: false,
280+
})
281+
mocks.queueGetJob.mockResolvedValue(providerJob)
282+
283+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
284+
285+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
286+
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
287+
expect(mocks.setJob).not.toHaveBeenCalled()
288+
expect(result).toEqual({ run: pushedRun, jobId: null })
289+
}
290+
)
291+
241292
it('re-enqueues when a stored Trigger.dev run no longer exists', async () => {
242293
mocks.claimAttempt
243294
.mockResolvedValueOnce({
@@ -342,6 +393,78 @@ describe('newsletter Resend queueing', () => {
342393
expect(mocks.markFailed).not.toHaveBeenCalled()
343394
})
344395

396+
it('stops when the authoritative post-reset read is already pushed', async () => {
397+
mocks.requireAttempt
398+
.mockResolvedValueOnce(run)
399+
.mockResolvedValueOnce({ ...run, status: 'pushed' })
400+
401+
await runNewsletterResendSync({
402+
runId: 'run-1',
403+
attempt: 2,
404+
requestedById: 'admin-1',
405+
})
406+
407+
expect(mocks.createSegment).not.toHaveBeenCalled()
408+
expect(mocks.setSegment).not.toHaveBeenCalled()
409+
expect(mocks.ensureProperties).not.toHaveBeenCalled()
410+
expect(mocks.markFailed).not.toHaveBeenCalled()
411+
})
412+
413+
it('uses the segment from the authoritative post-reset read', async () => {
414+
const staleRun = {
415+
...run,
416+
resendSegmentId: null,
417+
resendSegmentName: null,
418+
}
419+
const currentRun = {
420+
...run,
421+
resendSegmentId: 'segment-current',
422+
resendSegmentName: 'Current segment',
423+
}
424+
mocks.requireAttempt.mockResolvedValueOnce(staleRun).mockResolvedValueOnce(currentRun)
425+
mocks.setSegment.mockResolvedValue(currentRun)
426+
mocks.getExcludedEmails.mockResolvedValue(new Set())
427+
mocks.getPendingRecipients.mockResolvedValue([])
428+
mocks.countByStatus.mockResolvedValue({})
429+
430+
await runNewsletterResendSync({
431+
runId: 'run-1',
432+
attempt: 2,
433+
requestedById: 'admin-1',
434+
})
435+
436+
expect(mocks.createSegment).not.toHaveBeenCalled()
437+
expect(mocks.setSegment).toHaveBeenCalledWith('run-1', 2, 'segment-current', 'Current segment')
438+
expect(mocks.markPushed).toHaveBeenCalledWith('run-1', 2, 'segment-current', 'Current segment')
439+
})
440+
441+
it('stops when another worker pushes while a segment is being created', async () => {
442+
const runWithoutSegment = {
443+
...run,
444+
resendSegmentId: null,
445+
resendSegmentName: null,
446+
}
447+
mocks.requireAttempt.mockResolvedValue(runWithoutSegment)
448+
mocks.createSegment.mockResolvedValue({ id: 'segment-new', name: 'New segment' })
449+
mocks.setSegment.mockResolvedValue({
450+
...run,
451+
status: 'pushed',
452+
resendSegmentId: 'segment-winner',
453+
resendSegmentName: 'Winner segment',
454+
})
455+
456+
await runNewsletterResendSync({
457+
runId: 'run-1',
458+
attempt: 2,
459+
requestedById: 'admin-1',
460+
})
461+
462+
expect(mocks.setSegment).toHaveBeenCalledWith('run-1', 2, 'segment-new', 'New segment')
463+
expect(mocks.ensureProperties).not.toHaveBeenCalled()
464+
expect(mocks.markPushed).not.toHaveBeenCalled()
465+
expect(mocks.markFailed).not.toHaveBeenCalled()
466+
})
467+
345468
it('does not fail the newsletter attempt when its database claim is aborted', async () => {
346469
const controller = new AbortController()
347470
controller.abort('claim lost')

apps/sim/lib/newsletters/push-resend.ts

Lines changed: 36 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -74,17 +74,32 @@ export async function runNewsletterResendSync(
7474
throw new Error('Newsletter run must be finalized before pushing to Resend')
7575
}
7676
await resetFailedNewsletterRecipients(runId, attempt)
77-
await requireNewsletterRunAttempt(runId, attempt)
77+
const currentRun = await requireNewsletterRunAttempt(runId, attempt)
78+
if (currentRun.status === 'pushed') return
79+
if (currentRun.status !== 'pushing' && currentRun.status !== 'failed') {
80+
throw new Error('Newsletter run is not eligible to continue its Resend sync')
81+
}
7882

7983
signal?.throwIfAborted()
80-
const segment =
81-
run.resendSegmentId && run.resendSegmentName
82-
? { id: run.resendSegmentId, name: run.resendSegmentName }
83-
: await createNewsletterSegment(segmentNameForRun(run.name), { signal })
84+
const segmentCandidate =
85+
currentRun.resendSegmentId && currentRun.resendSegmentName
86+
? { id: currentRun.resendSegmentId, name: currentRun.resendSegmentName }
87+
: await createNewsletterSegment(segmentNameForRun(currentRun.name), { signal })
8488

8589
signal?.throwIfAborted()
86-
if (!run.resendSegmentId) {
87-
await setNewsletterRunResendSegment(runId, attempt, segment.id, segment.name)
90+
const segmentRun = await setNewsletterRunResendSegment(
91+
runId,
92+
attempt,
93+
segmentCandidate.id,
94+
segmentCandidate.name
95+
)
96+
if (segmentRun.status === 'pushed') return
97+
if (!segmentRun.resendSegmentId || !segmentRun.resendSegmentName) {
98+
throw new Error('Newsletter Resend segment tracking is incomplete')
99+
}
100+
const segment = {
101+
id: segmentRun.resendSegmentId,
102+
name: segmentRun.resendSegmentName,
88103
}
89104

90105
signal?.throwIfAborted()
@@ -189,17 +204,21 @@ export async function runNewsletterResendSync(
189204

190205
export async function enqueueNewsletterResendSync(runId: string, requestedById: string) {
191206
let claim = await claimNewsletterRunResendAttempt(runId)
207+
if (claim.run.status === 'pushed') {
208+
return { run: claim.run, jobId: claim.jobId }
209+
}
210+
192211
const queue = await getJobQueue()
193212
const backendType = getAsyncBackendType()
194213
if (!claim.shouldEnqueue && claim.jobId) {
195-
if (claim.run.status === 'pushed') {
196-
return { run: claim.run, jobId: claim.jobId }
197-
}
198214
const persistedJob = await queue.getJob(claim.jobId)
199215
if (persistedJob?.status === JOB_STATUS.COMPLETED) {
200216
const error = new Error('Newsletter sync job completed without finalizing the newsletter run')
201217
await markNewsletterRunPushFailed(runId, claim.attempt, error)
202218
claim = await claimNewsletterRunResendAttempt(runId)
219+
if (claim.run.status === 'pushed') {
220+
return { run: claim.run, jobId: claim.jobId }
221+
}
203222
} else if (backendType !== 'database') {
204223
if (
205224
persistedJob?.status === JOB_STATUS.PENDING ||
@@ -212,10 +231,16 @@ export async function enqueueNewsletterResendSync(runId: string, requestedById:
212231
)
213232
await markNewsletterRunPushFailed(runId, claim.attempt, error)
214233
claim = await claimNewsletterRunResendAttempt(runId)
234+
if (claim.run.status === 'pushed') {
235+
return { run: claim.run, jobId: claim.jobId }
236+
}
215237
} else if (persistedJob?.status === JOB_STATUS.FAILED) {
216238
const error = new Error(persistedJob.error ?? 'Newsletter sync job failed')
217239
await markNewsletterRunPushFailed(runId, claim.attempt, error)
218240
claim = await claimNewsletterRunResendAttempt(runId)
241+
if (claim.run.status === 'pushed') {
242+
return { run: claim.run, jobId: claim.jobId }
243+
}
219244
}
220245
}
221246

@@ -254,5 +279,5 @@ export async function enqueueNewsletterResendSync(runId: string, requestedById:
254279
{ cause: error }
255280
)
256281
}
257-
return { run: updatedRun, jobId }
282+
return { run: updatedRun, jobId: updatedRun.resendSyncJobId ?? jobId }
258283
}

apps/sim/lib/newsletters/resend.test.ts

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ import {
2222
ensureNewsletterContactProperties,
2323
getResendExcludedEmails,
2424
getResendSuppressedEmails,
25+
isNewsletterResendError,
26+
NewsletterResendError,
2527
} from '@/lib/newsletters/resend'
2628

2729
function jsonResponse(body: unknown, status = 200) {
@@ -50,6 +52,11 @@ describe('newsletter Resend service', () => {
5052
)
5153
})
5254

55+
it('classifies only typed Resend service failures', () => {
56+
expect(isNewsletterResendError(new NewsletterResendError('provider unavailable'))).toBe(true)
57+
expect(isNewsletterResendError(new Error('Resend text from unrelated code'))).toBe(false)
58+
})
59+
5360
it('does not make a Resend request when already aborted', async () => {
5461
const controller = new AbortController()
5562
controller.abort('cancelled')
@@ -60,6 +67,21 @@ describe('newsletter Resend service', () => {
6067
expect(fetchMock).not.toHaveBeenCalled()
6168
})
6269

70+
it('preserves cancellation while reading a Resend error response', async () => {
71+
const controller = new AbortController()
72+
const body = new ReadableStream({
73+
pull(streamController) {
74+
controller.abort('cancelled while reading')
75+
streamController.error(new Error('body read failed'))
76+
},
77+
})
78+
fetchMock.mockResolvedValueOnce(new Response(body, { status: 400 }))
79+
80+
await expect(createNewsletterSegment('Segment 1', { signal: controller.signal })).rejects.toBe(
81+
'cancelled while reading'
82+
)
83+
})
84+
6385
it('does not retry after cancellation during backoff', async () => {
6486
const controller = new AbortController()
6587
fetchMock.mockResolvedValueOnce(jsonResponse({ message: 'retry' }, 429))

0 commit comments

Comments
 (0)