Skip to content

Commit 5e44efe

Browse files
committed
fix(billing): lock the subscription before the outbox row in every sync retry path
1 parent 40ec5a8 commit 5e44efe

9 files changed

Lines changed: 185 additions & 30 deletions

File tree

‎apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts‎

Lines changed: 17 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import {
1414
} from '@/lib/billing/enterprise-outbox-events'
1515
import {
1616
isSubscriptionSyncEventType,
17+
lockSubscriptionForSyncRetry,
1718
recommitSubscriptionSync,
1819
} from '@/lib/billing/webhooks/subscription-sync'
1920
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
@@ -68,7 +69,17 @@ export const POST = withRouteHandler(
6869
const deliveryRevision = metadataIntent?.success
6970
? metadataIntent.data.deliveryRevision + 1
7071
: null
72+
const subscriptionId = toRecord(existing?.payload).subscriptionId
73+
const subscriptionSync =
74+
existing &&
75+
isSubscriptionSyncEventType(existing.eventType) &&
76+
typeof subscriptionId === 'string'
77+
? { eventType: existing.eventType, subscriptionId }
78+
: null
7179
const result = await db.transaction(async (tx) => {
80+
if (subscriptionSync) {
81+
await lockSubscriptionForSyncRetry(tx, subscriptionSync.subscriptionId)
82+
}
7283
const requeued = await tx
7384
.update(outboxEvent)
7485
.set({
@@ -86,13 +97,12 @@ export const POST = withRouteHandler(
8697
})
8798
.where(and(eq(outboxEvent.id, id), eq(outboxEvent.status, 'dead_letter')))
8899
.returning({ id: outboxEvent.id, eventType: outboxEvent.eventType })
89-
const subscriptionId = toRecord(existing?.payload).subscriptionId
90-
if (
91-
requeued.length > 0 &&
92-
isSubscriptionSyncEventType(requeued[0].eventType) &&
93-
typeof subscriptionId === 'string'
94-
) {
95-
await recommitSubscriptionSync(tx, requeued[0].eventType, subscriptionId)
100+
if (subscriptionSync && requeued.length > 0) {
101+
await recommitSubscriptionSync(
102+
tx,
103+
subscriptionSync.eventType,
104+
subscriptionSync.subscriptionId
105+
)
96106
}
97107
return requeued
98108
})

‎apps/sim/lib/admin/subscription-lifecycle.test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,7 @@ describe('admin subscription cancellation', () => {
103103

104104
it('requeues the same dead-lettered period-end cancellation operation', async () => {
105105
dbChainMockFns.returning.mockResolvedValueOnce([{ id: activeSubscription.id }])
106+
queueTableRows(outboxEvent, [{ subscriptionId: 'sub-row-1' }])
106107
queueTableRows(outboxEvent, [
107108
{
108109
id: 'outbox-1',
@@ -124,6 +125,10 @@ describe('admin subscription cancellation', () => {
124125
expect.objectContaining({ status: 'pending', attempts: 0, lastError: null })
125126
)
126127
expect(dbChainMockFns.set).toHaveBeenCalledWith({ cancelAtPeriodEnd: true })
128+
expect(billingSubscriptionSyncMockFns.mockLockSubscriptionForSyncRetry).toHaveBeenCalledWith(
129+
expect.anything(),
130+
'sub-row-1'
131+
)
127132
expect(billingSubscriptionSyncMockFns.mockRecommitSubscriptionSync).toHaveBeenCalledWith(
128133
expect.anything(),
129134
'stripe.sync-cancel-at-period-end',
@@ -136,6 +141,7 @@ describe('admin subscription cancellation', () => {
136141
})
137142

138143
it('replays an immediate cancellation after the webhook removed active entitlement', async () => {
144+
queueTableRows(outboxEvent, [])
139145
queueTableRows(outboxEvent, [
140146
{
141147
id: 'outbox-1',
@@ -159,6 +165,7 @@ describe('admin subscription cancellation', () => {
159165
})
160166

161167
it('rejects reuse of a cancellation operation id with different timing', async () => {
168+
queueTableRows(outboxEvent, [])
162169
queueTableRows(outboxEvent, [
163170
{
164171
id: 'outbox-1',

‎apps/sim/lib/admin/subscription-lifecycle.ts‎

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/util
1010
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
1111
import {
1212
enqueueCancelAtPeriodEndSync,
13+
lockSubscriptionForSyncRetry,
1314
recommitSubscriptionSync,
1415
} from '@/lib/billing/webhooks/subscription-sync'
1516
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
@@ -261,6 +262,24 @@ export async function requestDashboardSubscriptionCancellation({
261262
: 'admin-dashboard-cancel-at-period-end')
262263
const cancellation = await db.transaction(async (tx) => {
263264
await acquireOrganizationMutationLock(tx, organizationId)
265+
const isThisOperation = and(
266+
sql`${outboxEvent.payload} ->> 'operationId' = ${operationId}`,
267+
sql`${outboxEvent.payload} ->> 'organizationId' = ${organizationId}`
268+
)
269+
270+
const [retriedSync] = await tx
271+
.select({ subscriptionId: sql<string | null>`${outboxEvent.payload} ->> 'subscriptionId'` })
272+
.from(outboxEvent)
273+
.where(
274+
and(
275+
eq(outboxEvent.eventType, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END),
276+
isThisOperation
277+
)
278+
)
279+
.limit(1)
280+
if (retriedSync?.subscriptionId) {
281+
await lockSubscriptionForSyncRetry(tx, retriedSync.subscriptionId)
282+
}
264283

265284
const [existingOperation] = await tx
266285
.select({
@@ -277,8 +296,7 @@ export async function requestDashboardSubscriptionCancellation({
277296
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
278297
OUTBOX_EVENT_TYPES.STRIPE_CANCEL_SUBSCRIPTION_IMMEDIATELY,
279298
]),
280-
sql`${outboxEvent.payload} ->> 'operationId' = ${operationId}`,
281-
sql`${outboxEvent.payload} ->> 'organizationId' = ${organizationId}`
299+
isThisOperation
282300
)
283301
)
284302
.for('update')

‎apps/sim/lib/billing/enterprise-provisioning.ts‎

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,10 @@ import { TERMINAL_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/util
7777
import { countPendingSeatInvitations } from '@/lib/billing/validation/seat-management'
7878
import { withEnterpriseReconciliationLease } from '@/lib/billing/webhooks/enterprise-reconciliation-lease'
7979
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
80-
import { recommitSubscriptionSync } from '@/lib/billing/webhooks/subscription-sync'
80+
import {
81+
lockSubscriptionForSyncRetry,
82+
recommitSubscriptionSync,
83+
} from '@/lib/billing/webhooks/subscription-sync'
8184
import { env } from '@/lib/core/config/env'
8285
import {
8386
continueOutboxHandler,
@@ -2104,6 +2107,9 @@ export async function retryEnterpriseFollowUpJob(
21042107

21052108
const retried = await db.transaction(async (tx) => {
21062109
await acquireOrganizationMutationLock(tx, operationPayload.request.organizationId)
2110+
if (snapshotDetail.kind === 'personal_subscription_cancellation') {
2111+
await lockSubscriptionForSyncRetry(tx, snapshotDetail.subjectId)
2112+
}
21072113
const [row] = await tx
21082114
.select({
21092115
status: outboxEvent.status,
@@ -2120,7 +2126,9 @@ export async function retryEnterpriseFollowUpJob(
21202126
!detail ||
21212127
!getEnterpriseFollowUpOperationIds(row.eventType, row.payload).includes(operationId) ||
21222128
(detail.kind === 'member_reconciliation' &&
2123-
detail.subjectId !== operationPayload.request.organizationId)
2129+
detail.subjectId !== operationPayload.request.organizationId) ||
2130+
detail.kind !== snapshotDetail.kind ||
2131+
detail.subjectId !== snapshotDetail.subjectId
21242132
) {
21252133
throw new EnterpriseProvisioningError('Enterprise follow-up job not found')
21262134
}

‎apps/sim/lib/billing/webhooks/stripe-sync-convergence.integration.ts‎

Lines changed: 97 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ vi.mock('@sim/db', () => ({
4343
vi.mock('@/lib/billing/stripe-client', () => stripeClientMock)
4444
vi.mock('@/lib/core/config/env-flags', () => envFlagsMock)
4545

46+
import { requestDashboardSubscriptionCancellation } from '@/lib/admin/subscription-lifecycle'
4647
import { createSimAuthAdapter } from '@/lib/auth/sim-auth-adapter'
4748
import {
4849
pauseProSubscriptionForOrgCoverage,
@@ -132,6 +133,7 @@ beforeAll(async () => {
132133
'organization',
133134
'workspace',
134135
'permissions',
136+
'audit_log',
135137
]) {
136138
await connection.unsafe(`CREATE TABLE "${table}" (LIKE public."${table}" INCLUDING ALL)`)
137139
}
@@ -283,6 +285,18 @@ async function storedSubscription(subscriptionId: string) {
283285
return row
284286
}
285287

288+
/** Resolves once another backend is blocked on a lock, i.e. the racing transaction is parked. */
289+
async function untilAnotherTransactionWaitsOnALock() {
290+
for (let attempt = 0; attempt < 200; attempt++) {
291+
const [row] = await connection<{ waiting: number }[]>`
292+
select count(*)::int as waiting from pg_stat_activity
293+
where datname = current_database() and wait_event_type = 'Lock'`
294+
if (row.waiting > 0) return
295+
await new Promise<void>((resolve) => setImmediate(resolve))
296+
}
297+
throw new Error('The racing transaction never waited on a lock')
298+
}
299+
286300
describe('cancel_at_period_end sync', () => {
287301
it('pushes the latest value when an earlier sync lands in Stripe after a newer one', async () => {
288302
const pro = await createProUserInPaidOrganization()
@@ -618,18 +632,6 @@ describe('cancel_at_period_end sync', () => {
618632
})
619633

620634
describe('Team activation', () => {
621-
/** Resolves once another backend is blocked on a lock, i.e. the racing transaction is parked. */
622-
async function untilAnotherTransactionWaitsOnALock() {
623-
for (let attempt = 0; attempt < 200; attempt++) {
624-
const [row] = await connection<{ waiting: number }[]>`
625-
select count(*)::int as waiting from pg_stat_activity
626-
where datname = current_database() and wait_event_type = 'Lock'`
627-
if (row.waiting > 0) return
628-
await new Promise<void>((resolve) => setImmediate(resolve))
629-
}
630-
throw new Error('The racing transaction never waited on a lock')
631-
}
632-
633635
it('records the cleared cancellation when a cancel is committed while it activates Team', async () => {
634636
const owner = await createUser('owner')
635637
const subscriptionId = generateId()
@@ -688,6 +690,89 @@ describe('Team activation', () => {
688690
})
689691
})
690692

693+
describe('operator retry', () => {
694+
it('requeues a dead-lettered sync while a writer commits a new value for the subscription', async () => {
695+
const pro = await createProUserInPaidOrganization()
696+
await pauseProSubscriptionForOrgCoverage(pro.userId)
697+
const pauseSync = await latestOutboxEventId(
698+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
699+
pro.subscriptionId
700+
)
701+
await deadLetter(pauseSync)
702+
703+
let releaseWriter: () => void = () => {}
704+
const writerHeld = new Promise<void>((resolve) => {
705+
releaseWriter = resolve
706+
})
707+
const writing = testDatabase.transaction(async (tx) => {
708+
await tx
709+
.update(subscription)
710+
.set({ cancelAtPeriodEnd: false })
711+
.where(eq(subscription.id, pro.subscriptionId))
712+
await writerHeld
713+
await enqueueCancelAtPeriodEndSync(tx, {
714+
stripeSubscriptionId: pro.stripeSubscriptionId,
715+
subscriptionId: pro.subscriptionId,
716+
cancelAtPeriodEnd: false,
717+
reason: 'member-left-paid-org',
718+
})
719+
})
720+
const requeuing = requeueFromAdminApi(pauseSync)
721+
await untilAnotherTransactionWaitsOnALock()
722+
releaseWriter()
723+
724+
await expect(Promise.all([writing, requeuing])).resolves.toBeDefined()
725+
await deliverUnrelatedUpdate(pro.stripeSubscriptionId)
726+
expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false)
727+
})
728+
729+
it('retries a dead-lettered dashboard cancellation while a writer commits a new value', async () => {
730+
const org = await createOrganizationWithPlan('team')
731+
const operationId = generateId()
732+
const actor = { id: null, name: 'Admin', email: null }
733+
await requestDashboardSubscriptionCancellation({
734+
organizationId: org.organizationId,
735+
operationId,
736+
timing: 'period_end',
737+
actor,
738+
})
739+
const cancelSync = await latestOutboxEventId(
740+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
741+
org.subscriptionId
742+
)
743+
await deadLetter(cancelSync)
744+
745+
let releaseWriter: () => void = () => {}
746+
const writerHeld = new Promise<void>((resolve) => {
747+
releaseWriter = resolve
748+
})
749+
const writing = testDatabase.transaction(async (tx) => {
750+
await tx
751+
.update(subscription)
752+
.set({ cancelAtPeriodEnd: false })
753+
.where(eq(subscription.id, org.subscriptionId))
754+
await writerHeld
755+
await enqueueCancelAtPeriodEndSync(tx, {
756+
stripeSubscriptionId: org.stripeSubscriptionId,
757+
subscriptionId: org.subscriptionId,
758+
cancelAtPeriodEnd: false,
759+
reason: 'pro-to-team-conversion',
760+
})
761+
})
762+
const retrying = requestDashboardSubscriptionCancellation({
763+
organizationId: org.organizationId,
764+
operationId,
765+
timing: 'period_end',
766+
actor,
767+
})
768+
await untilAnotherTransactionWaitsOnALock()
769+
releaseWriter()
770+
771+
await expect(Promise.all([writing, retrying])).resolves.toBeDefined()
772+
expect((await storedSubscription(org.subscriptionId)).cancelAtPeriodEnd).toBe(true)
773+
})
774+
})
775+
691776
describe('customer contact sync', () => {
692777
it('pushes the current owner when an earlier sync lands in Stripe after a newer one', async () => {
693778
const [first, second, third] = await Promise.all([

‎apps/sim/lib/billing/webhooks/subscription-sync.ts‎

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import { requireStripeClient } from '@/lib/billing/stripe-client'
99
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
1010
import {
1111
enqueueOutboxEvent,
12+
INFLIGHT_OUTBOX_STATUSES,
1213
listRetryableOutboxEvents,
1314
patchRetryableOutboxEvents,
1415
} from '@/lib/core/outbox/service'
@@ -90,7 +91,8 @@ async function withCommittedAt<T extends SyncIntentFields>(
9091
* onto every event of that type that can still run: pending, processing, or dead-lettered (each
9192
* operator retry path resets dead letters to pending). No event that can run again ever carries
9293
* an older value for the webhook reconcile to restore. The caller must hold the subscription row
93-
* lock (`FOR UPDATE`, or the `UPDATE` itself).
94+
* lock (`FOR UPDATE`, or the `UPDATE` itself), per the lock order on
95+
* {@link lockSubscriptionForSyncRetry}.
9496
*/
9597
async function commitIntent<T extends SyncIntentFields>(
9698
tx: DbOrTx,
@@ -150,9 +152,32 @@ export async function recordCancelAtPeriodEnd(
150152
}
151153

152154
/**
153-
* Re-commits the subscription's current DB value onto its in-flight sync events, for a
155+
* Takes the subscription row lock for an operator retry of one of its sync events; call it
156+
* before touching the event, then {@link recommitSubscriptionSync} after resetting it.
157+
*
158+
* Lock order for every writer of a subscription's synced fields and their outbox events:
159+
* organization mutation lock (where taken) → subscription row → outbox rows. Committing a value
160+
* rewrites the subscription's retryable sync events, dead letters included, so a retry that
161+
* locked a dead-lettered event before the subscription would deadlock against any concurrent
162+
* writer.
163+
*/
164+
export async function lockSubscriptionForSyncRetry(
165+
tx: DbOrTx,
166+
subscriptionId: string
167+
): Promise<void> {
168+
await tx
169+
.select({ id: subscription.id })
170+
.from(subscription)
171+
.where(eq(subscription.id, subscriptionId))
172+
.for('update')
173+
.limit(1)
174+
}
175+
176+
/**
177+
* Re-commits the subscription's current DB value onto its sync events that can still run, for a
154178
* dead-lettered event that was just reset to `pending`: the retry then carries the latest value
155-
* rather than the one it failed with. Takes the subscription row lock itself.
179+
* rather than the one it failed with. The caller holds the lock from
180+
* {@link lockSubscriptionForSyncRetry}, taken before the reset.
156181
*/
157182
export async function recommitSubscriptionSync(
158183
tx: DbOrTx,
@@ -163,7 +188,6 @@ export async function recommitSubscriptionSync(
163188
.select({ cancelAtPeriodEnd: subscription.cancelAtPeriodEnd, seats: subscription.seats })
164189
.from(subscription)
165190
.where(eq(subscription.id, subscriptionId))
166-
.for('update')
167191
.limit(1)
168192
if (!current) return
169193

@@ -229,7 +253,7 @@ export function cancelAtPeriodEndSyncIdempotencyKey(eventId: string): string {
229253
*/
230254
type InflightIntent<T> = { status: 'none' } | { status: 'legacy' } | { status: 'value'; value: T }
231255

232-
const INFLIGHT_STATUSES = new Set(['pending', 'processing'])
256+
const INFLIGHT_STATUSES: ReadonlySet<string> = new Set(INFLIGHT_OUTBOX_STATUSES)
233257

234258
function latestIntent<T>(
235259
events: { eventType: string; status: string; payload: unknown }[],

‎apps/sim/lib/core/outbox/service.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -407,7 +407,7 @@ export async function findDeadLetteredEvents(
407407
}
408408

409409
/** Statuses of an event whose side effect may still run. */
410-
const INFLIGHT_OUTBOX_STATUSES = ['pending', 'processing'] as const
410+
export const INFLIGHT_OUTBOX_STATUSES = ['pending', 'processing'] as const
411411
/**
412412
* Statuses an event can still run from: in flight, or dead-lettered, which every operator retry
413413
* path resets to `pending`. A `completed` event never runs again.

‎packages/testing/src/mocks/billing-subscription-sync.mock.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ import { vi } from 'vitest'
1717
*/
1818
export const billingSubscriptionSyncMockFns = {
1919
mockEnqueueCancelAtPeriodEndSync: vi.fn(async () => 'cancel-at-period-end-sync-event'),
20+
mockLockSubscriptionForSyncRetry: vi.fn(async () => undefined),
2021
mockRecommitSubscriptionSync: vi.fn(async () => undefined),
2122
mockRecordCancelAtPeriodEnd: vi.fn(async () => undefined),
2223
mockIsSubscriptionSyncEventType: vi.fn(
@@ -42,6 +43,7 @@ export const billingSubscriptionSyncMockFns = {
4243
*/
4344
export const billingSubscriptionSyncMock = {
4445
enqueueCancelAtPeriodEndSync: billingSubscriptionSyncMockFns.mockEnqueueCancelAtPeriodEndSync,
46+
lockSubscriptionForSyncRetry: billingSubscriptionSyncMockFns.mockLockSubscriptionForSyncRetry,
4547
recommitSubscriptionSync: billingSubscriptionSyncMockFns.mockRecommitSubscriptionSync,
4648
recordCancelAtPeriodEnd: billingSubscriptionSyncMockFns.mockRecordCancelAtPeriodEnd,
4749
isSubscriptionSyncEventType: billingSubscriptionSyncMockFns.mockIsSubscriptionSyncEventType,

0 commit comments

Comments
 (0)