Skip to content

Commit 2bd650a

Browse files
committed
fix(billing): stamp every sync a later Stripe read confirms; ignore a moved cancel_at
1 parent 7515202 commit 2bd650a

3 files changed

Lines changed: 84 additions & 18 deletions

File tree

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

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -519,6 +519,66 @@ describe('cancel_at_period_end sync', () => {
519519
expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false)
520520
})
521521

522+
it('records a later Stripe read even when every sync already carries its value', async () => {
523+
const pro = await createProUserInPaidOrganization()
524+
await pauseProSubscriptionForOrgCoverage(pro.userId)
525+
const pauseSync = await latestOutboxEventId(
526+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
527+
pro.subscriptionId
528+
)
529+
stripe.failNextUpdateAfterApplying('subscriptions')
530+
await expect(processEvent(pauseSync)).resolves.toBe('pending')
531+
532+
stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false })
533+
const restoreRead = stripe.holdNextRequest('subscriptions.retrieve')
534+
const reconcilingRestore = deliver(stripe.events.at(-1) as Stripe.Event)
535+
await restoreRead.reached
536+
537+
stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true })
538+
const cancelRead = stripe.holdNextRequest('subscriptions.retrieve')
539+
const reconcilingCancel = deliver(stripe.events.at(-1) as Stripe.Event)
540+
await cancelRead.reached
541+
542+
cancelRead.release()
543+
await reconcilingCancel
544+
restoreRead.release()
545+
await reconcilingRestore
546+
547+
expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true)
548+
await makeDue(pauseSync)
549+
await expect(processEvent(pauseSync)).resolves.toBe('completed')
550+
expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true)
551+
})
552+
553+
it('keeps a pending value when an unrelated Stripe update moves the cancellation date', async () => {
554+
const pro = await createProUserInPaidOrganization()
555+
await pauseProSubscriptionForOrgCoverage(pro.userId)
556+
const pauseSync = await latestOutboxEventId(
557+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
558+
pro.subscriptionId
559+
)
560+
await expect(processEvent(pauseSync)).resolves.toBe('completed')
561+
const periodEnd = Math.floor(Date.now() / 1000) + 30 * 24 * 60 * 60
562+
stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: periodEnd })
563+
await deliver(stripe.events.at(-1) as Stripe.Event)
564+
565+
await leaveOrganization(pro.userId, pro.paidOrganization.organizationId)
566+
await restoreUserProSubscription(pro.userId)
567+
const restoreSync = await latestOutboxEventId(
568+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END,
569+
pro.subscriptionId
570+
)
571+
572+
stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: periodEnd + 335 * 24 * 60 * 60 })
573+
const moved = stripe.events.at(-1) as Stripe.Event
574+
expect(Object.keys(moved.data.previous_attributes ?? {})).toEqual(['cancel_at'])
575+
await deliver(moved)
576+
577+
expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false)
578+
await expect(processEvent(restoreSync)).resolves.toBe('completed')
579+
expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false)
580+
})
581+
522582
it('lets a change made in Stripe while a sync is pending win over the pending value', async () => {
523583
const pro = await createProUserInPaidOrganization()
524584
await pauseProSubscriptionForOrgCoverage(pro.userId)

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

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -101,8 +101,9 @@ async function readDatabaseClock(executor: DbOrTx): Promise<number> {
101101
*
102102
* A Sim commit is stamped with the clock under that lock and rewrites every such event. A value
103103
* taken from Stripe passes `observedAt`, the clock read just before Stripe was read, so it orders
104-
* by when it was observed (a slower reconcile of an earlier Stripe read cannot outrank a later
105-
* one), and rewrites only the events that do not already carry it.
104+
* by when it was observed: it rewrites (value and stamp) only the events last stamped before that
105+
* observation, including ones already holding the value, so a slower reconcile of an earlier
106+
* Stripe read cannot outrank a later one and never overwrites a newer commit.
106107
*/
107108
async function commitIntent<T extends SyncIntentFields>(
108109
tx: DbOrTx,
@@ -120,7 +121,7 @@ async function commitIntent<T extends SyncIntentFields>(
120121
eventType,
121122
subscriptionSubject(subscriptionId),
122123
committed,
123-
observedAt === undefined ? undefined : fields
124+
observedAt === undefined ? undefined : 'committedAt'
124125
)
125126
return committed
126127
}
@@ -358,15 +359,18 @@ async function readSyncIntents(executor: DbOrTx, subscriptionId: string) {
358359
}
359360

360361
/**
361-
* True when the event records a cancellation change made in Stripe, not by Sim's sync. A change
362-
* to `cancel_at` counts too: Better Auth's restore clears `cancel_at` when it is set, and Stripe
363-
* may then list only `cancel_at` among the previous attributes.
362+
* True when the event records a cancellation change made in Stripe, not by Sim's sync. A
363+
* `cancel_at` that was set or cleared counts too: Better Auth's restore clears `cancel_at` when it
364+
* is set, and Stripe may then list only `cancel_at` among the previous attributes. A `cancel_at`
365+
* that only moved (e.g. a billing-interval switch on a subscription already ending) is not a
366+
* cancellation change.
364367
*/
365368
function isCancellationChangedInStripe(event: Stripe.Event): boolean {
366369
const previousAttributes = toRecord(event.data.previous_attributes)
367-
if (!('cancel_at_period_end' in previousAttributes) && !('cancel_at' in previousAttributes)) {
368-
return false
369-
}
370+
const scheduledOrCleared =
371+
'cancel_at' in previousAttributes &&
372+
(previousAttributes.cancel_at == null) !== (toRecord(event.data.object).cancel_at == null)
373+
if (!('cancel_at_period_end' in previousAttributes) && !scheduledOrCleared) return false
370374
const idempotencyKey = event.request?.idempotency_key
371375
const issuedBySimSync =
372376
idempotencyKey?.startsWith(CANCEL_AT_PERIOD_END_SYNC_KEY_PREFIX) ||
@@ -410,9 +414,10 @@ function cancelAtPeriodEndSource(
410414
* - `cancelAtPeriodEnd`: while a cancel sync is in flight, its committed value wins over
411415
* snapshots and over echoes of Sim's own writes. A change made in Stripe itself (customer
412416
* portal, dashboard, Better Auth's cancel/restore endpoints), recognised by a non-Sim request
413-
* changing `cancel_at_period_end` or `cancel_at`, wins and is committed onto every
414-
* sync that can still run, unless Sim committed a newer value after Stripe was read. With no
415-
* sync in flight Stripe wins, read live so out-of-order delivery cannot regress it.
417+
* changing `cancel_at_period_end` or setting or clearing `cancel_at`, wins and is committed
418+
* onto every sync that can still run, unless Sim committed a newer value after Stripe was
419+
* read. With no sync in flight Stripe wins, read live so out-of-order delivery cannot regress
420+
* it.
416421
* - Precedence across the two systems is arrival order, not wall-clock order: a Stripe-side
417422
* change whose webhook is processed after a Sim commit wins even if the customer made it
418423
* earlier. Stripe's `event.created` is not compared with the database clock, because skew

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

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -461,16 +461,17 @@ export async function listInflightOutboxEvents(
461461

462462
/**
463463
* Shallow-merges `patch` into the payload of every `pending`, `processing`, or `dead_letter`
464-
* event of the type for one subject, skipping events whose payload already contains
465-
* `unlessPayloadContains`. One UPDATE; nothing is read into memory. Callers serialize writers
466-
* for the subject with their domain lock.
464+
* event of the type for one subject. With `onlyIfOlderThanPatch`, naming a numeric payload key
465+
* that `patch` sets, an event whose own value for that key is already at least the patch's is
466+
* left alone, so a stale writer never overwrites a newer one. One UPDATE; nothing is read into
467+
* memory. Callers serialize writers for the subject with their domain lock.
467468
*/
468469
export async function patchRetryableOutboxEvents(
469470
executor: Pick<typeof db, 'update'>,
470471
eventType: string,
471472
subject: OutboxPayloadSubject,
472473
patch: Record<string, unknown>,
473-
unlessPayloadContains?: Record<string, unknown>
474+
onlyIfOlderThanPatch?: string
474475
): Promise<number> {
475476
const patched = await executor
476477
.update(outboxEvent)
@@ -480,8 +481,8 @@ export async function patchRetryableOutboxEvents(
480481
.where(
481482
and(
482483
eventsForSubject([eventType], subject, RETRYABLE_OUTBOX_STATUSES),
483-
unlessPayloadContains
484-
? sql`not (${outboxEvent.payload}::jsonb @> ${JSON.stringify(unlessPayloadContains)}::jsonb)`
484+
onlyIfOlderThanPatch
485+
? sql`coalesce((${outboxEvent.payload} ->> ${onlyIfOlderThanPatch})::numeric, -1) < ${String(patch[onlyIfOlderThanPatch])}::numeric`
485486
: undefined
486487
)
487488
)

0 commit comments

Comments
 (0)