Skip to content

Commit 886cfa2

Browse files
authored
fix(billing): converge Stripe subscription syncs and stop webhook echoes overwriting unsynced changes (#8764)
* fix(billing): converge Stripe cancel/contact syncs and guard webhook echoes over unsynced DB changes * fix(billing): record each committed sync value on every in-flight event so no stale intent can resurface * fix(billing): close remaining stale-intent paths (enterprise retry, legacy echoes, customer restore) * fix(billing): record Team activation's cleared cancellation under the row lock; drop the settled-event scan * fix(billing): lock the subscription before the outbox row in every sync retry path * fix(billing): push each sync's recorded value and reject Stripe reads older than a newer Sim commit * fix(billing): order Stripe-accepted reconcile values by when Stripe was read * fix(billing): recommit the in-flight value, include the plan in seat convergence, bound the reconcile's outbox reads * fix(billing): stamp every sync a later Stripe read confirms; ignore a moved cancel_at * fix(billing): a fresh idempotency key per seat write so a returning value is applied, not replayed * fix(billing): compare membership-driven seat and cancel changes against the committed value, not a possibly-stale row * fix(billing): a membership writer skips only when both the row and the committed value already match * fix(billing): a seat-row repair records no seat-change audit; deterministic latest-sync test helper * test(billing): latest-sync helper fails loudly on an enqueue-time tie
1 parent 6df74d6 commit 886cfa2

22 files changed

Lines changed: 2927 additions & 243 deletions

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

Lines changed: 46 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { db } from '@sim/db'
22
import { outboxEvent } from '@sim/db/schema'
33
import { createLogger } from '@sim/logger'
44
import { toError } from '@sim/utils/errors'
5+
import { toRecord } from '@sim/utils/object'
56
import { and, eq, sql } from 'drizzle-orm'
67
import { NextResponse } from 'next/server'
78
import { adminV1RequeueOutboxEventContract } from '@/lib/api/contracts/v1/admin'
@@ -11,6 +12,11 @@ import {
1112
ENTERPRISE_METADATA_SYNC_EVENT_TYPE,
1213
ENTERPRISE_PROVISION_EVENT_TYPE,
1314
} from '@/lib/billing/enterprise-outbox-events'
15+
import {
16+
isSubscriptionSyncEventType,
17+
lockSubscriptionForSyncRetry,
18+
recommitSubscriptionSync,
19+
} from '@/lib/billing/webhooks/subscription-sync'
1420
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
1521
import { withAdminAuthParams } from '@/app/api/v1/admin/middleware'
1622

@@ -28,7 +34,9 @@ const invalidOutboxEventResponse = (message: string) =>
2834
* will retry it. Resets `attempts`, `lastError`, and `availableAt` so
2935
* the next poll picks it up. Only dead-lettered events can be
3036
* requeued — completed/pending/processing rows are rejected to avoid
31-
* operator errors.
37+
* operator errors. A Stripe subscription sync is re-committed with the
38+
* subscription's current value, so the retry never revives the value it
39+
* failed with.
3240
*/
3341
export const POST = withRouteHandler(
3442
withAdminAuthParams<{ id: string }>(async (request, context) => {
@@ -61,23 +69,43 @@ export const POST = withRouteHandler(
6169
const deliveryRevision = metadataIntent?.success
6270
? metadataIntent.data.deliveryRevision + 1
6371
: null
64-
const result = await db
65-
.update(outboxEvent)
66-
.set({
67-
status: 'pending',
68-
attempts: 0,
69-
lastError: null,
70-
availableAt: new Date(),
71-
lockedAt: null,
72-
processedAt: null,
73-
...(deliveryRevision === null
74-
? {}
75-
: {
76-
payload: sql`(${outboxEvent.payload}::jsonb || ${JSON.stringify({ deliveryRevision })}::jsonb)::json`,
77-
}),
78-
})
79-
.where(and(eq(outboxEvent.id, id), eq(outboxEvent.status, 'dead_letter')))
80-
.returning({ id: outboxEvent.id, eventType: outboxEvent.eventType })
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
79+
const result = await db.transaction(async (tx) => {
80+
if (subscriptionSync) {
81+
await lockSubscriptionForSyncRetry(tx, subscriptionSync.subscriptionId)
82+
}
83+
const requeued = await tx
84+
.update(outboxEvent)
85+
.set({
86+
status: 'pending',
87+
attempts: 0,
88+
lastError: null,
89+
availableAt: new Date(),
90+
lockedAt: null,
91+
processedAt: null,
92+
...(deliveryRevision === null
93+
? {}
94+
: {
95+
payload: sql`(${outboxEvent.payload}::jsonb || ${JSON.stringify({ deliveryRevision })}::jsonb)::json`,
96+
}),
97+
})
98+
.where(and(eq(outboxEvent.id, id), eq(outboxEvent.status, 'dead_letter')))
99+
.returning({ id: outboxEvent.id, eventType: outboxEvent.eventType })
100+
if (subscriptionSync && requeued.length > 0) {
101+
await recommitSubscriptionSync(
102+
tx,
103+
subscriptionSync.eventType,
104+
subscriptionSync.subscriptionId
105+
)
106+
}
107+
return requeued
108+
})
81109

82110
if (result.length === 0) {
83111
return NextResponse.json(

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

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,7 @@ import {
3333
} from '@/lib/api/contracts/v1/admin'
3434
import { parseRequest } from '@/lib/api/server'
3535
import { requireStripeClient } from '@/lib/billing/stripe-client'
36-
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
37-
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
36+
import { enqueueCancelAtPeriodEndSync } from '@/lib/billing/webhooks/subscription-sync'
3837
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
3938
import { withAdminAuthParams } from '@/app/api/v1/admin/middleware'
4039
import {
@@ -113,15 +112,17 @@ export const DELETE = withRouteHandler(
113112
}
114113

115114
if (atPeriodEnd) {
115+
const stripeSubscriptionId = existing.stripeSubscriptionId
116116
await db.transaction(async (tx) => {
117117
await tx
118118
.update(subscription)
119119
.set({ cancelAtPeriodEnd: true })
120120
.where(eq(subscription.id, subscriptionId))
121121

122-
await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, {
123-
stripeSubscriptionId: existing.stripeSubscriptionId,
122+
await enqueueCancelAtPeriodEndSync(tx, {
123+
stripeSubscriptionId,
124124
subscriptionId: existing.id,
125+
cancelAtPeriodEnd: true,
125126
reason: reason ?? 'admin-cancel-at-period-end',
126127
})
127128
})

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { outboxEvent, subscription } from '@sim/db/schema'
22
import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing'
33
import { auditMock, auditMockFns } from '@sim/testing/mocks/audit.mock'
44
import { billingOutboxHandlersMock } from '@sim/testing/mocks/billing-outbox-handlers.mock'
5+
import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock'
56
import { organizationMembershipMock } from '@sim/testing/mocks/organization-membership.mock'
67
import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock'
78
import { stripeClientMock } from '@sim/testing/mocks/stripe.mock'
@@ -24,6 +25,7 @@ vi.mock('@/lib/billing/organizations/membership', () => organizationMembershipMo
2425
vi.mock('@/lib/billing/stripe-client', () => stripeClientMock)
2526
vi.mock('@/lib/billing/webhooks/outbox-handlers', () => billingOutboxHandlersMock)
2627
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)
28+
vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock)
2729

2830
import {
2931
refundDashboardSubscriptionPayment,
@@ -98,6 +100,7 @@ describe('admin subscription cancellation', () => {
98100

99101
it('requeues the same dead-lettered period-end cancellation operation', async () => {
100102
dbChainMockFns.returning.mockResolvedValueOnce([{ id: activeSubscription.id }])
103+
queueTableRows(outboxEvent, [{ subscriptionId: 'sub-row-1' }])
101104
queueTableRows(outboxEvent, [
102105
{
103106
id: 'outbox-1',
@@ -126,6 +129,7 @@ describe('admin subscription cancellation', () => {
126129
})
127130

128131
it('replays an immediate cancellation after the webhook removed active entitlement', async () => {
132+
queueTableRows(outboxEvent, [])
129133
queueTableRows(outboxEvent, [
130134
{
131135
id: 'outbox-1',
@@ -149,6 +153,7 @@ describe('admin subscription cancellation', () => {
149153
})
150154

151155
it('rejects reuse of a cancellation operation id with different timing', async () => {
156+
queueTableRows(outboxEvent, [])
152157
queueTableRows(outboxEvent, [
153158
{
154159
id: 'outbox-1',

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

Lines changed: 36 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,11 @@ import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/mem
88
import { requireStripeClient } from '@/lib/billing/stripe-client'
99
import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/utils'
1010
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
11+
import {
12+
enqueueCancelAtPeriodEndSync,
13+
lockSubscriptionForSyncRetry,
14+
recordCancelAtPeriodEnd,
15+
} from '@/lib/billing/webhooks/subscription-sync'
1116
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
1217

1318
const RECENT_INVOICE_LIMIT = 12
@@ -257,6 +262,24 @@ export async function requestDashboardSubscriptionCancellation({
257262
: 'admin-dashboard-cancel-at-period-end')
258263
const cancellation = await db.transaction(async (tx) => {
259264
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+
}
260283

261284
const [existingOperation] = await tx
262285
.select({
@@ -273,8 +296,7 @@ export async function requestDashboardSubscriptionCancellation({
273296
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
274297
OUTBOX_EVENT_TYPES.STRIPE_CANCEL_SUBSCRIPTION_IMMEDIATELY,
275298
]),
276-
sql`${outboxEvent.payload} ->> 'operationId' = ${operationId}`,
277-
sql`${outboxEvent.payload} ->> 'organizationId' = ${organizationId}`
299+
isThisOperation
278300
)
279301
)
280302
.for('update')
@@ -315,6 +337,9 @@ export async function requestDashboardSubscriptionCancellation({
315337
.where(
316338
and(eq(outboxEvent.id, existingOperation.id), eq(outboxEvent.status, 'dead_letter'))
317339
)
340+
if (existingOperation.eventType === OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END) {
341+
await recordCancelAtPeriodEnd(tx, existingOperation.subscriptionId, true)
342+
}
318343
return {
319344
operationId,
320345
outboxEventId: existingOperation.id,
@@ -384,18 +409,15 @@ export async function requestDashboardSubscriptionCancellation({
384409
.set({ cancelAtPeriodEnd: true })
385410
.where(eq(subscription.id, subscriptionRow.id))
386411
}
387-
const eventId = await enqueueOutboxEvent(
388-
tx,
389-
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
390-
{
391-
operationId,
392-
organizationId,
393-
subscriptionId: subscriptionRow.id,
394-
stripeSubscriptionId: subscriptionRow.stripeSubscriptionId,
395-
reason: normalizedReason,
396-
requestedBy: actor,
397-
}
398-
)
412+
const eventId = await enqueueCancelAtPeriodEndSync(tx, {
413+
operationId,
414+
organizationId,
415+
subscriptionId: subscriptionRow.id,
416+
stripeSubscriptionId: subscriptionRow.stripeSubscriptionId,
417+
cancelAtPeriodEnd: true,
418+
reason: normalizedReason,
419+
requestedBy: actor,
420+
})
399421
return {
400422
operationId,
401423
outboxEventId: eventId,

‎apps/sim/lib/auth/auth.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,10 @@ import {
103103
handleSubscriptionCreated,
104104
handleSubscriptionDeleted,
105105
} from '@/lib/billing/webhooks/subscription'
106+
import {
107+
reconcileSubscriptionSyncFromStripe,
108+
recordCustomerRestoreAfterHook,
109+
} from '@/lib/billing/webhooks/subscription-sync'
106110
import { handleSubscriptionUsageUpdate } from '@/lib/billing/webhooks/subscription-usage'
107111
import { env } from '@/lib/core/config/env'
108112
import {
@@ -1100,6 +1104,8 @@ export const auth = betterAuth({
11001104
return
11011105
}),
11021106
after: createAuthMiddleware(async (ctx) => {
1107+
if (isBillingEnabled) await recordCustomerRestoreAfterHook(ctx)
1108+
11031109
if (isBillingEnabled && ctx.path === '/subscription/upgrade') {
11041110
const checkoutContext = ctx as typeof ctx & {
11051111
billingCheckoutAdmissionClaim?: CheckoutAdmissionClaim
@@ -1727,6 +1733,7 @@ export const auth = betterAuth({
17271733
case 'customer.subscription.created':
17281734
case 'customer.subscription.updated': {
17291735
await handleManualEnterpriseSubscription(event)
1736+
await reconcileSubscriptionSyncFromStripe(event)
17301737
await handleSubscriptionUsageUpdate(event)
17311738
break
17321739
}

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

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -77,6 +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 {
81+
lockSubscriptionForSyncRetry,
82+
recommitSubscriptionSync,
83+
} from '@/lib/billing/webhooks/subscription-sync'
8084
import { env } from '@/lib/core/config/env'
8185
import {
8286
continueOutboxHandler,
@@ -2103,6 +2107,9 @@ export async function retryEnterpriseFollowUpJob(
21032107

21042108
const retried = await db.transaction(async (tx) => {
21052109
await acquireOrganizationMutationLock(tx, operationPayload.request.organizationId)
2110+
if (snapshotDetail.kind === 'personal_subscription_cancellation') {
2111+
await lockSubscriptionForSyncRetry(tx, snapshotDetail.subjectId)
2112+
}
21062113
const [row] = await tx
21072114
.select({
21082115
status: outboxEvent.status,
@@ -2119,7 +2126,9 @@ export async function retryEnterpriseFollowUpJob(
21192126
!detail ||
21202127
!getEnterpriseFollowUpOperationIds(row.eventType, row.payload).includes(operationId) ||
21212128
(detail.kind === 'member_reconciliation' &&
2122-
detail.subjectId !== operationPayload.request.organizationId)
2129+
detail.subjectId !== operationPayload.request.organizationId) ||
2130+
detail.kind !== snapshotDetail.kind ||
2131+
detail.subjectId !== snapshotDetail.subjectId
21232132
) {
21242133
throw new EnterpriseProvisioningError('Enterprise follow-up job not found')
21252134
}
@@ -2135,6 +2144,13 @@ export async function retryEnterpriseFollowUpJob(
21352144
processedAt: null,
21362145
})
21372146
.where(eq(outboxEvent.id, jobEventId))
2147+
if (detail.kind === 'personal_subscription_cancellation') {
2148+
await recommitSubscriptionSync(
2149+
tx,
2150+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
2151+
detail.subjectId
2152+
)
2153+
}
21382154
return true
21392155
})
21402156

‎apps/sim/lib/billing/organizations/lock-order.test.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ import {
1717
workspace,
1818
} from '@sim/db/schema'
1919
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
20+
import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock'
2021
import { outboxServiceMock } from '@sim/testing/mocks/outbox-service.mock'
2122
import { beforeEach, describe, expect, it, vi } from 'vitest'
2223

@@ -42,6 +43,7 @@ import type { DbOrTx } from '@/lib/db/types'
4243
import { attachOwnedWorkspacesToOrganizationTx } from '@/lib/workspaces/organization-workspaces'
4344

4445
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)
46+
vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock)
4547

4648
/**
4749
* A superset row that satisfies every read in the join path: a paid org sub, a

0 commit comments

Comments
 (0)