Skip to content

Commit fa5e995

Browse files
committed
fix(billing): compare membership-driven seat and cancel changes against the committed value, not a possibly-stale row
1 parent dc34ffa commit fa5e995

6 files changed

Lines changed: 128 additions & 15 deletions

File tree

‎apps/sim/lib/billing/organizations/membership.ts‎

Lines changed: 28 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,10 @@ import {
4848
import { toDecimal, toNumber } from '@/lib/billing/utils/decimal'
4949
import { validateSeatAvailability } from '@/lib/billing/validation/seat-management'
5050
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
51-
import { enqueueCancelAtPeriodEndSync } from '@/lib/billing/webhooks/subscription-sync'
51+
import {
52+
enqueueCancelAtPeriodEndSync,
53+
readCommittedCancelAtPeriodEnd,
54+
} from '@/lib/billing/webhooks/subscription-sync'
5255
import { isBillingEnabled } from '@/lib/core/config/env-flags'
5356
import { OrchestrationError } from '@/lib/core/orchestration/types'
5457
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
@@ -259,7 +262,13 @@ export async function restoreUserProSubscription(userId: string): Promise<Restor
259262
.for('update')
260263
.limit(1)
261264

262-
if (!personalPro?.cancelAtPeriodEnd || !personalPro.stripeSubscriptionId) return
265+
if (!personalPro?.stripeSubscriptionId) return
266+
const pausing = await readCommittedCancelAtPeriodEnd(
267+
tx,
268+
personalPro.id,
269+
Boolean(personalPro.cancelAtPeriodEnd)
270+
)
271+
if (!pausing) return
263272
result.subscriptionId = personalPro.id
264273

265274
const organizationMemberships = await tx
@@ -405,7 +414,15 @@ export async function pauseProSubscriptionForOrgCoverage(
405414

406415
result.subscriptionId = personalPro.id
407416

408-
if (personalPro.cancelAtPeriodEnd) return
417+
if (
418+
await readCommittedCancelAtPeriodEnd(
419+
tx,
420+
personalPro.id,
421+
Boolean(personalPro.cancelAtPeriodEnd)
422+
)
423+
) {
424+
return
425+
}
409426

410427
await tx
411428
.update(subscriptionTable)
@@ -855,7 +872,14 @@ async function applyPaidOrgJoinBillingTx(
855872
.for('update')
856873
.limit(1)
857874

858-
if (personalPro && !personalPro.cancelAtPeriodEnd) {
875+
const alreadyPausing =
876+
personalPro &&
877+
(await readCommittedCancelAtPeriodEnd(
878+
tx,
879+
personalPro.id,
880+
Boolean(personalPro.cancelAtPeriodEnd)
881+
))
882+
if (personalPro && !alreadyPausing) {
859883
await tx
860884
.update(subscriptionTable)
861885
.set({ cancelAtPeriodEnd: true })

‎apps/sim/lib/billing/organizations/provision-seat.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils'
1919
import {
2020
enqueueCancelAtPeriodEndSync,
2121
enqueueSubscriptionSeatsSync,
22+
readCommittedCancelAtPeriodEnd,
23+
readCommittedSeats,
2224
recordCancelAtPeriodEnd,
2325
} from '@/lib/billing/webhooks/subscription-sync'
2426
import type { DbOrTx, DbTransaction } from '@/lib/db/types'
@@ -267,13 +269,13 @@ async function activateTeamSubscription(
267269
if (planChanged) {
268270
await enqueueSubscriptionSeatsSync(tx, {
269271
subscriptionId: sub.id,
270-
seats: locked?.seats ?? 1,
272+
seats: await readCommittedSeats(tx, sub.id, locked?.seats ?? 1),
271273
reason: 'pro-to-team-conversion',
272274
})
273275
}
274276

275277
if (!locked?.stripeSubscriptionId) return
276-
if (locked.cancelAtPeriodEnd) {
278+
if (await readCommittedCancelAtPeriodEnd(tx, sub.id, Boolean(locked.cancelAtPeriodEnd))) {
277279
await enqueueCancelAtPeriodEndSync(tx, {
278280
stripeSubscriptionId: locked.stripeSubscriptionId,
279281
subscriptionId: sub.id,

‎apps/sim/lib/billing/organizations/seats.ts‎

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,10 @@ import { and, count, desc, eq, inArray } from 'drizzle-orm'
66
import { syncSubscriptionUsageLimits } from '@/lib/billing/organization'
77
import { isTeam } from '@/lib/billing/plan-helpers'
88
import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/utils'
9-
import { enqueueSubscriptionSeatsSync } from '@/lib/billing/webhooks/subscription-sync'
9+
import {
10+
enqueueSubscriptionSeatsSync,
11+
readCommittedSeats,
12+
} from '@/lib/billing/webhooks/subscription-sync'
1013
import { isBillingEnabled } from '@/lib/core/config/env-flags'
1114
import { captureServerEvent } from '@/lib/posthog/server'
1215

@@ -100,7 +103,11 @@ export async function reconcileOrganizationSeats({
100103
.where(eq(member.organizationId, organizationId))
101104

102105
const targetSeats = Math.max(1, memberCountRow?.value ?? 1)
103-
const currentSeats = orgSubscription.seats ?? 1
106+
const currentSeats = await readCommittedSeats(
107+
tx,
108+
orgSubscription.id,
109+
orgSubscription.seats ?? 1
110+
)
104111

105112
if (targetSeats === currentSeats) {
106113
return { kind: 'noop', seats: currentSeats }

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

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -853,6 +853,26 @@ describe('cancel_at_period_end sync', () => {
853853
expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true)
854854
})
855855

856+
it('restores the personal Pro when its member leaves while the plugin has overwritten the row', async () => {
857+
const pro = await createProUserInPaidOrganization()
858+
await pauseProSubscriptionForOrgCoverage(pro.userId)
859+
const pauseSync = await latestOutboxEventId(
860+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
861+
pro.subscriptionId
862+
)
863+
864+
beforeReconcile = async () => {
865+
beforeReconcile = undefined
866+
await leaveOrganization(pro.userId, pro.paidOrganization.organizationId)
867+
await restoreUserProSubscription(pro.userId)
868+
}
869+
await deliverUnrelatedUpdate(pro.stripeSubscriptionId)
870+
871+
expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false)
872+
await expect(processEvent(pauseSync)).resolves.toBe('completed')
873+
expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false)
874+
})
875+
856876
it('does not revive an older value when a retry path resets its sync without re-committing', async () => {
857877
const pro = await createProUserInPaidOrganization()
858878
await pauseProSubscriptionForOrgCoverage(pro.userId)
@@ -1237,6 +1257,32 @@ describe('Team seat sync', () => {
12371257
expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(2)
12381258
})
12391259

1260+
it('drops a seat when a member leaves while the plugin has overwritten the row', async () => {
1261+
const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')])
1262+
const org = await createOrganizationWithPlan('team', 1)
1263+
await addMember(org.organizationId, owner.id, 'owner')
1264+
await addMember(org.organizationId, joiner.id)
1265+
await reconcileOrganizationSeats({ organizationId: org.organizationId, reason: 'member-added' })
1266+
const growSync = await latestOutboxEventId(
1267+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS,
1268+
org.subscriptionId
1269+
)
1270+
1271+
beforeReconcile = async () => {
1272+
beforeReconcile = undefined
1273+
await leaveOrganization(joiner.id, org.organizationId)
1274+
await reconcileOrganizationSeats({
1275+
organizationId: org.organizationId,
1276+
reason: 'member-removed',
1277+
})
1278+
}
1279+
await deliverUnrelatedUpdate(org.stripeSubscriptionId)
1280+
1281+
expect((await storedSubscription(org.subscriptionId)).seats).toBe(1)
1282+
await expect(processEvent(growSync)).resolves.toBe('completed')
1283+
expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(1)
1284+
})
1285+
12401286
it('does not revive an older seat count when its dead-lettered sync is requeued', async () => {
12411287
const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')])
12421288
const org = await createOrganizationWithPlan('team', 1)

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

Lines changed: 31 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -215,18 +215,18 @@ export async function recommitSubscriptionSync(
215215
.limit(1)
216216
if (!current) return
217217

218-
const pending = await readSyncIntents(tx, subscriptionId)
219218
if (eventType === CANCEL_SYNC) {
220-
const intent = pending.cancelAtPeriodEnd
221219
await commitIntent(tx, eventType, subscriptionId, {
222-
cancelAtPeriodEnd:
223-
intent.status === 'value' ? intent.value : Boolean(current.cancelAtPeriodEnd),
220+
cancelAtPeriodEnd: await readCommittedCancelAtPeriodEnd(
221+
tx,
222+
subscriptionId,
223+
Boolean(current.cancelAtPeriodEnd)
224+
),
224225
})
225226
return
226227
}
227-
const intent = pending.seats
228228
await commitIntent(tx, eventType, subscriptionId, {
229-
seats: intent.status === 'value' ? intent.value : (current.seats ?? 1),
229+
seats: await readCommittedSeats(tx, subscriptionId, current.seats ?? 1),
230230
})
231231
}
232232

@@ -304,6 +304,31 @@ export async function readRecordedSyncValue(
304304
}
305305
}
306306

307+
/**
308+
* The subscription's latest committed `cancelAtPeriodEnd`: the newest value an in-flight sync
309+
* records, else `stored` (the row). Until the reconcile step runs, the row can hold the Stripe
310+
* plugin's stale webhook payload, so a writer deciding whether a change is needed compares
311+
* against this, never the row alone. The caller holds the subscription row lock.
312+
*/
313+
export async function readCommittedCancelAtPeriodEnd(
314+
tx: DbOrTx,
315+
subscriptionId: string,
316+
stored: boolean
317+
): Promise<boolean> {
318+
const intent = (await readSyncIntents(tx, subscriptionId)).cancelAtPeriodEnd
319+
return intent.status === 'value' ? intent.value : stored
320+
}
321+
322+
/** The seat-count counterpart of {@link readCommittedCancelAtPeriodEnd}. */
323+
export async function readCommittedSeats(
324+
tx: DbOrTx,
325+
subscriptionId: string,
326+
stored: number
327+
): Promise<number> {
328+
const intent = (await readSyncIntents(tx, subscriptionId)).seats
329+
return intent.status === 'value' ? intent.value : stored
330+
}
331+
307332
/** A fresh key per Stripe write: the SDK reuses it across its own network retries of that call. */
308333
export function cancelAtPeriodEndSyncIdempotencyKey(eventId: string): string {
309334
return `${CANCEL_AT_PERIOD_END_SYNC_KEY_PREFIX}${eventId}:${generateShortId()}`

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

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,8 @@ import { vi } from 'vitest'
33
/**
44
* Controllable mock functions for `@/lib/billing/webhooks/subscription-sync`. The enqueue
55
* functions resolve to a fixed event id; drive them with `mockResolvedValueOnce`.
6-
* `mockIsSubscriptionSyncEventType` keeps the real logic.
6+
* `mockIsSubscriptionSyncEventType` keeps the real logic; the `mockReadCommitted*` readers
7+
* return the stored value they are given, as when nothing is in flight.
78
*
89
* @example
910
* ```ts
@@ -32,6 +33,12 @@ export const billingSubscriptionSyncMockFns = {
3233
mockReconcileSubscriptionSyncFromStripe: vi.fn(async () => undefined),
3334
mockRecordCustomerRestoreAfterHook: vi.fn(async () => undefined),
3435
mockReadRecordedSyncValue: vi.fn(async () => undefined),
36+
mockReadCommittedCancelAtPeriodEnd: vi.fn(
37+
async (_tx: unknown, _subscriptionId: string, stored: boolean) => stored
38+
),
39+
mockReadCommittedSeats: vi.fn(
40+
async (_tx: unknown, _subscriptionId: string, stored: number) => stored
41+
),
3542
}
3643

3744
/**
@@ -55,4 +62,6 @@ export const billingSubscriptionSyncMock = {
5562
billingSubscriptionSyncMockFns.mockReconcileSubscriptionSyncFromStripe,
5663
recordCustomerRestoreAfterHook: billingSubscriptionSyncMockFns.mockRecordCustomerRestoreAfterHook,
5764
readRecordedSyncValue: billingSubscriptionSyncMockFns.mockReadRecordedSyncValue,
65+
readCommittedCancelAtPeriodEnd: billingSubscriptionSyncMockFns.mockReadCommittedCancelAtPeriodEnd,
66+
readCommittedSeats: billingSubscriptionSyncMockFns.mockReadCommittedSeats,
5867
}

0 commit comments

Comments
 (0)