Skip to content

Commit 4e522cd

Browse files
Bill LeoutsakosBill Leoutsakos
authored andcommitted
fix(newsletters): handle review recovery cases
1 parent daa8db0 commit 4e522cd

8 files changed

Lines changed: 535 additions & 100 deletions

File tree

apps/sim/components/settings/account-settings-renderer.tsx

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,5 +49,6 @@ export function AccountSettingsRenderer({ section }: AccountSettingsRendererProp
4949
if (section === 'api-keys') return <ApiKeys scope='personal' />
5050
if (section === 'admin') return <Admin />
5151
if (section === 'mothership') return <Mothership />
52-
return <Newsletters />
52+
if (section === 'newsletters') return <Newsletters />
53+
return null
5354
}

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

Lines changed: 163 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ const mocks = vi.hoisted(() => ({
1515
isAsyncJobEnqueueError: vi.fn(),
1616
markFailed: vi.fn(),
1717
markPushed: vi.fn(),
18+
queueCancel: vi.fn(),
1819
queueEnqueue: vi.fn(),
1920
queueGetJob: vi.fn(),
2021
requireAttempt: vi.fn(),
@@ -31,6 +32,8 @@ vi.mock('@/lib/core/async-jobs', () => ({
3132
JOB_STATUS: {
3233
COMPLETED: 'completed',
3334
FAILED: 'failed',
35+
PENDING: 'pending',
36+
PROCESSING: 'processing',
3437
},
3538
}))
3639

@@ -70,6 +73,7 @@ describe('newsletter Resend queueing', () => {
7073
vi.clearAllMocks()
7174
mocks.getAsyncBackendType.mockReturnValue('trigger-dev')
7275
mocks.getJobQueue.mockResolvedValue({
76+
cancelJob: mocks.queueCancel,
7377
enqueue: mocks.queueEnqueue,
7478
getJob: mocks.queueGetJob,
7579
})
@@ -93,6 +97,7 @@ describe('newsletter Resend queueing', () => {
9397
expect.objectContaining({ jobId: 'newsletter_resend_run-1_2', maxAttempts: 3 })
9498
)
9599
expect(mocks.setJob).toHaveBeenCalledWith('run-1', 2, 'trigger-run-123')
100+
expect(mocks.resetFailedRecipients).not.toHaveBeenCalled()
96101
expect(result.jobId).toBe('trigger-run-123')
97102
})
98103

@@ -142,23 +147,151 @@ describe('newsletter Resend queueing', () => {
142147
)
143148
})
144149

145-
it('moves a newsletter run to failed when its persisted database job failed', async () => {
146-
mocks.getAsyncBackendType.mockReturnValue('database')
147-
mocks.claimAttempt.mockResolvedValue({
148-
attempt: 2,
149-
jobId: 'newsletter_resend_run-1_2',
150-
run,
151-
shouldEnqueue: false,
152-
})
150+
it('starts a new attempt when a persisted Trigger.dev job failed', async () => {
151+
mocks.claimAttempt
152+
.mockResolvedValueOnce({
153+
attempt: 2,
154+
jobId: 'trigger-run-failed',
155+
run,
156+
shouldEnqueue: false,
157+
})
158+
.mockResolvedValueOnce({
159+
attempt: 3,
160+
jobId: null,
161+
run: { ...run, resendSyncJobId: null },
162+
shouldEnqueue: true,
163+
})
153164
mocks.queueGetJob.mockResolvedValue({
154-
id: 'newsletter_resend_run-1_2',
165+
id: 'trigger-run-failed',
155166
status: 'failed',
156167
error: 'worker stopped',
157168
})
169+
mocks.queueEnqueue.mockResolvedValue('trigger-run-retry')
170+
mocks.setJob.mockResolvedValue({ ...run, resendSyncJobId: 'trigger-run-retry' })
171+
172+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
158173

159-
await expect(enqueueNewsletterResendSync('run-1', 'admin-1')).rejects.toThrow('worker stopped')
160174
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
161-
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
175+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
176+
'newsletter-resend-sync',
177+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
178+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
179+
)
180+
expect(result.jobId).toBe('trigger-run-retry')
181+
})
182+
183+
it('starts a new attempt when a completed job did not finalize the newsletter run', async () => {
184+
mocks.claimAttempt
185+
.mockResolvedValueOnce({
186+
attempt: 2,
187+
jobId: 'trigger-run-completed',
188+
run,
189+
shouldEnqueue: false,
190+
})
191+
.mockResolvedValueOnce({
192+
attempt: 3,
193+
jobId: null,
194+
run: { ...run, resendSyncJobId: null },
195+
shouldEnqueue: true,
196+
})
197+
mocks.queueGetJob.mockResolvedValue({
198+
id: 'trigger-run-completed',
199+
status: 'completed',
200+
})
201+
mocks.queueEnqueue.mockResolvedValue('trigger-run-retry')
202+
mocks.setJob.mockResolvedValue({ ...run, resendSyncJobId: 'trigger-run-retry' })
203+
204+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
205+
206+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
207+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
208+
'newsletter-resend-sync',
209+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
210+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
211+
)
212+
expect(result.jobId).toBe('trigger-run-retry')
213+
})
214+
215+
it.each([
216+
['completed', { id: 'trigger-run-existing', status: 'completed' }],
217+
['processing', { id: 'trigger-run-existing', status: 'processing' }],
218+
['missing', null],
219+
])(
220+
'keeps the stored job when the newsletter run is pushed and the provider job is %s',
221+
async (_providerState, providerJob) => {
222+
const pushedRun = { ...run, status: 'pushed' }
223+
mocks.claimAttempt.mockResolvedValue({
224+
attempt: 2,
225+
jobId: 'trigger-run-existing',
226+
run: pushedRun,
227+
shouldEnqueue: false,
228+
})
229+
mocks.queueGetJob.mockResolvedValue(providerJob)
230+
231+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
232+
233+
expect(mocks.queueGetJob).not.toHaveBeenCalled()
234+
expect(mocks.markFailed).not.toHaveBeenCalled()
235+
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
236+
expect(result).toEqual({ run: pushedRun, jobId: 'trigger-run-existing' })
237+
}
238+
)
239+
240+
it('re-enqueues when a stored Trigger.dev run no longer exists', async () => {
241+
mocks.claimAttempt
242+
.mockResolvedValueOnce({
243+
attempt: 2,
244+
jobId: 'trigger-run-missing',
245+
run,
246+
shouldEnqueue: false,
247+
})
248+
.mockResolvedValueOnce({
249+
attempt: 3,
250+
jobId: null,
251+
run,
252+
shouldEnqueue: true,
253+
})
254+
mocks.queueGetJob.mockResolvedValue(null)
255+
256+
await enqueueNewsletterResendSync('run-1', 'admin-1')
257+
258+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
259+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
260+
'newsletter-resend-sync',
261+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
262+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
263+
)
264+
})
265+
266+
it('cancels and replaces an active Trigger.dev run when an admin resumes it', async () => {
267+
mocks.claimAttempt
268+
.mockResolvedValueOnce({
269+
attempt: 2,
270+
jobId: 'trigger-run-active',
271+
run,
272+
shouldEnqueue: false,
273+
})
274+
.mockResolvedValueOnce({
275+
attempt: 3,
276+
jobId: null,
277+
run,
278+
shouldEnqueue: true,
279+
})
280+
mocks.queueGetJob.mockResolvedValue({
281+
id: 'trigger-run-active',
282+
status: 'processing',
283+
})
284+
285+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
286+
287+
expect(mocks.queueCancel).toHaveBeenCalledWith('trigger-run-active')
288+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
289+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
290+
'newsletter-resend-sync',
291+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
292+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
293+
)
294+
expect(result.jobId).toBe('trigger-run-123')
162295
})
163296

164297
it('resets failed recipients before a task retry', async () => {
@@ -173,10 +306,28 @@ describe('newsletter Resend queueing', () => {
173306
requestedById: 'admin-1',
174307
})
175308

176-
expect(mocks.resetFailedRecipients).toHaveBeenCalledWith('run-1')
309+
expect(mocks.resetFailedRecipients).toHaveBeenCalledWith('run-1', 2)
177310
expect(mocks.markPushed).toHaveBeenCalledWith('run-1', 2, 'segment-1', 'Segment 1')
178311
})
179312

313+
it('does not mark the run pushed while recipients remain pending', async () => {
314+
mocks.requireAttempt.mockResolvedValue(run)
315+
mocks.getExcludedEmails.mockResolvedValue(new Set())
316+
mocks.getPendingRecipients.mockResolvedValue([])
317+
mocks.countByStatus.mockResolvedValue({ pending: 1 })
318+
319+
await expect(
320+
runNewsletterResendSync({
321+
runId: 'run-1',
322+
attempt: 2,
323+
requestedById: 'admin-1',
324+
})
325+
).rejects.toThrow('1 remain pending')
326+
327+
expect(mocks.markPushed).not.toHaveBeenCalled()
328+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
329+
})
330+
180331
it('treats same-attempt worker re-entry after success as a no-op', async () => {
181332
mocks.requireAttempt.mockResolvedValue({ ...run, status: 'pushed' })
182333

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

Lines changed: 42 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -67,11 +67,13 @@ export async function runNewsletterResendSync(
6767
try {
6868
signal?.throwIfAborted()
6969
const run = await requireNewsletterRunAttempt(runId, attempt)
70+
signal?.throwIfAborted()
7071
if (run.status === 'pushed') return
7172
if (run.status !== 'finalized' && run.status !== 'pushing' && run.status !== 'failed') {
7273
throw new Error('Newsletter run must be finalized before pushing to Resend')
7374
}
74-
await resetFailedNewsletterRecipients(runId)
75+
await resetFailedNewsletterRecipients(runId, attempt)
76+
await requireNewsletterRunAttempt(runId, attempt)
7577

7678
signal?.throwIfAborted()
7779
const segment =
@@ -92,7 +94,11 @@ export async function runNewsletterResendSync(
9294

9395
while (true) {
9496
signal?.throwIfAborted()
95-
const recipients = await getPendingNewsletterRecipients(runId, NEWSLETTER_RESEND_BATCH_SIZE)
97+
const recipients = await getPendingNewsletterRecipients(
98+
runId,
99+
attempt,
100+
NEWSLETTER_RESEND_BATCH_SIZE
101+
)
96102
if (recipients.length === 0) break
97103

98104
await runWithConcurrency(
@@ -107,6 +113,7 @@ export async function runNewsletterResendSync(
107113
) {
108114
await updateRecipientSyncStatus(
109115
runId,
116+
attempt,
110117
recipient.snapshotVersion,
111118
recipient.email,
112119
'excluded'
@@ -125,6 +132,7 @@ export async function runNewsletterResendSync(
125132
signal?.throwIfAborted()
126133
await updateRecipientSyncStatus(
127134
runId,
135+
attempt,
128136
recipient.snapshotVersion,
129137
recipient.email,
130138
result.status,
@@ -136,6 +144,7 @@ export async function runNewsletterResendSync(
136144
signal?.throwIfAborted()
137145
await updateRecipientSyncStatus(
138146
runId,
147+
attempt,
139148
recipient.snapshotVersion,
140149
recipient.email,
141150
'failed',
@@ -157,11 +166,14 @@ export async function runNewsletterResendSync(
157166
}
158167

159168
signal?.throwIfAborted()
160-
const statusCounts = await countNewsletterRecipientsByStatus(runId)
169+
const statusCounts = await countNewsletterRecipientsByStatus(runId, attempt)
161170
signal?.throwIfAborted()
162171
const failed = statusCounts.failed ?? 0
163-
if (failed > 0) {
164-
throw new Error(`${failed} recipients failed to sync to Resend`)
172+
const pending = statusCounts.pending ?? 0
173+
if (failed > 0 || pending > 0) {
174+
throw new Error(
175+
`Newsletter sync incomplete: ${failed} recipients failed and ${pending} remain pending`
176+
)
165177
}
166178

167179
signal?.throwIfAborted()
@@ -174,28 +186,43 @@ export async function runNewsletterResendSync(
174186
}
175187

176188
export async function enqueueNewsletterResendSync(runId: string, requestedById: string) {
177-
const claim = await claimNewsletterRunResendAttempt(runId)
189+
let claim = await claimNewsletterRunResendAttempt(runId)
178190
const queue = await getJobQueue()
191+
const backendType = getAsyncBackendType()
179192
if (!claim.shouldEnqueue && claim.jobId) {
180-
if (getAsyncBackendType() !== 'database') {
193+
if (claim.run.status === 'pushed') {
181194
return { run: claim.run, jobId: claim.jobId }
182195
}
183-
184196
const persistedJob = await queue.getJob(claim.jobId)
185197
if (persistedJob?.status === JOB_STATUS.COMPLETED) {
186-
return { run: claim.run, jobId: claim.jobId }
187-
}
188-
if (persistedJob?.status === JOB_STATUS.FAILED) {
189-
const error = new Error(persistedJob.error ?? 'Newsletter sync database job failed')
198+
const error = new Error('Newsletter sync job completed without finalizing the newsletter run')
199+
await markNewsletterRunPushFailed(runId, claim.attempt, error)
200+
claim = await claimNewsletterRunResendAttempt(runId)
201+
} else if (backendType !== 'database') {
202+
if (
203+
persistedJob?.status === JOB_STATUS.PENDING ||
204+
persistedJob?.status === JOB_STATUS.PROCESSING
205+
) {
206+
await queue.cancelJob(claim.jobId)
207+
}
208+
const error = new Error(
209+
persistedJob?.error ?? 'Newsletter sync was resumed with a fresh background job'
210+
)
211+
await markNewsletterRunPushFailed(runId, claim.attempt, error)
212+
claim = await claimNewsletterRunResendAttempt(runId)
213+
} else if (persistedJob?.status === JOB_STATUS.FAILED) {
214+
const error = new Error(persistedJob.error ?? 'Newsletter sync job failed')
190215
await markNewsletterRunPushFailed(runId, claim.attempt, error)
191-
throw error
216+
claim = await claimNewsletterRunResendAttempt(runId)
192217
}
193218
}
194219

195-
const enqueueKey = claim.jobId ?? `newsletter_resend_${runId}_${claim.attempt}`
220+
const enqueueKey =
221+
backendType === 'database' && claim.jobId
222+
? claim.jobId
223+
: `newsletter_resend_${runId}_${claim.attempt}`
196224
let jobId: string
197225
try {
198-
await resetFailedNewsletterRecipients(runId)
199226
jobId = await queue.enqueue<NewsletterResendSyncPayload>(
200227
'newsletter-resend-sync',
201228
{ runId, attempt: claim.attempt, requestedById },

0 commit comments

Comments
 (0)