Skip to content

Commit 39f069b

Browse files
committed
fix(billing): a fresh idempotency key per seat write so a returning value is applied, not replayed
1 parent 1c32175 commit 39f069b

3 files changed

Lines changed: 65 additions & 12 deletions

File tree

‎apps/sim/lib/billing/webhooks/outbox-handlers.ts‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,8 +27,11 @@ const logger = createLogger('BillingOutboxHandlers')
2727
*/
2828
const MAX_SYNC_ATTEMPTS = 2
2929

30-
/** A fresh key per Stripe write: the SDK reuses it across its own network retries of that call. */
31-
function customerContactSyncIdempotencyKey(eventId: string): string {
30+
/**
31+
* A fresh key per Stripe write: the SDK reuses it across its own network retries of that call.
32+
* A key derived from the pushed value would be replayed, unapplied, once that value comes back.
33+
*/
34+
function syncWriteIdempotencyKey(eventId: string): string {
3235
return `outbox:${eventId}:${generateShortId()}`
3336
}
3437

@@ -293,7 +296,7 @@ const stripeSyncSubscriptionSeats: OutboxHandler<SubscriptionSeatsSyncPayload> =
293296
],
294297
proration_behavior: 'always_invoice',
295298
},
296-
{ idempotencyKey: `outbox:${ctx.eventId}:${row.plan}:${desiredSeats}` }
299+
{ idempotencyKey: syncWriteIdempotencyKey(ctx.eventId) }
297300
)
298301
}
299302

@@ -503,7 +506,7 @@ const stripeSyncCustomerContact: OutboxHandler<StripeSyncCustomerContactPayload>
503506
email: contact.email,
504507
...(contact.name ? { name: contact.name } : {}),
505508
},
506-
{ idempotencyKey: customerContactSyncIdempotencyKey(ctx.eventId) }
509+
{ idempotencyKey: syncWriteIdempotencyKey(ctx.eventId) }
507510
)
508511
}
509512

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

Lines changed: 53 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -25,12 +25,23 @@ import postgres from 'postgres'
2525
import type Stripe from 'stripe'
2626
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
2727

28-
const ADMIN_API_KEY = vi.hoisted(() => {
29-
const key = 'integration-fixture-admin-key'
30-
process.env.ADMIN_API_KEY = key
31-
process.env.STRIPE_PRICE_TEAM_25_MO = 'price_team_pro_tier_month'
32-
process.env.STRIPE_PRICE_TEAM_100_MO = 'price_team_max_tier_month'
33-
return key
28+
const { ADMIN_API_KEY, restoreEnvironment } = vi.hoisted(() => {
29+
const fixture: Record<string, string> = {
30+
ADMIN_API_KEY: 'integration-fixture-admin-key',
31+
STRIPE_PRICE_TEAM_25_MO: 'price_team_pro_tier_month',
32+
STRIPE_PRICE_TEAM_100_MO: 'price_team_max_tier_month',
33+
}
34+
const previous = Object.fromEntries(Object.keys(fixture).map((key) => [key, process.env[key]]))
35+
Object.assign(process.env, fixture)
36+
return {
37+
ADMIN_API_KEY: fixture.ADMIN_API_KEY,
38+
restoreEnvironment() {
39+
for (const [key, value] of Object.entries(previous)) {
40+
if (value === undefined) delete process.env[key]
41+
else process.env[key] = value
42+
}
43+
},
44+
}
3445
})
3546

3647
const database = vi.hoisted(() => ({
@@ -166,6 +177,7 @@ beforeEach(() => {
166177

167178
afterAll(async () => {
168179
resetEnvFlagsMock()
180+
restoreEnvironment()
169181
try {
170182
await connection`DROP SCHEMA ${connection(schemaName)} CASCADE`
171183
} finally {
@@ -1190,6 +1202,41 @@ describe('Team seat sync', () => {
11901202
)
11911203
})
11921204

1205+
it('pushes a seat count that changes away and back while its sync retries', async () => {
1206+
const org = await createOrganizationWithPlan('team', 1)
1207+
async function commitSeats(seats: number) {
1208+
await testDatabase.transaction(async (tx) => {
1209+
await tx.update(subscription).set({ seats }).where(eq(subscription.id, org.subscriptionId))
1210+
await enqueueSubscriptionSeatsSync(tx, {
1211+
subscriptionId: org.subscriptionId,
1212+
seats,
1213+
reason: 'member-change',
1214+
})
1215+
})
1216+
return latestOutboxEventId(
1217+
OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS,
1218+
org.subscriptionId
1219+
)
1220+
}
1221+
1222+
const seatSync = await commitSeats(2)
1223+
const firstPush = stripe.holdNextRequest('subscriptions.update')
1224+
const syncing = processEvent(seatSync)
1225+
await firstPush.reached
1226+
await commitSeats(3)
1227+
const secondPush = stripe.holdNextRequest('subscriptions.update')
1228+
firstPush.release()
1229+
await secondPush.reached
1230+
await commitSeats(2)
1231+
secondPush.release()
1232+
await expect(syncing).resolves.toBe('pending')
1233+
expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(3)
1234+
1235+
await makeDue(seatSync)
1236+
await expect(processEvent(seatSync)).resolves.toBe('completed')
1237+
expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(2)
1238+
})
1239+
11931240
it('does not revive an older seat count when its dead-lettered sync is requeued', async () => {
11941241
const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')])
11951242
const org = await createOrganizationWithPlan('team', 1)

‎packages/testing/src/mocks/stripe.mock.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -180,8 +180,11 @@ export function createInMemoryStripe() {
180180
}
181181
}
182182
if (params.metadata) {
183-
previousAttributes.metadata = current.metadata
184-
next.metadata = { ...current.metadata, ...params.metadata }
183+
const metadata = { ...current.metadata, ...params.metadata }
184+
if (JSON.stringify(metadata) !== JSON.stringify(current.metadata)) {
185+
previousAttributes.metadata = current.metadata
186+
next.metadata = metadata
187+
}
185188
}
186189
for (const item of params.items ?? []) {
187190
const target = next.items.data.find((existing) => existing.id === item.id)

0 commit comments

Comments
 (0)