Skip to content

Commit 3a141dd

Browse files
committed
chore(webhooks): key blocked-run claims by gate and centralize the polling utils mock
Two gates that fail without a code (a ban lookup error and a usage lookup error) no longer share one throttle claim, so neither hides the other's row. The polling utils module gets one central mock in @sim/testing, replacing the partial importOriginal mocks and the hand-rolled factory in the table trigger test.
1 parent 317463a commit 3a141dd

8 files changed

Lines changed: 155 additions & 48 deletions

File tree

‎apps/sim/lib/execution/preprocessing.test.ts‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1065,6 +1065,16 @@ describe('preprocessExecution admission rejection codes and blocked-run log thro
10651065
expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalledTimes(2)
10661066
})
10671067

1068+
it('writes a row for each gate whose check fails without a code', async () => {
1069+
const workflowId = 'workflow-1'
1070+
mockGetActivelyBannedUserIds.mockRejectedValueOnce(new Error('ban lookup failed'))
1071+
await refuse(workflowId, { throttleErrorLogs: true })
1072+
mockCheckAttributedUsageLimits.mockRejectedValueOnce(new Error('usage lookup failed'))
1073+
await refuse(workflowId, { throttleErrorLogs: true })
1074+
1075+
expect(loggingSessionMockFns.mockSafeCompleteWithError).toHaveBeenCalledTimes(2)
1076+
})
1077+
10681078
it('writes every row when the caller does not ask for throttling', async () => {
10691079
const workflowId = 'workflow-1'
10701080
await refuse(workflowId)

‎apps/sim/lib/execution/preprocessing.ts‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -382,8 +382,9 @@ export async function preprocessExecution(
382382
const isFailureLogSuppressed = (failure: PreprocessExecutionError): boolean =>
383383
suppressRetryableFailureLogs && failure.statusCode >= 500 && failure.retryable === true
384384

385-
/** Records an admission gate's error row, at most once per window when throttled. */
385+
/** Records an admission gate's error row, at most once per gate and outcome per window when throttled. */
386386
const recordGateFailure = async (
387+
gate: 'ban' | 'usage' | 'rate-limit' | 'reservation',
387388
failure: PreprocessExecutionError,
388389
record: Parameters<typeof logPreprocessingError>[0]
389390
): Promise<void> => {
@@ -392,7 +393,7 @@ export async function preprocessExecution(
392393
throttleErrorLogs &&
393394
logPreprocessingErrors &&
394395
!providedLoggingSession &&
395-
!(await claimBlockedRunLog(workflowId, failure.code ?? String(failure.statusCode)))
396+
!(await claimBlockedRunLog(workflowId, `${gate}:${failure.code ?? failure.statusCode}`))
396397
) {
397398
return
398399
}
@@ -919,15 +920,23 @@ export async function preprocessExecution(
919920
const readGateFailure = banFailure ?? usageResult.failure
920921
if (readGateFailure) {
921922
if (readGateFailure.recordError) {
922-
await recordGateFailure(readGateFailure.response.error, readGateFailure.recordError)
923+
await recordGateFailure(
924+
banFailure ? 'ban' : 'usage',
925+
readGateFailure.response.error,
926+
readGateFailure.recordError
927+
)
923928
}
924929
return readGateFailure.response
925930
}
926931

927932
const rateLimitFailure = await runRateLimitGate()
928933
if (rateLimitFailure) {
929934
if (rateLimitFailure.recordError) {
930-
await recordGateFailure(rateLimitFailure.response.error, rateLimitFailure.recordError)
935+
await recordGateFailure(
936+
'rate-limit',
937+
rateLimitFailure.response.error,
938+
rateLimitFailure.recordError
939+
)
931940
}
932941
return rateLimitFailure.response
933942
}
@@ -993,7 +1002,7 @@ export async function preprocessExecution(
9931002
constraint: reservation.reason,
9941003
},
9951004
}
996-
await recordGateFailure(failure, {
1005+
await recordGateFailure('reservation', failure, {
9971006
workflowId,
9981007
executionId,
9991008
triggerType,

‎apps/sim/lib/table/trigger.test.ts‎

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -4,25 +4,24 @@
44
* test would pass green, so every assertion runs on the captured payload AFTER
55
* the `await`, never inside a mock factory.
66
*/
7+
import {
8+
webhooksPollingUtilsMock,
9+
webhooksPollingUtilsMockFns,
10+
} from '@sim/testing/mocks/webhooks-polling-utils.mock'
711
import {
812
webhooksProcessorMock,
913
webhooksProcessorMockFns,
1014
} from '@sim/testing/mocks/webhooks-processor.mock'
1115
import { beforeEach, describe, expect, it, vi } from 'vitest'
1216

13-
const { mockFetchActiveWebhooks } = vi.hoisted(() => ({
14-
mockFetchActiveWebhooks: vi.fn(),
15-
}))
16-
17-
vi.mock('@/lib/webhooks/polling/utils', () => ({
18-
fetchActiveWebhooks: mockFetchActiveWebhooks,
19-
}))
17+
vi.mock('@/lib/webhooks/polling/utils', () => webhooksPollingUtilsMock)
2018
vi.mock('@/lib/webhooks/processor', () => webhooksProcessorMock)
2119

2220
import { fireTableTrigger } from '@/lib/table/trigger'
2321
import type { RowData, TableSchema } from '@/lib/table/types'
2422

2523
const mockProcessPolledWebhookEvent = webhooksProcessorMockFns.mockProcessPolledWebhookEvent
24+
const { mockFetchActiveWebhooks } = webhooksPollingUtilsMockFns
2625

2726
const schema: TableSchema = {
2827
columns: [

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

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,16 @@
11
import { createLogger } from '@sim/logger'
22
import { createWorkflowRecord } from '@sim/testing'
33
import { jsonResponse } from '@sim/testing/helpers/http'
4+
import {
5+
webhooksPollingUtilsMock,
6+
webhooksPollingUtilsMockFns,
7+
} from '@sim/testing/mocks/webhooks-polling-utils.mock'
48
import {
59
webhooksProcessorMock,
610
webhooksProcessorMockFns,
711
} from '@sim/testing/mocks/webhooks-processor.mock'
812
import { beforeEach, describe, expect, it, vi } from 'vitest'
913

10-
const { mockUpdateConfig, mockMarkFailed } = vi.hoisted(() => ({
11-
mockUpdateConfig: vi.fn(),
12-
mockMarkFailed: vi.fn(),
13-
}))
14-
1514
vi.mock('@/lib/core/idempotency/service', () => ({
1615
pollingIdempotency: {
1716
executeWithIdempotency: vi.fn(
@@ -22,19 +21,18 @@ vi.mock('@/lib/core/idempotency/service', () => ({
2221

2322
vi.mock('@/lib/webhooks/processor', () => webhooksProcessorMock)
2423

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-
}))
24+
vi.mock('@/lib/webhooks/polling/utils', () => webhooksPollingUtilsMock)
3225

3326
import { ADMISSION_REJECTION_CODE } from '@/lib/core/admission/rejection'
3427
import { googleCalendarPollingHandler } from '@/lib/webhooks/polling/google-calendar'
3528
import type { PollWebhookContext, WebhookRecord } from '@/lib/webhooks/polling/types'
3629

3730
const mockProcessEvent = webhooksProcessorMockFns.mockProcessPolledWebhookEvent
31+
const {
32+
mockUpdateWebhookProviderConfig: mockUpdateConfig,
33+
mockMarkWebhookFailed: mockMarkFailed,
34+
mockResolveOAuthCredential,
35+
} = webhooksPollingUtilsMockFns
3836

3937
function context(): PollWebhookContext {
4038
const webhookData: WebhookRecord = {
@@ -72,6 +70,7 @@ function context(): PollWebhookContext {
7270

7371
describe('Google Calendar polling when execution admission refuses events', () => {
7472
beforeEach(() => {
73+
mockResolveOAuthCredential.mockResolvedValue('access-token')
7574
const events = ['event-1', 'event-2'].map((id) => ({
7675
id,
7776
status: 'confirmed',

‎apps/sim/lib/webhooks/polling/imap.test.ts‎

Lines changed: 6 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,20 @@
11
import { dbChainMockFns } from '@sim/testing'
2+
import {
3+
webhooksPollingUtilsMock,
4+
webhooksPollingUtilsMockFns,
5+
} from '@sim/testing/mocks/webhooks-polling-utils.mock'
26
import { webhooksProcessorMock } from '@sim/testing/mocks/webhooks-processor.mock'
37
import { beforeEach, describe, expect, it, vi } from 'vitest'
48

59
const {
610
mockCreateSecureImapClient,
711
mockHasImapEnvironmentReferences,
812
mockLogger,
9-
mockMarkWebhookFailed,
1013
mockResolveImapConnectionForActor,
1114
} = vi.hoisted(() => ({
1215
mockCreateSecureImapClient: vi.fn(),
1316
mockHasImapEnvironmentReferences: vi.fn(),
1417
mockLogger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() },
15-
mockMarkWebhookFailed: vi.fn(),
1618
mockResolveImapConnectionForActor: vi.fn(),
1719
}))
1820

@@ -27,23 +29,18 @@ vi.mock('@/lib/imap/connection.server', () => ({
2729
resolveImapConnectionForActor: mockResolveImapConnectionForActor,
2830
}))
2931

30-
vi.mock('@/lib/webhooks/polling/utils', async (importOriginal) => ({
31-
...(await importOriginal<typeof import('@/lib/webhooks/polling/utils')>()),
32-
markWebhookFailed: mockMarkWebhookFailed,
33-
markWebhookSuccess: vi.fn(),
34-
updateWebhookProviderConfig: vi.fn(),
35-
}))
32+
vi.mock('@/lib/webhooks/polling/utils', () => webhooksPollingUtilsMock)
3633

3734
vi.mock('@/lib/webhooks/processor', () => webhooksProcessorMock)
3835

3936
import { imapPollingHandler } from '@/lib/webhooks/polling/imap'
4037

4138
const mockDbSelect = dbChainMockFns.select
39+
const { mockMarkWebhookFailed } = webhooksPollingUtilsMockFns
4240

4341
describe('IMAP runtime polling policy', () => {
4442
beforeEach(() => {
4543
mockHasImapEnvironmentReferences.mockReturnValue(true)
46-
mockMarkWebhookFailed.mockResolvedValue(undefined)
4744
})
4845

4946
it('fails closed before resolution, DNS, or ImapFlow when referenced auth has no deployment actor', async () => {

‎apps/sim/lib/webhooks/polling/rss.test.ts‎

Lines changed: 12 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -4,18 +4,16 @@ import {
44
inputValidationMock,
55
inputValidationMockFns,
66
} from '@sim/testing/mocks/input-validation.mock'
7+
import {
8+
webhooksPollingUtilsMock,
9+
webhooksPollingUtilsMockFns,
10+
} from '@sim/testing/mocks/webhooks-polling-utils.mock'
711
import {
812
webhooksProcessorMock,
913
webhooksProcessorMockFns,
1014
} from '@sim/testing/mocks/webhooks-processor.mock'
1115
import { beforeEach, describe, expect, it, vi } from 'vitest'
1216

13-
const { mockUpdateConfig, mockMarkFailed, mockRecordPollSourceFailure } = vi.hoisted(() => ({
14-
mockUpdateConfig: vi.fn(),
15-
mockMarkFailed: vi.fn(),
16-
mockRecordPollSourceFailure: vi.fn(),
17-
}))
18-
1917
vi.mock('@/lib/core/security/input-validation.server', () => inputValidationMock)
2018
const mockFetch = inputValidationMockFns.mockSecureFetchWithPinnedIP
2119
const mockValidateUrl = inputValidationMockFns.mockValidateUrlWithDNS
@@ -30,20 +28,19 @@ vi.mock('@/lib/core/idempotency/service', () => ({
3028

3129
vi.mock('@/lib/webhooks/processor', () => webhooksProcessorMock)
3230

33-
vi.mock('@/lib/webhooks/polling/utils', async (importOriginal) => ({
34-
...(await importOriginal<typeof import('@/lib/webhooks/polling/utils')>()),
35-
markWebhookSuccess: vi.fn(),
36-
markWebhookFailed: mockMarkFailed,
37-
recordPollSourceFailure: mockRecordPollSourceFailure,
38-
updateWebhookProviderConfig: mockUpdateConfig,
39-
}))
31+
vi.mock('@/lib/webhooks/polling/utils', () => webhooksPollingUtilsMock)
4032

4133
import { ADMISSION_REJECTION_CODE } from '@/lib/core/admission/rejection'
4234
import { rssPollingHandler } from '@/lib/webhooks/polling/rss'
4335
import type { PollWebhookContext, WebhookRecord } from '@/lib/webhooks/polling/types'
4436
import { PollFetchError } from '@/lib/webhooks/polling/utils'
4537

4638
const mockProcessEvent = webhooksProcessorMockFns.mockProcessPolledWebhookEvent
39+
const {
40+
mockUpdateWebhookProviderConfig: mockUpdateConfig,
41+
mockMarkWebhookFailed: mockMarkFailed,
42+
mockRecordPollSourceFailure,
43+
} = webhooksPollingUtilsMockFns
4744

4845
const SUBSCRIBED_AT = new Date('2026-08-27T18:36:16.000Z')
4946
const LAST_CHECKED_AT = '2026-09-11T23:26:27.000Z'
@@ -158,7 +155,7 @@ describe('RSS polling against refusals and rate limits', () => {
158155
expect(mockMarkFailed).not.toHaveBeenCalled()
159156
})
160157

161-
it("records a rate-limited fetch as one failure carrying the source's requested wait", async () => {
158+
it('records a rate-limited fetch as one source failure carrying its status', async () => {
162159
mockFetch.mockResolvedValue(
163160
new Response('Too Many Requests', {
164161
status: 429,
@@ -172,6 +169,6 @@ describe('RSS polling against refusals and rate limits', () => {
172169
expect(mockRecordPollSourceFailure).toHaveBeenCalledOnce()
173170
const [, , error] = mockRecordPollSourceFailure.mock.calls[0]
174171
expect(error).toBeInstanceOf(PollFetchError)
175-
expect(error).toMatchObject({ status: 429, retryAfterMs: 12_000 })
172+
expect(error).toMatchObject({ status: 429 })
176173
})
177174
})

‎packages/testing/src/mocks/index.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -938,6 +938,10 @@ export {
938938
v2RateLimiterModuleMock,
939939
v2RouteMocks,
940940
} from './v2-route.mock'
941+
export {
942+
webhooksPollingUtilsMock,
943+
webhooksPollingUtilsMockFns,
944+
} from './webhooks-polling-utils.mock'
941945
export {
942946
webhooksProcessorMock,
943947
webhooksProcessorMockFns,
Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,92 @@
1+
import { vi } from 'vitest'
2+
3+
/** Faithful copy of the production deterministic admission codes (`lib/core/admission/rejection`). */
4+
const DETERMINISTIC_ADMISSION_REJECTION_CODES = new Set([
5+
'USAGE_LIMIT_EXCEEDED',
6+
'ACCOUNT_SUSPENDED',
7+
])
8+
9+
/** Faithful copy of the production `PollAdmissionRefusedError`. */
10+
class PollAdmissionRefusedError extends Error {
11+
constructor(result: { statusCode?: number; error?: string }) {
12+
super(`Execution admission refused (${result.statusCode}): ${result.error}`)
13+
this.name = 'PollAdmissionRefusedError'
14+
}
15+
}
16+
17+
/** Faithful copy of the production `PollFetchError`. */
18+
class PollFetchError extends Error {
19+
readonly status: number
20+
readonly retryAfterMs: number | null
21+
22+
constructor(message: string, status: number, retryAfterMs: number | null) {
23+
super(message)
24+
this.name = 'PollFetchError'
25+
this.status = status
26+
this.retryAfterMs = retryAfterMs
27+
}
28+
}
29+
30+
/**
31+
* Controllable mock functions for `@/lib/webhooks/polling/utils`.
32+
*
33+
* State writes (`markWebhookFailed`, `markWebhookSuccess`, `updateWebhookProviderConfig`,
34+
* `recordPollSourceFailure`) resolve `undefined`; `fetchActiveWebhooks` resolves `[]`;
35+
* `resolveOAuthCredential` is a bare `vi.fn()`. `throwIfAdmissionRefused` and
36+
* `skipAdmissionRefusedPoll` keep production behavior so a poller's refusal path runs as it
37+
* does in production. Defaults are `vi.fn(impl)`, so `mockReset()` restores them.
38+
*
39+
* @example
40+
* ```ts
41+
* import { webhooksPollingUtilsMockFns } from '@sim/testing/mocks/webhooks-polling-utils.mock'
42+
*
43+
* webhooksPollingUtilsMockFns.mockResolveOAuthCredential.mockResolvedValue('access-token')
44+
* ```
45+
*/
46+
export const webhooksPollingUtilsMockFns = {
47+
mockIsPollBackedOff: vi.fn((_providerConfig: unknown, _now: number): boolean => false),
48+
mockThrowIfAdmissionRefused: vi.fn(
49+
(result: { code?: string; statusCode?: number; error?: string }): void => {
50+
if (result.code && DETERMINISTIC_ADMISSION_REJECTION_CODES.has(result.code)) {
51+
throw new PollAdmissionRefusedError(result)
52+
}
53+
}
54+
),
55+
mockSkipAdmissionRefusedPoll: vi.fn((..._args: unknown[]): 'skipped' => 'skipped'),
56+
mockReadPollRetryAfterMs: vi.fn((_header: string | null, _body: string): number | null => null),
57+
mockClearPollBackoff: vi.fn((_providerConfig: unknown): Record<string, undefined> => ({})),
58+
mockRecordPollSourceFailure: vi.fn(async (..._args: unknown[]): Promise<void> => {}),
59+
mockMarkWebhookFailed: vi.fn(async (..._args: unknown[]): Promise<void> => {}),
60+
mockMarkWebhookSuccess: vi.fn(async (..._args: unknown[]): Promise<void> => {}),
61+
mockFetchActiveWebhooks: vi.fn(async (..._args: unknown[]): Promise<unknown[]> => []),
62+
mockRunWithConcurrency: vi.fn(),
63+
mockUpdateWebhookProviderConfig: vi.fn(async (..._args: unknown[]): Promise<void> => {}),
64+
mockResolveOAuthCredential: vi.fn(),
65+
}
66+
67+
/**
68+
* Static mock module for `@/lib/webhooks/polling/utils`. Covers every runtime export; the
69+
* error classes and `CONCURRENCY` are faithful copies of production.
70+
*
71+
* @example
72+
* ```ts
73+
* vi.mock('@/lib/webhooks/polling/utils', () => webhooksPollingUtilsMock)
74+
* ```
75+
*/
76+
export const webhooksPollingUtilsMock = {
77+
CONCURRENCY: 10,
78+
PollAdmissionRefusedError,
79+
PollFetchError,
80+
isPollBackedOff: webhooksPollingUtilsMockFns.mockIsPollBackedOff,
81+
throwIfAdmissionRefused: webhooksPollingUtilsMockFns.mockThrowIfAdmissionRefused,
82+
skipAdmissionRefusedPoll: webhooksPollingUtilsMockFns.mockSkipAdmissionRefusedPoll,
83+
readPollRetryAfterMs: webhooksPollingUtilsMockFns.mockReadPollRetryAfterMs,
84+
clearPollBackoff: webhooksPollingUtilsMockFns.mockClearPollBackoff,
85+
recordPollSourceFailure: webhooksPollingUtilsMockFns.mockRecordPollSourceFailure,
86+
markWebhookFailed: webhooksPollingUtilsMockFns.mockMarkWebhookFailed,
87+
markWebhookSuccess: webhooksPollingUtilsMockFns.mockMarkWebhookSuccess,
88+
fetchActiveWebhooks: webhooksPollingUtilsMockFns.mockFetchActiveWebhooks,
89+
runWithConcurrency: webhooksPollingUtilsMockFns.mockRunWithConcurrency,
90+
updateWebhookProviderConfig: webhooksPollingUtilsMockFns.mockUpdateWebhookProviderConfig,
91+
resolveOAuthCredential: webhooksPollingUtilsMockFns.mockResolveOAuthCredential,
92+
}

0 commit comments

Comments
 (0)