Skip to content

Commit 40ec5a8

Browse files
committed
fix(billing): record Team activation's cleared cancellation under the row lock; drop the settled-event scan
1 parent 95194c8 commit 40ec5a8

7 files changed

Lines changed: 284 additions & 119 deletions

File tree

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

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -58,18 +58,24 @@ const mockGetHighestPriorityPersonalSubscription =
5858
const { mockEnqueueSubscriptionSeatsSync, mockEnqueueCancelAtPeriodEndSync } =
5959
billingSubscriptionSyncMockFns
6060

61-
function testExecutor(onUpdate: () => void = () => {}) {
61+
/** The subscription row as the activation re-reads it under its lock. */
62+
function testExecutor(onSubscriptionLock: () => void = () => {}) {
63+
const lockedRow = { cancelAtPeriodEnd: false, seats: 1 }
6264
return {
65+
select: () => ({
66+
from: () => ({
67+
where: () => ({
68+
for: () => {
69+
onSubscriptionLock()
70+
return { limit: () => Promise.resolve([lockedRow]) }
71+
},
72+
}),
73+
}),
74+
}),
6375
update: () => ({
6476
set: (values: Record<string, unknown>) => {
65-
onUpdate()
6677
updateCalls.value.push(values)
67-
return {
68-
where: () =>
69-
Object.assign(Promise.resolve([]), {
70-
returning: () => Promise.resolve([{ seats: 1 }]),
71-
}),
72-
}
78+
return { where: () => Promise.resolve([]) }
7379
},
7480
}),
7581
} as never

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

Lines changed: 35 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils'
1919
import {
2020
enqueueCancelAtPeriodEndSync,
2121
enqueueSubscriptionSeatsSync,
22+
recordCancelAtPeriodEnd,
2223
} from '@/lib/billing/webhooks/subscription-sync'
2324
import type { DbOrTx, DbTransaction } from '@/lib/db/types'
2425

@@ -238,43 +239,49 @@ async function convertPersonalSubscriptionToTeam(
238239
* the post-join seat reconcile is skipped or fails. Any scheduled cancellation
239240
* is cleared (DB + Stripe) so a freshly-activated Team is not left scheduled to
240241
* cancel, including the legacy personal-scoped Team case where the plan is
241-
* unchanged.
242+
* unchanged. The row is read under its lock, so a cancellation committed after
243+
* the caller's earlier read is still cleared and recorded.
242244
*/
243245
async function activateTeamSubscription(
244-
sub: { id: string; cancelAtPeriodEnd?: boolean | null; stripeSubscriptionId: string | null },
246+
sub: { id: string; stripeSubscriptionId: string | null },
245247
targetPlan: string,
246248
{ planChanged }: { planChanged: boolean },
247-
executor: DbOrTx
249+
tx: DbOrTx
248250
): Promise<void> {
249-
const shouldClearCancellation =
250-
Boolean(sub.cancelAtPeriodEnd) && Boolean(sub.stripeSubscriptionId)
251-
252-
const apply = async (tx: DbOrTx) => {
253-
const [activated] = await tx
254-
.update(subscriptionTable)
255-
.set({ plan: targetPlan, cancelAtPeriodEnd: false })
256-
.where(eq(subscriptionTable.id, sub.id))
257-
.returning({ seats: subscriptionTable.seats })
251+
const [locked] = await tx
252+
.select({
253+
cancelAtPeriodEnd: subscriptionTable.cancelAtPeriodEnd,
254+
seats: subscriptionTable.seats,
255+
})
256+
.from(subscriptionTable)
257+
.where(eq(subscriptionTable.id, sub.id))
258+
.for('update')
259+
.limit(1)
258260

259-
if (planChanged) {
260-
await enqueueSubscriptionSeatsSync(tx, {
261-
subscriptionId: sub.id,
262-
seats: activated?.seats ?? 1,
263-
reason: 'pro-to-team-conversion',
264-
})
265-
}
261+
await tx
262+
.update(subscriptionTable)
263+
.set({ plan: targetPlan, cancelAtPeriodEnd: false })
264+
.where(eq(subscriptionTable.id, sub.id))
266265

267-
if (shouldClearCancellation) {
268-
await enqueueCancelAtPeriodEndSync(tx, {
269-
stripeSubscriptionId: sub.stripeSubscriptionId as string,
270-
subscriptionId: sub.id,
271-
cancelAtPeriodEnd: false,
272-
reason: 'pro-to-team-conversion',
273-
})
274-
}
266+
if (planChanged) {
267+
await enqueueSubscriptionSeatsSync(tx, {
268+
subscriptionId: sub.id,
269+
seats: locked?.seats ?? 1,
270+
reason: 'pro-to-team-conversion',
271+
})
275272
}
276273

277-
await apply(executor)
274+
if (!sub.stripeSubscriptionId) return
275+
if (locked?.cancelAtPeriodEnd) {
276+
await enqueueCancelAtPeriodEndSync(tx, {
277+
stripeSubscriptionId: sub.stripeSubscriptionId,
278+
subscriptionId: sub.id,
279+
cancelAtPeriodEnd: false,
280+
reason: 'pro-to-team-conversion',
281+
})
282+
} else {
283+
await recordCancelAtPeriodEnd(tx, sub.id, false)
284+
}
278285
}
279286

280287
/**

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

Lines changed: 152 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,13 +48,15 @@ import {
4848
pauseProSubscriptionForOrgCoverage,
4949
restoreUserProSubscription,
5050
} from '@/lib/billing/organizations/membership'
51+
import { ensureTeamOrganizationForAcceptance } from '@/lib/billing/organizations/provision-seat'
5152
import { reconcileOrganizationSeats } from '@/lib/billing/organizations/seats'
5253
import { isTeam } from '@/lib/billing/plan-helpers'
5354
import { syncSeatsFromStripeQuantity } from '@/lib/billing/validation/seat-management'
5455
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
5556
import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers'
5657
import {
5758
commitCustomerRestoredSubscription,
59+
enqueueCancelAtPeriodEndSync,
5860
reconcileSubscriptionSyncFromStripe,
5961
} from '@/lib/billing/webhooks/subscription-sync'
6062
import { enqueueOutboxEvent, processOutboxEventById } from '@/lib/core/outbox/service'
@@ -122,7 +124,15 @@ let deliver: ReturnType<typeof createWebhookEndpoint>
122124

123125
beforeAll(async () => {
124126
await connection`CREATE SCHEMA ${connection(schemaName)}`
125-
for (const table of ['subscription', 'outbox_event', 'member', 'user', 'organization']) {
127+
for (const table of [
128+
'subscription',
129+
'outbox_event',
130+
'member',
131+
'user',
132+
'organization',
133+
'workspace',
134+
'permissions',
135+
]) {
126136
await connection.unsafe(`CREATE TABLE "${table}" (LIKE public."${table}" INCLUDING ALL)`)
127137
}
128138
database.current = testDatabase
@@ -607,6 +617,77 @@ describe('cancel_at_period_end sync', () => {
607617
})
608618
})
609619

620+
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+
633+
it('records the cleared cancellation when a cancel is committed while it activates Team', async () => {
634+
const owner = await createUser('owner')
635+
const subscriptionId = generateId()
636+
const stripeSubscriptionId = `sub_${subscriptionId}`
637+
await testDatabase.insert(subscription).values({
638+
id: subscriptionId,
639+
plan: 'team',
640+
referenceId: owner.id,
641+
status: 'active',
642+
seats: 1,
643+
stripeSubscriptionId,
644+
stripeCustomerId: `cus_${subscriptionId}`,
645+
cancelAtPeriodEnd: false,
646+
})
647+
stripe.addSubscription({ id: stripeSubscriptionId, customer: `cus_${subscriptionId}` })
648+
649+
let releaseCancel: () => void = () => {}
650+
const cancelHeld = new Promise<void>((resolve) => {
651+
releaseCancel = resolve
652+
})
653+
const cancelling = testDatabase.transaction(async (tx) => {
654+
await tx
655+
.update(subscription)
656+
.set({ cancelAtPeriodEnd: true })
657+
.where(eq(subscription.id, subscriptionId))
658+
await enqueueCancelAtPeriodEndSync(tx, {
659+
stripeSubscriptionId,
660+
subscriptionId,
661+
cancelAtPeriodEnd: true,
662+
reason: 'admin-cancel-at-period-end',
663+
})
664+
await cancelHeld
665+
})
666+
const activating = testDatabase.transaction((tx) =>
667+
ensureTeamOrganizationForAcceptance({
668+
billingOwnerUserId: owner.id,
669+
workspaceOrganizationId: null,
670+
executor: tx,
671+
workspaceIdsToAttach: [],
672+
})
673+
)
674+
await untilAnotherTransactionWaitsOnALock()
675+
releaseCancel()
676+
await cancelling
677+
await expect(activating).resolves.toMatchObject({ success: true })
678+
expect((await storedSubscription(subscriptionId)).cancelAtPeriodEnd).toBe(false)
679+
680+
await deliverUnrelatedUpdate(stripeSubscriptionId)
681+
expect((await storedSubscription(subscriptionId)).cancelAtPeriodEnd).toBe(false)
682+
const cancelSync = await latestOutboxEventId(
683+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
684+
subscriptionId
685+
)
686+
await expect(processEvent(cancelSync)).resolves.toBe('completed')
687+
expect(stripe.subscription(stripeSubscriptionId).cancel_at_period_end).toBe(false)
688+
})
689+
})
690+
610691
describe('customer contact sync', () => {
611692
it('pushes the current owner when an earlier sync lands in Stripe after a newer one', async () => {
612693
const [first, second, third] = await Promise.all([
@@ -704,3 +785,73 @@ describe('Team seat sync', () => {
704785
expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(1)
705786
})
706787
})
788+
789+
describe('webhook reconcile cost', () => {
790+
interface QueryPlan {
791+
'Node Type': string
792+
'Relation Name'?: string
793+
'Shared Hit Blocks': number
794+
'Shared Read Blocks': number
795+
Plans?: QueryPlan[]
796+
}
797+
const planNodes = (plan: QueryPlan): QueryPlan[] => [
798+
plan,
799+
...(plan.Plans ?? []).flatMap(planNodes),
800+
]
801+
802+
it('reads only the syncs that can still run, however many have completed', async () => {
803+
const pro = await createProUserInPaidOrganization()
804+
await pauseProSubscriptionForOrgCoverage(pro.userId)
805+
await connection`
806+
INSERT INTO outbox_event (id, event_type, payload, status, available_at, created_at, processed_at)
807+
SELECT ${generateId()} || ':' || n, ${OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END},
808+
json_build_object(
809+
'subscriptionId', ${pro.subscriptionId}::text,
810+
'cancelAtPeriodEnd', n % 2 = 0,
811+
'committedAt', n
812+
),
813+
'completed', now(), now(), now()
814+
FROM generate_series(1, 20000) AS n`
815+
await connection`ANALYZE outbox_event`
816+
817+
const issued: { query: string; parameters: unknown[] }[] = []
818+
const traced = postgres(
819+
readTestDatabaseUrl(),
820+
withUtcTimestamps({
821+
max: 2,
822+
prepare: false,
823+
fetch_types: false,
824+
connection: { search_path: schemaName },
825+
onnotice: () => {},
826+
debug: (_connection: number, query: string, parameters: unknown[]) => {
827+
issued.push({ query, parameters })
828+
},
829+
})
830+
)
831+
database.current = drizzle(traced, { schema })
832+
try {
833+
await deliverUnrelatedUpdate(pro.stripeSubscriptionId)
834+
} finally {
835+
database.current = testDatabase
836+
await traced.end()
837+
}
838+
expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true)
839+
840+
const outboxReads = issued.filter(({ query }) => /from "outbox_event"/i.test(query))
841+
expect(outboxReads.length).toBeGreaterThan(0)
842+
for (const { query, parameters } of outboxReads) {
843+
const [explained] = await connection.unsafe(
844+
`EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${query}`,
845+
parameters as never[]
846+
)
847+
const plan = (explained['QUERY PLAN'] as { Plan: QueryPlan }[])[0].Plan
848+
const nodes = planNodes(plan)
849+
expect(
850+
nodes.some(
851+
(node) => node['Node Type'] === 'Seq Scan' && node['Relation Name'] === 'outbox_event'
852+
)
853+
).toBe(false)
854+
expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(100)
855+
}
856+
})
857+
})

0 commit comments

Comments
 (0)