Skip to content

Commit 94db255

Browse files
committed
fix(webhooks): stop every poller's batch on a deterministic admission refusal
Only RSS stopped at a refused item; Gmail, Outlook, and IMAP advanced their cursors past refused emails, and every poller counted the refusals toward auto-disable. A shared PollAdmissionRefusedError now leaves the idempotency callback, stops the batch, and returns skipped before any cursor update or failure count; items that already ran replay as idempotent no-ops. Source backoff goes back to RSS only, where the rate-limited feed was: the other pollers' fetch helpers do not carry status or Retry-After, so routing their failures through it would back off on a guess. The block-missing 404 also tells Slack not to redeliver.
1 parent 7c08bec commit 94db255

12 files changed

Lines changed: 234 additions & 95 deletions

File tree

‎apps/sim/lib/webhooks/polling/gmail.ts‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -4,13 +4,16 @@ import { pollingIdempotency } from '@/lib/core/idempotency/service'
44
import {
55
getProviderConfig,
66
type PollingProviderHandler,
7+
type PollOutcome,
78
type PollWebhookContext,
89
} from '@/lib/webhooks/polling/types'
910
import {
1011
markWebhookFailed,
1112
markWebhookSuccess,
12-
recordPollSourceFailure,
13+
PollAdmissionRefusedError,
1314
resolveOAuthCredential,
15+
skipAdmissionRefusedPoll,
16+
throwIfAdmissionRefused,
1417
updateWebhookProviderConfig,
1518
} from '@/lib/webhooks/polling/utils'
1619
import { processPolledWebhookEvent } from '@/lib/webhooks/processor'
@@ -64,10 +67,9 @@ export const gmailPollingHandler: PollingProviderHandler = {
6467
provider: 'gmail',
6568
label: 'Gmail',
6669

67-
async pollWebhook(ctx: PollWebhookContext) {
70+
async pollWebhook(ctx: PollWebhookContext): Promise<PollOutcome> {
6871
const { webhookData, workflowData, requestId, logger } = ctx
6972
const webhookId = webhookData.id
70-
const pollStartedAt = Date.now()
7173

7274
try {
7375
const accessToken = await resolveOAuthCredential(webhookData, 'google-email', requestId)
@@ -135,13 +137,11 @@ export const gmailPollingHandler: PollingProviderHandler = {
135137
)
136138
return 'success'
137139
} catch (error) {
138-
await recordPollSourceFailure(
139-
webhookData,
140-
pollStartedAt,
141-
error,
142-
`[${requestId}] Error polling Gmail webhook ${webhookId}`,
143-
logger
144-
)
140+
if (error instanceof PollAdmissionRefusedError) {
141+
return skipAdmissionRefusedPoll(logger, requestId, webhookId)
142+
}
143+
logger.error(`[${requestId}] Error processing Gmail webhook ${webhookId}:`, error)
144+
await markWebhookFailed(webhookId, logger)
145145
return 'failure'
146146
}
147147
},
@@ -546,6 +546,7 @@ async function processEmails(
546546
)
547547

548548
if (!result.success) {
549+
throwIfAdmissionRefused(result)
549550
logger.error(
550551
`[${requestId}] Failed to process webhook for email ${email.id}:`,
551552
result.statusCode,
@@ -567,6 +568,7 @@ async function processEmails(
567568
)
568569
processedCount++
569570
} catch (error) {
571+
if (error instanceof PollAdmissionRefusedError) throw error
570572
const errorMessage = getErrorMessage(error, 'Unknown error')
571573
logger.error(`[${requestId}] Error processing email ${email.id}:`, errorMessage)
572574
failedCount++
Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
import { createLogger } from '@sim/logger'
2+
import { createWorkflowRecord } from '@sim/testing'
3+
import { jsonResponse } from '@sim/testing/helpers/http'
4+
import {
5+
webhooksProcessorMock,
6+
webhooksProcessorMockFns,
7+
} from '@sim/testing/mocks/webhooks-processor.mock'
8+
import { beforeEach, describe, expect, it, vi } from 'vitest'
9+
10+
const { mockUpdateConfig, mockMarkFailed } = vi.hoisted(() => ({
11+
mockUpdateConfig: vi.fn(),
12+
mockMarkFailed: vi.fn(),
13+
}))
14+
15+
vi.mock('@/lib/core/idempotency/service', () => ({
16+
pollingIdempotency: {
17+
executeWithIdempotency: vi.fn(
18+
async (_provider: string, _key: string, execute: () => Promise<unknown>) => execute()
19+
),
20+
},
21+
}))
22+
23+
vi.mock('@/lib/webhooks/processor', () => webhooksProcessorMock)
24+
25+
vi.mock('@/lib/webhooks/polling/utils', async (importOriginal) => ({
26+
...(await importOriginal<typeof import('@/lib/webhooks/polling/utils')>()),
27+
resolveOAuthCredential: vi.fn().mockResolvedValue('access-token'),
28+
markWebhookSuccess: vi.fn(),
29+
markWebhookFailed: mockMarkFailed,
30+
updateWebhookProviderConfig: mockUpdateConfig,
31+
}))
32+
33+
import { ADMISSION_REJECTION_CODE } from '@/lib/core/admission/rejection'
34+
import { googleCalendarPollingHandler } from '@/lib/webhooks/polling/google-calendar'
35+
import type { PollWebhookContext, WebhookRecord } from '@/lib/webhooks/polling/types'
36+
37+
const mockProcessEvent = webhooksProcessorMockFns.mockProcessPolledWebhookEvent
38+
39+
function context(): PollWebhookContext {
40+
const webhookData: WebhookRecord = {
41+
id: 'calendar-webhook',
42+
workflowId: 'calendar-listener',
43+
deploymentVersionId: null,
44+
registrationStatus: null,
45+
registrationGeneration: null,
46+
configFingerprint: null,
47+
preparedAt: null,
48+
blockId: null,
49+
path: 'calendar-listener',
50+
routingKey: null,
51+
provider: 'google-calendar',
52+
providerConfig: {
53+
calendarId: 'primary',
54+
lastCheckedTimestamp: '2026-10-09T11:00:00.000Z',
55+
},
56+
isActive: true,
57+
failedCount: 0,
58+
lastFailedAt: null,
59+
archivedAt: null,
60+
createdAt: new Date('2026-10-01T00:00:00.000Z'),
61+
updatedAt: new Date('2026-10-09T11:00:00.000Z'),
62+
}
63+
return {
64+
webhookData,
65+
workflowData: createWorkflowRecord({
66+
id: 'calendar-listener',
67+
}) as PollWebhookContext['workflowData'],
68+
requestId: 'calendar-request',
69+
logger: createLogger('GoogleCalendarTest'),
70+
}
71+
}
72+
73+
describe('Google Calendar polling when execution admission refuses events', () => {
74+
beforeEach(() => {
75+
const events = ['event-1', 'event-2'].map((id) => ({
76+
id,
77+
status: 'confirmed',
78+
created: '2026-10-09T11:30:00.000Z',
79+
updated: '2026-10-09T11:30:00.000Z',
80+
}))
81+
vi.stubGlobal('fetch', vi.fn().mockResolvedValue(jsonResponse({ items: events }, 200)))
82+
mockProcessEvent.mockResolvedValue({
83+
success: false,
84+
statusCode: 402,
85+
error: 'Usage limit exceeded',
86+
code: ADMISSION_REJECTION_CODE.USAGE_LIMIT_EXCEEDED,
87+
retryable: false,
88+
})
89+
})
90+
91+
it('stops the batch without advancing its cursor or counting a failure', async () => {
92+
expect(await googleCalendarPollingHandler.pollWebhook(context())).toBe('skipped')
93+
94+
expect(mockProcessEvent).toHaveBeenCalledOnce()
95+
const cursorUpdates = mockUpdateConfig.mock.calls.filter(
96+
([, update]) => 'lastCheckedTimestamp' in (update as Record<string, unknown>)
97+
)
98+
expect(cursorUpdates).toEqual([])
99+
expect(mockMarkFailed).not.toHaveBeenCalled()
100+
})
101+
})

‎apps/sim/lib/webhooks/polling/google-calendar.ts‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,16 @@ import { readCanonicalTriggerValue } from '@/lib/webhooks/polling/canonical'
55
import {
66
getProviderConfig,
77
type PollingProviderHandler,
8+
type PollOutcome,
89
type PollWebhookContext,
910
} from '@/lib/webhooks/polling/types'
1011
import {
1112
markWebhookFailed,
1213
markWebhookSuccess,
13-
recordPollSourceFailure,
14+
PollAdmissionRefusedError,
1415
resolveOAuthCredential,
16+
skipAdmissionRefusedPoll,
17+
throwIfAdmissionRefused,
1518
updateWebhookProviderConfig,
1619
} from '@/lib/webhooks/polling/utils'
1720
import { processPolledWebhookEvent } from '@/lib/webhooks/processor'
@@ -95,10 +98,9 @@ export const googleCalendarPollingHandler: PollingProviderHandler = {
9598
provider: 'google-calendar',
9699
label: 'Google Calendar',
97100

98-
async pollWebhook(ctx: PollWebhookContext) {
101+
async pollWebhook(ctx: PollWebhookContext): Promise<PollOutcome> {
99102
const { webhookData, workflowData, requestId, logger } = ctx
100103
const webhookId = webhookData.id
101-
const pollStartedAt = Date.now()
102104

103105
try {
104106
const accessToken = await resolveOAuthCredential(webhookData, 'google-calendar', requestId)
@@ -172,13 +174,11 @@ export const googleCalendarPollingHandler: PollingProviderHandler = {
172174
)
173175
return 'success'
174176
} catch (error) {
175-
await recordPollSourceFailure(
176-
webhookData,
177-
pollStartedAt,
178-
error,
179-
`[${requestId}] Error polling Google Calendar webhook ${webhookId}`,
180-
logger
181-
)
177+
if (error instanceof PollAdmissionRefusedError) {
178+
return skipAdmissionRefusedPoll(logger, requestId, webhookId)
179+
}
180+
logger.error(`[${requestId}] Error processing Google Calendar webhook ${webhookId}:`, error)
181+
await markWebhookFailed(webhookId, logger)
182182
return 'failure'
183183
}
184184
},
@@ -336,6 +336,7 @@ async function processEvents(
336336
)
337337

338338
if (!result.success) {
339+
throwIfAdmissionRefused(result)
339340
logger.error(
340341
`[${requestId}] Failed to process webhook for event ${event.id}:`,
341342
result.statusCode,
@@ -353,6 +354,7 @@ async function processEvents(
353354
)
354355
processedCount++
355356
} catch (error) {
357+
if (error instanceof PollAdmissionRefusedError) throw error
356358
const errorMessage = getErrorMessage(error, 'Unknown error')
357359
logger.error(`[${requestId}] Error processing event ${event.id}:`, errorMessage)
358360
failedCount++

‎apps/sim/lib/webhooks/polling/google-drive.ts‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,16 @@ import { readCanonicalTriggerValue } from '@/lib/webhooks/polling/canonical'
55
import {
66
getProviderConfig,
77
type PollingProviderHandler,
8+
type PollOutcome,
89
type PollWebhookContext,
910
} from '@/lib/webhooks/polling/types'
1011
import {
1112
markWebhookFailed,
1213
markWebhookSuccess,
13-
recordPollSourceFailure,
14+
PollAdmissionRefusedError,
1415
resolveOAuthCredential,
16+
skipAdmissionRefusedPoll,
17+
throwIfAdmissionRefused,
1518
updateWebhookProviderConfig,
1619
} from '@/lib/webhooks/polling/utils'
1720
import { processPolledWebhookEvent } from '@/lib/webhooks/processor'
@@ -83,10 +86,9 @@ export const googleDrivePollingHandler: PollingProviderHandler = {
8386
provider: 'google-drive',
8487
label: 'Google Drive',
8588

86-
async pollWebhook(ctx: PollWebhookContext) {
89+
async pollWebhook(ctx: PollWebhookContext): Promise<PollOutcome> {
8790
const { webhookData, workflowData, requestId, logger } = ctx
8891
const webhookId = webhookData.id
89-
const pollStartedAt = Date.now()
9092

9193
try {
9294
const accessToken = await resolveOAuthCredential(webhookData, 'google-drive', requestId)
@@ -171,6 +173,9 @@ export const googleDrivePollingHandler: PollingProviderHandler = {
171173
)
172174
return 'success'
173175
} catch (error) {
176+
if (error instanceof PollAdmissionRefusedError) {
177+
return skipAdmissionRefusedPoll(logger, requestId, webhookId)
178+
}
174179
if (error instanceof Error && error.name === 'DrivePageTokenInvalidError') {
175180
await updateWebhookProviderConfig(webhookId, { pageToken: undefined }, logger)
176181
await markWebhookSuccess(webhookId, logger)
@@ -186,13 +191,8 @@ export const googleDrivePollingHandler: PollingProviderHandler = {
186191
)
187192
return 'success'
188193
}
189-
await recordPollSourceFailure(
190-
webhookData,
191-
pollStartedAt,
192-
error,
193-
`[${requestId}] Error polling Google Drive webhook ${webhookId}`,
194-
logger
195-
)
194+
logger.error(`[${requestId}] Error processing Google Drive webhook ${webhookId}:`, error)
195+
await markWebhookFailed(webhookId, logger)
196196
return 'failure'
197197
}
198198
},
@@ -404,6 +404,7 @@ async function processChanges(
404404
)
405405

406406
if (!result.success) {
407+
throwIfAdmissionRefused(result)
407408
logger.error(
408409
`[${requestId}] Failed to process webhook for file ${change.fileId}:`,
409410
result.statusCode,
@@ -420,6 +421,7 @@ async function processChanges(
420421
)
421422
processedCount++
422423
} catch (error) {
424+
if (error instanceof PollAdmissionRefusedError) throw error
423425
const errorMessage = getErrorMessage(error, 'Unknown error')
424426
logger.error(
425427
`[${requestId}] Error processing change for file ${change.fileId}:`,

‎apps/sim/lib/webhooks/polling/google-sheets.ts‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,16 @@ import { readCanonicalTriggerValue } from '@/lib/webhooks/polling/canonical'
55
import {
66
getProviderConfig,
77
type PollingProviderHandler,
8+
type PollOutcome,
89
type PollWebhookContext,
910
} from '@/lib/webhooks/polling/types'
1011
import {
1112
markWebhookFailed,
1213
markWebhookSuccess,
13-
recordPollSourceFailure,
14+
PollAdmissionRefusedError,
1415
resolveOAuthCredential,
16+
skipAdmissionRefusedPoll,
17+
throwIfAdmissionRefused,
1518
updateWebhookProviderConfig,
1619
} from '@/lib/webhooks/polling/utils'
1720
import { processPolledWebhookEvent } from '@/lib/webhooks/processor'
@@ -52,10 +55,9 @@ export const googleSheetsPollingHandler: PollingProviderHandler = {
5255
provider: 'google-sheets',
5356
label: 'Google Sheets',
5457

55-
async pollWebhook(ctx: PollWebhookContext) {
58+
async pollWebhook(ctx: PollWebhookContext): Promise<PollOutcome> {
5659
const { webhookData, workflowData, requestId, logger } = ctx
5760
const webhookId = webhookData.id
58-
const pollStartedAt = Date.now()
5961

6062
try {
6163
const accessToken = await resolveOAuthCredential(webhookData, 'google-sheets', requestId)
@@ -229,13 +231,11 @@ export const googleSheetsPollingHandler: PollingProviderHandler = {
229231
)
230232
return 'success'
231233
} catch (error) {
232-
await recordPollSourceFailure(
233-
webhookData,
234-
pollStartedAt,
235-
error,
236-
`[${requestId}] Error polling Google Sheets webhook ${webhookId}`,
237-
logger
238-
)
234+
if (error instanceof PollAdmissionRefusedError) {
235+
return skipAdmissionRefusedPoll(logger, requestId, webhookId)
236+
}
237+
logger.error(`[${requestId}] Error processing Google Sheets webhook ${webhookId}:`, error)
238+
await markWebhookFailed(webhookId, logger)
239239
return 'failure'
240240
}
241241
},
@@ -433,6 +433,7 @@ async function processRows(
433433
)
434434

435435
if (!result.success) {
436+
throwIfAdmissionRefused(result)
436437
logger.error(
437438
`[${requestId}] Failed to process webhook for row ${rowNumber}:`,
438439
result.statusCode,
@@ -450,6 +451,7 @@ async function processRows(
450451
)
451452
processedCount++
452453
} catch (error) {
454+
if (error instanceof PollAdmissionRefusedError) throw error
453455
const errorMessage = getErrorMessage(error, 'Unknown error')
454456
logger.error(`[${requestId}] Error processing row ${rowNumber}:`, errorMessage)
455457
failedCount++

0 commit comments

Comments
 (0)