Skip to content

Commit 849323e

Browse files
Bill LeoutsakosBill Leoutsakos
authored andcommitted
fix(newsletters): handle review recovery cases
1 parent 19613d7 commit 849323e

5 files changed

Lines changed: 235 additions & 37 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: 143 additions & 11 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
})
@@ -142,23 +146,151 @@ describe('newsletter Resend queueing', () => {
142146
)
143147
})
144148

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-
})
149+
it('starts a new attempt when a persisted Trigger.dev job failed', async () => {
150+
mocks.claimAttempt
151+
.mockResolvedValueOnce({
152+
attempt: 2,
153+
jobId: 'trigger-run-failed',
154+
run,
155+
shouldEnqueue: false,
156+
})
157+
.mockResolvedValueOnce({
158+
attempt: 3,
159+
jobId: null,
160+
run: { ...run, resendSyncJobId: null },
161+
shouldEnqueue: true,
162+
})
153163
mocks.queueGetJob.mockResolvedValue({
154-
id: 'newsletter_resend_run-1_2',
164+
id: 'trigger-run-failed',
155165
status: 'failed',
156166
error: 'worker stopped',
157167
})
168+
mocks.queueEnqueue.mockResolvedValue('trigger-run-retry')
169+
mocks.setJob.mockResolvedValue({ ...run, resendSyncJobId: 'trigger-run-retry' })
170+
171+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
172+
173+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
174+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
175+
'newsletter-resend-sync',
176+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
177+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
178+
)
179+
expect(result.jobId).toBe('trigger-run-retry')
180+
})
181+
182+
it('starts a new attempt when a completed job did not finalize the newsletter run', async () => {
183+
mocks.claimAttempt
184+
.mockResolvedValueOnce({
185+
attempt: 2,
186+
jobId: 'trigger-run-completed',
187+
run,
188+
shouldEnqueue: false,
189+
})
190+
.mockResolvedValueOnce({
191+
attempt: 3,
192+
jobId: null,
193+
run: { ...run, resendSyncJobId: null },
194+
shouldEnqueue: true,
195+
})
196+
mocks.queueGetJob.mockResolvedValue({
197+
id: 'trigger-run-completed',
198+
status: 'completed',
199+
})
200+
mocks.queueEnqueue.mockResolvedValue('trigger-run-retry')
201+
mocks.setJob.mockResolvedValue({ ...run, resendSyncJobId: 'trigger-run-retry' })
202+
203+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
204+
205+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
206+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
207+
'newsletter-resend-sync',
208+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
209+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
210+
)
211+
expect(result.jobId).toBe('trigger-run-retry')
212+
})
213+
214+
it.each([
215+
['completed', { id: 'trigger-run-existing', status: 'completed' }],
216+
['processing', { id: 'trigger-run-existing', status: 'processing' }],
217+
['missing', null],
218+
])(
219+
'keeps the stored job when the newsletter run is pushed and the provider job is %s',
220+
async (_providerState, providerJob) => {
221+
const pushedRun = { ...run, status: 'pushed' }
222+
mocks.claimAttempt.mockResolvedValue({
223+
attempt: 2,
224+
jobId: 'trigger-run-existing',
225+
run: pushedRun,
226+
shouldEnqueue: false,
227+
})
228+
mocks.queueGetJob.mockResolvedValue(providerJob)
229+
230+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
231+
232+
expect(mocks.queueGetJob).not.toHaveBeenCalled()
233+
expect(mocks.markFailed).not.toHaveBeenCalled()
234+
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
235+
expect(result).toEqual({ run: pushedRun, jobId: 'trigger-run-existing' })
236+
}
237+
)
238+
239+
it('re-enqueues when a stored Trigger.dev run no longer exists', async () => {
240+
mocks.claimAttempt
241+
.mockResolvedValueOnce({
242+
attempt: 2,
243+
jobId: 'trigger-run-missing',
244+
run,
245+
shouldEnqueue: false,
246+
})
247+
.mockResolvedValueOnce({
248+
attempt: 3,
249+
jobId: null,
250+
run,
251+
shouldEnqueue: true,
252+
})
253+
mocks.queueGetJob.mockResolvedValue(null)
254+
255+
await enqueueNewsletterResendSync('run-1', 'admin-1')
256+
257+
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
258+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
259+
'newsletter-resend-sync',
260+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
261+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
262+
)
263+
})
264+
265+
it('cancels and replaces an active Trigger.dev run when an admin resumes it', async () => {
266+
mocks.claimAttempt
267+
.mockResolvedValueOnce({
268+
attempt: 2,
269+
jobId: 'trigger-run-active',
270+
run,
271+
shouldEnqueue: false,
272+
})
273+
.mockResolvedValueOnce({
274+
attempt: 3,
275+
jobId: null,
276+
run,
277+
shouldEnqueue: true,
278+
})
279+
mocks.queueGetJob.mockResolvedValue({
280+
id: 'trigger-run-active',
281+
status: 'processing',
282+
})
158283

159-
await expect(enqueueNewsletterResendSync('run-1', 'admin-1')).rejects.toThrow('worker stopped')
284+
const result = await enqueueNewsletterResendSync('run-1', 'admin-1')
285+
286+
expect(mocks.queueCancel).toHaveBeenCalledWith('trigger-run-active')
160287
expect(mocks.markFailed).toHaveBeenCalledWith('run-1', 2, expect.any(Error))
161-
expect(mocks.queueEnqueue).not.toHaveBeenCalled()
288+
expect(mocks.queueEnqueue).toHaveBeenCalledWith(
289+
'newsletter-resend-sync',
290+
{ runId: 'run-1', attempt: 3, requestedById: 'admin-1' },
291+
expect.objectContaining({ jobId: 'newsletter_resend_run-1_3' })
292+
)
293+
expect(result.jobId).toBe('trigger-run-123')
162294
})
163295

164296
it('resets failed recipients before a task retry', async () => {

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

Lines changed: 25 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -174,25 +174,41 @@ export async function runNewsletterResendSync(
174174
}
175175

176176
export async function enqueueNewsletterResendSync(runId: string, requestedById: string) {
177-
const claim = await claimNewsletterRunResendAttempt(runId)
177+
let claim = await claimNewsletterRunResendAttempt(runId)
178178
const queue = await getJobQueue()
179+
const backendType = getAsyncBackendType()
179180
if (!claim.shouldEnqueue && claim.jobId) {
180-
if (getAsyncBackendType() !== 'database') {
181+
if (claim.run.status === 'pushed') {
181182
return { run: claim.run, jobId: claim.jobId }
182183
}
183-
184184
const persistedJob = await queue.getJob(claim.jobId)
185185
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')
186+
const error = new Error('Newsletter sync job completed without finalizing the newsletter run')
187+
await markNewsletterRunPushFailed(runId, claim.attempt, error)
188+
claim = await claimNewsletterRunResendAttempt(runId)
189+
} else if (backendType !== 'database') {
190+
if (
191+
persistedJob?.status === JOB_STATUS.PENDING ||
192+
persistedJob?.status === JOB_STATUS.PROCESSING
193+
) {
194+
await queue.cancelJob(claim.jobId)
195+
}
196+
const error = new Error(
197+
persistedJob?.error ?? 'Newsletter sync was resumed with a fresh background job'
198+
)
199+
await markNewsletterRunPushFailed(runId, claim.attempt, error)
200+
claim = await claimNewsletterRunResendAttempt(runId)
201+
} else if (persistedJob?.status === JOB_STATUS.FAILED) {
202+
const error = new Error(persistedJob.error ?? 'Newsletter sync job failed')
190203
await markNewsletterRunPushFailed(runId, claim.attempt, error)
191-
throw error
204+
claim = await claimNewsletterRunResendAttempt(runId)
192205
}
193206
}
194207

195-
const enqueueKey = claim.jobId ?? `newsletter_resend_${runId}_${claim.attempt}`
208+
const enqueueKey =
209+
backendType === 'database' && claim.jobId
210+
? claim.jobId
211+
: `newsletter_resend_${runId}_${claim.attempt}`
196212
let jobId: string
197213
try {
198214
await resetFailedNewsletterRecipients(runId)

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

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -127,21 +127,53 @@ describe('newsletter Resend service', () => {
127127
it('normalizes suppressed email addresses', async () => {
128128
fetchMock.mockResolvedValueOnce(
129129
jsonResponse({
130-
data: [{ email: ' First@Example.com ' }, { email: 'second@example.com' }],
130+
data: [
131+
{ id: 'suppression-1', email: ' First@Example.com ' },
132+
{ id: 'suppression-2', email: 'second@example.com' },
133+
],
131134
has_more: false,
132135
})
133136
)
134137

135138
const emails = await getResendSuppressedEmails()
136139

137140
expect(emails).toEqual(new Set(['first@example.com', 'second@example.com']))
141+
expect(fetchMock).toHaveBeenCalledWith(
142+
'https://api.resend.com/suppressions?limit=100',
143+
expect.objectContaining({ method: 'GET' })
144+
)
145+
})
146+
147+
it('paginates through all suppressed email addresses', async () => {
148+
fetchMock
149+
.mockResolvedValueOnce(
150+
jsonResponse({
151+
data: [{ id: 'suppression-1', email: 'first@example.com' }],
152+
has_more: true,
153+
})
154+
)
155+
.mockResolvedValueOnce(
156+
jsonResponse({
157+
data: [{ id: 'suppression-2', email: 'second@example.com' }],
158+
has_more: false,
159+
})
160+
)
161+
162+
const emails = await getResendSuppressedEmails()
163+
164+
expect(emails).toEqual(new Set(['first@example.com', 'second@example.com']))
165+
expect(fetchMock).toHaveBeenNthCalledWith(
166+
2,
167+
'https://api.resend.com/suppressions?limit=100&after=suppression-1',
168+
expect.objectContaining({ method: 'GET' })
169+
)
138170
})
139171

140172
it('combines suppressions with globally unsubscribed contacts', async () => {
141173
fetchMock
142174
.mockResolvedValueOnce(
143175
jsonResponse({
144-
data: [{ email: 'suppressed@example.com' }],
176+
data: [{ id: 'suppression-1', email: 'suppressed@example.com' }],
145177
has_more: false,
146178
})
147179
)

apps/sim/lib/newsletters/resend.ts

Lines changed: 31 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@ const RESEND_CONTACT_PROPERTY_PAGE_LIMIT = 100
1313
const RESEND_CONTACT_PROPERTY_MAX_PAGES = 10
1414
const RESEND_CONTACT_PAGE_LIMIT = 100
1515
const RESEND_CONTACT_MAX_PAGES = 1000
16+
const RESEND_SUPPRESSION_PAGE_LIMIT = 100
17+
const RESEND_SUPPRESSION_MAX_PAGES = 1000
1618
const NEWSLETTER_CONTACT_PROPERTY_KEYS = ['sim_user_id', 'newsletter_run_id'] as const
1719

1820
interface ResendErrorBody {
@@ -43,7 +45,7 @@ interface ResendRequestOptions {
4345
}
4446

4547
const resendSuppressionListSchema = z.object({
46-
data: z.array(z.object({ email: z.string().min(1) })),
48+
data: z.array(z.object({ id: z.string().min(1), email: z.string().min(1) })),
4749
has_more: z.boolean(),
4850
})
4951

@@ -116,22 +118,37 @@ async function resendRequest<T>(path: string, options: ResendRequestOptions = {}
116118
}
117119

118120
export async function getResendSuppressedEmails(options?: { required?: boolean }) {
119-
const rawResponse = await resendRequest<unknown>('/suppressions', {
120-
required: options?.required ?? true,
121-
})
122-
if (rawResponse === undefined) return new Set<string>()
123-
const parsedResponse = resendSuppressionListSchema.safeParse(rawResponse)
124-
if (!parsedResponse.success) {
125-
throw new Error('Resend suppression list response was malformed')
126-
}
127-
const response = parsedResponse.data
128121
const emails = new Set<string>()
129-
if (response?.has_more) {
130-
throw new Error('Resend suppression list was incomplete')
122+
let after: string | null = null
123+
let hasMore = false
124+
125+
for (let page = 0; page < RESEND_SUPPRESSION_MAX_PAGES; page++) {
126+
const query = new URLSearchParams({ limit: String(RESEND_SUPPRESSION_PAGE_LIMIT) })
127+
if (after) query.set('after', after)
128+
const rawResponse = await resendRequest<unknown>(`/suppressions?${query.toString()}`, {
129+
required: options?.required ?? true,
130+
})
131+
if (rawResponse === undefined) return emails
132+
const parsedResponse = resendSuppressionListSchema.safeParse(rawResponse)
133+
if (!parsedResponse.success) {
134+
throw new Error('Resend suppression list response was malformed')
135+
}
136+
const response = parsedResponse.data
137+
138+
for (const suppression of response.data) {
139+
emails.add(normalizeEmail(suppression.email))
140+
}
141+
142+
hasMore = response.has_more
143+
if (!hasMore) break
144+
after = response.data.at(-1)?.id ?? null
145+
if (!after) throw new Error('Resend suppression pagination returned no cursor')
131146
}
132-
for (const suppression of response.data) {
133-
emails.add(normalizeEmail(suppression.email))
147+
148+
if (hasMore) {
149+
throw new Error('Resend suppression list exceeded the pagination safety limit')
134150
}
151+
135152
return emails
136153
}
137154

0 commit comments

Comments
 (0)