diff --git a/apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts b/apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts index 944c0b94036..08ba707cba2 100644 --- a/apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts +++ b/apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts @@ -2,6 +2,7 @@ import { db } from '@sim/db' import { outboxEvent } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { toError } from '@sim/utils/errors' +import { toRecord } from '@sim/utils/object' import { and, eq, sql } from 'drizzle-orm' import { NextResponse } from 'next/server' import { adminV1RequeueOutboxEventContract } from '@/lib/api/contracts/v1/admin' @@ -11,6 +12,11 @@ import { ENTERPRISE_METADATA_SYNC_EVENT_TYPE, ENTERPRISE_PROVISION_EVENT_TYPE, } from '@/lib/billing/enterprise-outbox-events' +import { + isSubscriptionSyncEventType, + lockSubscriptionForSyncRetry, + recommitSubscriptionSync, +} from '@/lib/billing/webhooks/subscription-sync' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' import { withAdminAuthParams } from '@/app/api/v1/admin/middleware' @@ -28,7 +34,9 @@ const invalidOutboxEventResponse = (message: string) => * will retry it. Resets `attempts`, `lastError`, and `availableAt` so * the next poll picks it up. Only dead-lettered events can be * requeued — completed/pending/processing rows are rejected to avoid - * operator errors. + * operator errors. A Stripe subscription sync is re-committed with the + * subscription's current value, so the retry never revives the value it + * failed with. */ export const POST = withRouteHandler( withAdminAuthParams<{ id: string }>(async (request, context) => { @@ -61,23 +69,43 @@ export const POST = withRouteHandler( const deliveryRevision = metadataIntent?.success ? metadataIntent.data.deliveryRevision + 1 : null - const result = await db - .update(outboxEvent) - .set({ - status: 'pending', - attempts: 0, - lastError: null, - availableAt: new Date(), - lockedAt: null, - processedAt: null, - ...(deliveryRevision === null - ? {} - : { - payload: sql`(${outboxEvent.payload}::jsonb || ${JSON.stringify({ deliveryRevision })}::jsonb)::json`, - }), - }) - .where(and(eq(outboxEvent.id, id), eq(outboxEvent.status, 'dead_letter'))) - .returning({ id: outboxEvent.id, eventType: outboxEvent.eventType }) + const subscriptionId = toRecord(existing?.payload).subscriptionId + const subscriptionSync = + existing && + isSubscriptionSyncEventType(existing.eventType) && + typeof subscriptionId === 'string' + ? { eventType: existing.eventType, subscriptionId } + : null + const result = await db.transaction(async (tx) => { + if (subscriptionSync) { + await lockSubscriptionForSyncRetry(tx, subscriptionSync.subscriptionId) + } + const requeued = await tx + .update(outboxEvent) + .set({ + status: 'pending', + attempts: 0, + lastError: null, + availableAt: new Date(), + lockedAt: null, + processedAt: null, + ...(deliveryRevision === null + ? {} + : { + payload: sql`(${outboxEvent.payload}::jsonb || ${JSON.stringify({ deliveryRevision })}::jsonb)::json`, + }), + }) + .where(and(eq(outboxEvent.id, id), eq(outboxEvent.status, 'dead_letter'))) + .returning({ id: outboxEvent.id, eventType: outboxEvent.eventType }) + if (subscriptionSync && requeued.length > 0) { + await recommitSubscriptionSync( + tx, + subscriptionSync.eventType, + subscriptionSync.subscriptionId + ) + } + return requeued + }) if (result.length === 0) { return NextResponse.json( diff --git a/apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts b/apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts index ffc61afe9ba..0e88a8b1a3d 100644 --- a/apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts +++ b/apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts @@ -33,8 +33,7 @@ import { } from '@/lib/api/contracts/v1/admin' import { parseRequest } from '@/lib/api/server' import { requireStripeClient } from '@/lib/billing/stripe-client' -import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' -import { enqueueOutboxEvent } from '@/lib/core/outbox/service' +import { enqueueCancelAtPeriodEndSync } from '@/lib/billing/webhooks/subscription-sync' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' import { withAdminAuthParams } from '@/app/api/v1/admin/middleware' import { @@ -113,15 +112,17 @@ export const DELETE = withRouteHandler( } if (atPeriodEnd) { + const stripeSubscriptionId = existing.stripeSubscriptionId await db.transaction(async (tx) => { await tx .update(subscription) .set({ cancelAtPeriodEnd: true }) .where(eq(subscription.id, subscriptionId)) - await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, { - stripeSubscriptionId: existing.stripeSubscriptionId, + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId, subscriptionId: existing.id, + cancelAtPeriodEnd: true, reason: reason ?? 'admin-cancel-at-period-end', }) }) diff --git a/apps/sim/lib/admin/subscription-lifecycle.test.ts b/apps/sim/lib/admin/subscription-lifecycle.test.ts index 89a7f37e4b6..dc990730e2c 100644 --- a/apps/sim/lib/admin/subscription-lifecycle.test.ts +++ b/apps/sim/lib/admin/subscription-lifecycle.test.ts @@ -2,6 +2,7 @@ import { outboxEvent, subscription } from '@sim/db/schema' import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' import { auditMock, auditMockFns } from '@sim/testing/mocks/audit.mock' import { billingOutboxHandlersMock } from '@sim/testing/mocks/billing-outbox-handlers.mock' +import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock' import { organizationMembershipMock } from '@sim/testing/mocks/organization-membership.mock' import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock' import { stripeClientMock } from '@sim/testing/mocks/stripe.mock' @@ -24,6 +25,7 @@ vi.mock('@/lib/billing/organizations/membership', () => organizationMembershipMo vi.mock('@/lib/billing/stripe-client', () => stripeClientMock) vi.mock('@/lib/billing/webhooks/outbox-handlers', () => billingOutboxHandlersMock) vi.mock('@/lib/core/outbox/service', () => outboxServiceMock) +vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) import { refundDashboardSubscriptionPayment, @@ -98,6 +100,7 @@ describe('admin subscription cancellation', () => { it('requeues the same dead-lettered period-end cancellation operation', async () => { dbChainMockFns.returning.mockResolvedValueOnce([{ id: activeSubscription.id }]) + queueTableRows(outboxEvent, [{ subscriptionId: 'sub-row-1' }]) queueTableRows(outboxEvent, [ { id: 'outbox-1', @@ -126,6 +129,7 @@ describe('admin subscription cancellation', () => { }) it('replays an immediate cancellation after the webhook removed active entitlement', async () => { + queueTableRows(outboxEvent, []) queueTableRows(outboxEvent, [ { id: 'outbox-1', @@ -149,6 +153,7 @@ describe('admin subscription cancellation', () => { }) it('rejects reuse of a cancellation operation id with different timing', async () => { + queueTableRows(outboxEvent, []) queueTableRows(outboxEvent, [ { id: 'outbox-1', diff --git a/apps/sim/lib/admin/subscription-lifecycle.ts b/apps/sim/lib/admin/subscription-lifecycle.ts index 72335997f92..406eb5fadf0 100644 --- a/apps/sim/lib/admin/subscription-lifecycle.ts +++ b/apps/sim/lib/admin/subscription-lifecycle.ts @@ -8,6 +8,11 @@ import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/mem import { requireStripeClient } from '@/lib/billing/stripe-client' import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/utils' import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + enqueueCancelAtPeriodEndSync, + lockSubscriptionForSyncRetry, + recordCancelAtPeriodEnd, +} from '@/lib/billing/webhooks/subscription-sync' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' const RECENT_INVOICE_LIMIT = 12 @@ -257,6 +262,24 @@ export async function requestDashboardSubscriptionCancellation({ : 'admin-dashboard-cancel-at-period-end') const cancellation = await db.transaction(async (tx) => { await acquireOrganizationMutationLock(tx, organizationId) + const isThisOperation = and( + sql`${outboxEvent.payload} ->> 'operationId' = ${operationId}`, + sql`${outboxEvent.payload} ->> 'organizationId' = ${organizationId}` + ) + + const [retriedSync] = await tx + .select({ subscriptionId: sql`${outboxEvent.payload} ->> 'subscriptionId'` }) + .from(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END), + isThisOperation + ) + ) + .limit(1) + if (retriedSync?.subscriptionId) { + await lockSubscriptionForSyncRetry(tx, retriedSync.subscriptionId) + } const [existingOperation] = await tx .select({ @@ -273,8 +296,7 @@ export async function requestDashboardSubscriptionCancellation({ OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, OUTBOX_EVENT_TYPES.STRIPE_CANCEL_SUBSCRIPTION_IMMEDIATELY, ]), - sql`${outboxEvent.payload} ->> 'operationId' = ${operationId}`, - sql`${outboxEvent.payload} ->> 'organizationId' = ${organizationId}` + isThisOperation ) ) .for('update') @@ -315,6 +337,9 @@ export async function requestDashboardSubscriptionCancellation({ .where( and(eq(outboxEvent.id, existingOperation.id), eq(outboxEvent.status, 'dead_letter')) ) + if (existingOperation.eventType === OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END) { + await recordCancelAtPeriodEnd(tx, existingOperation.subscriptionId, true) + } return { operationId, outboxEventId: existingOperation.id, @@ -384,18 +409,15 @@ export async function requestDashboardSubscriptionCancellation({ .set({ cancelAtPeriodEnd: true }) .where(eq(subscription.id, subscriptionRow.id)) } - const eventId = await enqueueOutboxEvent( - tx, - OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, - { - operationId, - organizationId, - subscriptionId: subscriptionRow.id, - stripeSubscriptionId: subscriptionRow.stripeSubscriptionId, - reason: normalizedReason, - requestedBy: actor, - } - ) + const eventId = await enqueueCancelAtPeriodEndSync(tx, { + operationId, + organizationId, + subscriptionId: subscriptionRow.id, + stripeSubscriptionId: subscriptionRow.stripeSubscriptionId, + cancelAtPeriodEnd: true, + reason: normalizedReason, + requestedBy: actor, + }) return { operationId, outboxEventId: eventId, diff --git a/apps/sim/lib/auth/auth.ts b/apps/sim/lib/auth/auth.ts index 65ac5d554e0..4b317187d5a 100644 --- a/apps/sim/lib/auth/auth.ts +++ b/apps/sim/lib/auth/auth.ts @@ -103,6 +103,10 @@ import { handleSubscriptionCreated, handleSubscriptionDeleted, } from '@/lib/billing/webhooks/subscription' +import { + reconcileSubscriptionSyncFromStripe, + recordCustomerRestoreAfterHook, +} from '@/lib/billing/webhooks/subscription-sync' import { handleSubscriptionUsageUpdate } from '@/lib/billing/webhooks/subscription-usage' import { env } from '@/lib/core/config/env' import { @@ -1100,6 +1104,8 @@ export const auth = betterAuth({ return }), after: createAuthMiddleware(async (ctx) => { + if (isBillingEnabled) await recordCustomerRestoreAfterHook(ctx) + if (isBillingEnabled && ctx.path === '/subscription/upgrade') { const checkoutContext = ctx as typeof ctx & { billingCheckoutAdmissionClaim?: CheckoutAdmissionClaim @@ -1727,6 +1733,7 @@ export const auth = betterAuth({ case 'customer.subscription.created': case 'customer.subscription.updated': { await handleManualEnterpriseSubscription(event) + await reconcileSubscriptionSyncFromStripe(event) await handleSubscriptionUsageUpdate(event) break } diff --git a/apps/sim/lib/billing/enterprise-provisioning.ts b/apps/sim/lib/billing/enterprise-provisioning.ts index bb1d5dce165..cc3962f5a79 100644 --- a/apps/sim/lib/billing/enterprise-provisioning.ts +++ b/apps/sim/lib/billing/enterprise-provisioning.ts @@ -77,6 +77,10 @@ import { TERMINAL_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/util import { countPendingSeatInvitations } from '@/lib/billing/validation/seat-management' import { withEnterpriseReconciliationLease } from '@/lib/billing/webhooks/enterprise-reconciliation-lease' import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + lockSubscriptionForSyncRetry, + recommitSubscriptionSync, +} from '@/lib/billing/webhooks/subscription-sync' import { env } from '@/lib/core/config/env' import { continueOutboxHandler, @@ -2103,6 +2107,9 @@ export async function retryEnterpriseFollowUpJob( const retried = await db.transaction(async (tx) => { await acquireOrganizationMutationLock(tx, operationPayload.request.organizationId) + if (snapshotDetail.kind === 'personal_subscription_cancellation') { + await lockSubscriptionForSyncRetry(tx, snapshotDetail.subjectId) + } const [row] = await tx .select({ status: outboxEvent.status, @@ -2119,7 +2126,9 @@ export async function retryEnterpriseFollowUpJob( !detail || !getEnterpriseFollowUpOperationIds(row.eventType, row.payload).includes(operationId) || (detail.kind === 'member_reconciliation' && - detail.subjectId !== operationPayload.request.organizationId) + detail.subjectId !== operationPayload.request.organizationId) || + detail.kind !== snapshotDetail.kind || + detail.subjectId !== snapshotDetail.subjectId ) { throw new EnterpriseProvisioningError('Enterprise follow-up job not found') } @@ -2135,6 +2144,13 @@ export async function retryEnterpriseFollowUpJob( processedAt: null, }) .where(eq(outboxEvent.id, jobEventId)) + if (detail.kind === 'personal_subscription_cancellation') { + await recommitSubscriptionSync( + tx, + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + detail.subjectId + ) + } return true }) diff --git a/apps/sim/lib/billing/organizations/lock-order.test.ts b/apps/sim/lib/billing/organizations/lock-order.test.ts index 108bb594f63..b6234a37a9a 100644 --- a/apps/sim/lib/billing/organizations/lock-order.test.ts +++ b/apps/sim/lib/billing/organizations/lock-order.test.ts @@ -17,6 +17,7 @@ import { workspace, } from '@sim/db/schema' import { dbChainMockFns, resetDbChainMock } from '@sim/testing' +import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock' import { outboxServiceMock } from '@sim/testing/mocks/outbox-service.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' @@ -42,6 +43,7 @@ import type { DbOrTx } from '@/lib/db/types' import { attachOwnedWorkspacesToOrganizationTx } from '@/lib/workspaces/organization-workspaces' vi.mock('@/lib/core/outbox/service', () => outboxServiceMock) +vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) /** * A superset row that satisfies every read in the join path: a paid org sub, a diff --git a/apps/sim/lib/billing/organizations/membership.ts b/apps/sim/lib/billing/organizations/membership.ts index e93f4bcd0e2..c7a023759ac 100644 --- a/apps/sim/lib/billing/organizations/membership.ts +++ b/apps/sim/lib/billing/organizations/membership.ts @@ -48,6 +48,10 @@ import { import { toDecimal, toNumber } from '@/lib/billing/utils/decimal' import { validateSeatAvailability } from '@/lib/billing/validation/seat-management' import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + enqueueCancelAtPeriodEndSync, + isCancelAtPeriodEndSettled, +} from '@/lib/billing/webhooks/subscription-sync' import { isBillingEnabled } from '@/lib/core/config/env-flags' import { OrchestrationError } from '@/lib/core/orchestration/types' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' @@ -258,7 +262,17 @@ export async function restoreUserProSubscription(userId: string): Promise ({ @@ -8,11 +9,9 @@ vi.mock('@/lib/billing/storage/payer-transfer', () => ({ changeWorkspaceStoragePayersInTx: vi.fn(), })) vi.mock('@/lib/core/outbox/service', () => outboxServiceMock) +vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) import { pauseProSubscriptionForOrgCoverage } from '@/lib/billing/organizations/membership' -import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' - -const mockEnqueueOutboxEvent = outboxServiceMockFns.mockEnqueueOutboxEvent const ACTIVE_PERSONAL_PRO = { id: 'sub-personal', @@ -50,7 +49,7 @@ describe('pauseProSubscriptionForOrgCoverage', () => { resetDbChainMock() }) - it('pauses the personal Pro and queues the Stripe sync when an entitled paid org covers the user', async () => { + it('pauses the personal Pro when an entitled paid org covers the user', async () => { queueWhereResponses([ [{ organizationId: 'org-1' }], [{ plan: 'team_6000', referenceId: 'org-1' }], @@ -68,15 +67,6 @@ describe('pauseProSubscriptionForOrgCoverage', () => { organizationId: 'org-1', }) expect(dbChainMockFns.set).toHaveBeenCalledWith({ cancelAtPeriodEnd: true }) - expect(mockEnqueueOutboxEvent).toHaveBeenCalledWith( - expect.anything(), - OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, - { - stripeSubscriptionId: 'stripe-sub-personal', - subscriptionId: 'sub-personal', - reason: 'covered-by-organization', - } - ) }) it('reports covered even when no entitled personal Pro row exists', async () => { @@ -94,7 +84,6 @@ describe('pauseProSubscriptionForOrgCoverage', () => { organizationId: 'org-1', }) expect(dbChainMockFns.update).not.toHaveBeenCalled() - expect(mockEnqueueOutboxEvent).not.toHaveBeenCalled() }) it('reports covered without pausing again when the personal Pro is already pausing', async () => { @@ -113,6 +102,5 @@ describe('pauseProSubscriptionForOrgCoverage', () => { organizationId: 'org-1', }) expect(dbChainMockFns.update).not.toHaveBeenCalled() - expect(mockEnqueueOutboxEvent).not.toHaveBeenCalled() }) }) diff --git a/apps/sim/lib/billing/organizations/provision-seat.test.ts b/apps/sim/lib/billing/organizations/provision-seat.test.ts index 6a38cc50eb6..2129fb485f6 100644 --- a/apps/sim/lib/billing/organizations/provision-seat.test.ts +++ b/apps/sim/lib/billing/organizations/provision-seat.test.ts @@ -1,12 +1,13 @@ import { billingCoreMock, billingCoreMockFns } from '@sim/testing/mocks/billing-core.mock' import { billingOutboxHandlersMock } from '@sim/testing/mocks/billing-outbox-handlers.mock' import { billingPlanMock, billingPlanMockFns } from '@sim/testing/mocks/billing-plan.mock' +import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock' import { dbChainMockFns, resetDbChainMock } from '@sim/testing/mocks/database.mock' import { organizationMembershipMock, organizationMembershipMockFns, } from '@sim/testing/mocks/organization-membership.mock' -import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock' +import { outboxServiceMock } from '@sim/testing/mocks/outbox-service.mock' import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' const { @@ -41,6 +42,8 @@ vi.mock('@/lib/billing/plans', () => ({ vi.mock('@/lib/core/outbox/service', () => outboxServiceMock) +vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) + vi.mock('@/lib/billing/webhooks/outbox-handlers', () => billingOutboxHandlersMock) import { ensureTeamOrganizationForAcceptance } from '@/lib/billing/organizations/provision-seat' @@ -49,13 +52,23 @@ const { mockAcquireOrganizationMutationLock } = organizationMembershipMockFns const mockGetOrganizationSubscription = billingCoreMockFns.mockGetOrganizationSubscription const mockGetHighestPriorityPersonalSubscription = billingPlanMockFns.mockGetHighestPriorityPersonalSubscription -const enqueueMock = outboxServiceMockFns.mockEnqueueOutboxEvent -function testExecutor(onUpdate: () => void = () => {}) { +/** The subscription row as the activation re-reads it under its lock. */ +function testExecutor(onSubscriptionLock: () => void = () => {}) { + const lockedRow = { cancelAtPeriodEnd: false, seats: 1, stripeSubscriptionId: 'stripe_sub' } return { + select: () => ({ + from: () => ({ + where: () => ({ + for: () => { + onSubscriptionLock() + return { limit: () => Promise.resolve([lockedRow]) } + }, + }), + }), + }), update: () => ({ set: (values: Record) => { - onUpdate() updateCalls.value.push(values) return { where: () => Promise.resolve([]) } }, @@ -112,12 +125,6 @@ describe('ensureTeamOrganizationForAcceptance', () => { }, }) expect(updateCalls.value).toContainEqual(expect.objectContaining({ plan: 'team_6000' })) - // The Pro→Team price migration is durably enqueued at conversion time. - expect(enqueueMock).toHaveBeenCalledWith( - executor, - 'stripe.sync-subscription-seats', - expect.objectContaining({ subscriptionId: 'sub-pro' }) - ) expect(mockGetOrganizationSubscription).toHaveBeenCalledWith( 'org-1', expect.objectContaining({ executor }) @@ -202,20 +209,8 @@ describe('ensureTeamOrganizationForAcceptance', () => { executor, expect.objectContaining({ plan: 'team_6000', referenceId: 'owner-1' }) ) - // The plan change enqueues the price seat-sync... - expect(enqueueMock).toHaveBeenCalledWith( - executor, - 'stripe.sync-subscription-seats', - expect.objectContaining({ subscriptionId: 'sub-pro' }) - ) expect(dbChainMockFns.transaction).not.toHaveBeenCalled() expect(lockOrder).toEqual(['organization', 'subscription']) - // ...but with no scheduled cancellation there is no cancel-sync event. - expect(enqueueMock).not.toHaveBeenCalledWith( - expect.anything(), - 'stripe.sync-cancel-at-period-end', - expect.anything() - ) }) it('blocks personal Pro conversion when the reused organization has unresolved Enterprise', async () => { @@ -243,7 +238,6 @@ describe('ensureTeamOrganizationForAcceptance', () => { }) ).rejects.toThrow('Enterprise issuance is unfinished') expect(updateCalls.value).toHaveLength(0) - expect(enqueueMock).not.toHaveBeenCalled() }) it('provisions an org for a legacy personal-scoped Team subscription without a plan change', async () => { @@ -270,8 +264,6 @@ describe('ensureTeamOrganizationForAcceptance', () => { expect.anything(), expect.objectContaining({ plan: 'team', referenceId: 'owner-1' }) ) - // No plan change and no scheduled cancellation: nothing to push to Stripe. - expect(enqueueMock).not.toHaveBeenCalled() }) it('returns upgrade-required (no downgrade) when no eligible Team tier exists', async () => { diff --git a/apps/sim/lib/billing/organizations/provision-seat.ts b/apps/sim/lib/billing/organizations/provision-seat.ts index 0d5e55d79c5..4594e257e30 100644 --- a/apps/sim/lib/billing/organizations/provision-seat.ts +++ b/apps/sim/lib/billing/organizations/provision-seat.ts @@ -16,8 +16,13 @@ import { } from '@/lib/billing/plan-helpers' import { getPlanByName } from '@/lib/billing/plans' import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils' -import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' -import { enqueueOutboxEvent } from '@/lib/core/outbox/service' +import { + enqueueCancelAtPeriodEndSync, + enqueueSubscriptionSeatsSync, + isCancelAtPeriodEndSettled, + readCommittedSeats, + recordCancelAtPeriodEnd, +} from '@/lib/billing/webhooks/subscription-sync' import type { DbOrTx, DbTransaction } from '@/lib/db/types' const logger = createLogger('ProvisionSeat') @@ -236,40 +241,50 @@ async function convertPersonalSubscriptionToTeam( * the post-join seat reconcile is skipped or fails. Any scheduled cancellation * is cleared (DB + Stripe) so a freshly-activated Team is not left scheduled to * cancel, including the legacy personal-scoped Team case where the plan is - * unchanged. + * unchanged. The row is read under its lock, so a cancellation committed after + * the caller's earlier read is still cleared and recorded. */ async function activateTeamSubscription( - sub: { id: string; cancelAtPeriodEnd?: boolean | null; stripeSubscriptionId: string | null }, + sub: { id: string }, targetPlan: string, { planChanged }: { planChanged: boolean }, - executor: DbOrTx + tx: DbOrTx ): Promise { - const shouldClearCancellation = - Boolean(sub.cancelAtPeriodEnd) && Boolean(sub.stripeSubscriptionId) - - const apply = async (tx: DbOrTx) => { - await tx - .update(subscriptionTable) - .set({ plan: targetPlan, cancelAtPeriodEnd: false }) - .where(eq(subscriptionTable.id, sub.id)) + const [locked] = await tx + .select({ + cancelAtPeriodEnd: subscriptionTable.cancelAtPeriodEnd, + seats: subscriptionTable.seats, + stripeSubscriptionId: subscriptionTable.stripeSubscriptionId, + }) + .from(subscriptionTable) + .where(eq(subscriptionTable.id, sub.id)) + .for('update') + .limit(1) - if (planChanged) { - await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, { - subscriptionId: sub.id, - reason: 'pro-to-team-conversion', - }) - } + await tx + .update(subscriptionTable) + .set({ plan: targetPlan, cancelAtPeriodEnd: false }) + .where(eq(subscriptionTable.id, sub.id)) - if (shouldClearCancellation) { - await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, { - stripeSubscriptionId: sub.stripeSubscriptionId as string, - subscriptionId: sub.id, - reason: 'pro-to-team-conversion', - }) - } + if (planChanged) { + await enqueueSubscriptionSeatsSync(tx, { + subscriptionId: sub.id, + seats: await readCommittedSeats(tx, sub.id, locked?.seats ?? 1), + reason: 'pro-to-team-conversion', + }) } - await apply(executor) + if (!locked?.stripeSubscriptionId) return + if (!(await isCancelAtPeriodEndSettled(tx, sub.id, Boolean(locked.cancelAtPeriodEnd), false))) { + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId: locked.stripeSubscriptionId, + subscriptionId: sub.id, + cancelAtPeriodEnd: false, + reason: 'pro-to-team-conversion', + }) + } else { + await recordCancelAtPeriodEnd(tx, sub.id, false) + } } /** diff --git a/apps/sim/lib/billing/organizations/seats.test.ts b/apps/sim/lib/billing/organizations/seats.test.ts index d76bf2a02f9..51756140367 100644 --- a/apps/sim/lib/billing/organizations/seats.test.ts +++ b/apps/sim/lib/billing/organizations/seats.test.ts @@ -8,7 +8,7 @@ import { setEnvFlags, } from '@sim/testing' import { billingOutboxHandlersMock } from '@sim/testing/mocks/billing-outbox-handlers.mock' -import { outboxServiceMock, outboxServiceMockFns } from '@sim/testing/mocks/outbox-service.mock' +import { billingSubscriptionSyncMock } from '@sim/testing/mocks/billing-subscription-sync.mock' import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' const { mockSyncSubscriptionUsageLimits } = vi.hoisted(() => ({ @@ -19,7 +19,7 @@ vi.mock('@/lib/billing/organization', () => ({ syncSubscriptionUsageLimits: mockSyncSubscriptionUsageLimits, })) -vi.mock('@/lib/core/outbox/service', () => outboxServiceMock) +vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) vi.mock('@/lib/billing/webhooks/outbox-handlers', () => billingOutboxHandlersMock) @@ -27,8 +27,6 @@ vi.mock('@sim/audit', () => auditMock) import { reconcileOrganizationSeats } from '@/lib/billing/organizations/seats' -const enqueueMock = outboxServiceMockFns.mockEnqueueOutboxEvent - const teamSub = { id: 'sub-1', plan: 'team_6000', @@ -48,7 +46,6 @@ afterAll(resetEnvFlagsMock) describe('reconcileOrganizationSeats', () => { beforeEach(() => { resetDbChainMock() - enqueueMock.mockResolvedValue('evt-1') setEnvFlags({ isBillingEnabled: true }) }) @@ -56,7 +53,7 @@ describe('reconcileOrganizationSeats', () => { resetDbChainMock() }) - it('grows seats to the member count and enqueues a Stripe sync', async () => { + it('grows seats to the member count', async () => { queueReconcileReads([teamSub], [{ value: 2 }]) const result = await reconcileOrganizationSeats({ @@ -69,13 +66,9 @@ describe('reconcileOrganizationSeats', () => { previousSeats: 1, seats: 2, reason: undefined, - outboxEventId: 'evt-1', + outboxEventId: 'subscription-seats-sync-event', }) expect(dbChainMockFns.set).toHaveBeenCalledWith({ seats: 2 }) - expect(enqueueMock).toHaveBeenCalledWith(expect.anything(), 'stripe.sync-subscription-seats', { - subscriptionId: 'sub-1', - reason: 'member-accepted-invite', - }) expect(mockSyncSubscriptionUsageLimits).toHaveBeenCalledWith( expect.objectContaining({ id: 'sub-1', referenceId: 'org-1', seats: 2 }) ) @@ -91,7 +84,6 @@ describe('reconcileOrganizationSeats', () => { expect(result.changed).toBe(true) expect(dbChainMockFns.set).toHaveBeenCalledWith({ seats: 2 }) - expect(enqueueMock).toHaveBeenCalledOnce() }) it('still records the seat audit when the post-commit usage-limit sync fails', async () => { @@ -126,7 +118,6 @@ describe('reconcileOrganizationSeats', () => { expect(result.changed).toBe(true) expect(result.seats).toBe(2) expect(dbChainMockFns.set).toHaveBeenCalledWith({ seats: 2 }) - expect(enqueueMock).toHaveBeenCalled() }) it('never drops below one seat', async () => { diff --git a/apps/sim/lib/billing/organizations/seats.ts b/apps/sim/lib/billing/organizations/seats.ts index 7b2d02e70d6..7e7b804cae2 100644 --- a/apps/sim/lib/billing/organizations/seats.ts +++ b/apps/sim/lib/billing/organizations/seats.ts @@ -6,14 +6,21 @@ import { and, count, desc, eq, inArray } from 'drizzle-orm' import { syncSubscriptionUsageLimits } from '@/lib/billing/organization' import { isTeam } from '@/lib/billing/plan-helpers' import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/utils' -import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + enqueueSubscriptionSeatsSync, + readCommittedSeats, +} from '@/lib/billing/webhooks/subscription-sync' import { isBillingEnabled } from '@/lib/core/config/env-flags' -import { enqueueOutboxEvent } from '@/lib/core/outbox/service' import { captureServerEvent } from '@/lib/posthog/server' const logger = createLogger('OrganizationSeats') export interface ReconcileOrganizationSeatsResult { + /** + * True only when the seat count changed. Repairing a row the Stripe plugin left stale, back to + * the committed count, still rewrites the row and re-records the sync (`outboxEventId` is set) + * but reports false and records no seat audit or analytics event. + */ changed: boolean previousSeats?: number seats?: number @@ -101,9 +108,13 @@ export async function reconcileOrganizationSeats({ .where(eq(member.organizationId, organizationId)) const targetSeats = Math.max(1, memberCountRow?.value ?? 1) - const currentSeats = orgSubscription.seats ?? 1 + const currentSeats = await readCommittedSeats( + tx, + orgSubscription.id, + orgSubscription.seats ?? 1 + ) - if (targetSeats === currentSeats) { + if (targetSeats === currentSeats && targetSeats === (orgSubscription.seats ?? 1)) { return { kind: 'noop', seats: currentSeats } } @@ -112,14 +123,11 @@ export async function reconcileOrganizationSeats({ .set({ seats: targetSeats }) .where(eq(subscription.id, orgSubscription.id)) - const outboxEventId = await enqueueOutboxEvent( - tx, - OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, - { - subscriptionId: orgSubscription.id, - reason, - } - ) + const outboxEventId = await enqueueSubscriptionSeatsSync(tx, { + subscriptionId: orgSubscription.id, + seats: targetSeats, + reason, + }) return { kind: 'changed', @@ -168,6 +176,15 @@ export async function reconcileOrganizationSeats({ outboxEventId: outcome.outboxEventId, }) + if (outcome.seats === outcome.previousSeats) { + return { + changed: false, + previousSeats: outcome.previousSeats, + seats: outcome.seats, + outboxEventId: outcome.outboxEventId, + } + } + const increased = outcome.seats > outcome.previousSeats if (actorId) { recordAudit({ diff --git a/apps/sim/lib/billing/webhooks/outbox-events.ts b/apps/sim/lib/billing/webhooks/outbox-events.ts index 5a7b53307c5..7b027d32c4e 100644 --- a/apps/sim/lib/billing/webhooks/outbox-events.ts +++ b/apps/sim/lib/billing/webhooks/outbox-events.ts @@ -1,19 +1,14 @@ export const OUTBOX_EVENT_TYPES = { - /** - * Sync a subscription's `cancel_at_period_end` flag from our DB to - * Stripe. The handler reads the current DB value at processing time - * — so rapid cancel→uncancel→cancel sequences always converge on - * the last-committed DB state regardless of outbox ordering. Callers - * enqueue this event after every DB change to `cancelAtPeriodEnd`. - */ + /** Sync `cancelAtPeriodEnd` from our DB to Stripe; enqueue via `enqueueCancelAtPeriodEndSync`. */ STRIPE_SYNC_CANCEL_AT_PERIOD_END: 'stripe.sync-cancel-at-period-end', /** Cancel in Stripe; the verified deletion webhook remains the only DB entitlement authority. */ STRIPE_CANCEL_SUBSCRIPTION_IMMEDIATELY: 'stripe.cancel-subscription-immediately', /** * Sync a Team subscription's price and seat quantity from our DB to - * Stripe. The handler reads the current DB plan + seats at processing - * time and reconciles the Stripe item's price (e.g. after a Pro→Team - * conversion) and quantity, charging the proration via `always_invoice`. + * Stripe. Enqueue through `enqueueSubscriptionSeatsSync`. The handler + * reads the current DB plan + seats at processing time and reconciles + * the Stripe item's price (e.g. after a Pro→Team conversion) and + * quantity, charging the proration via `always_invoice`. * A failed charge surfaces through Stripe dunning and the existing * billing-blocked system, never under the synchronous accept path. */ diff --git a/apps/sim/lib/billing/webhooks/outbox-handlers.ts b/apps/sim/lib/billing/webhooks/outbox-handlers.ts index 1d3aa61109e..6c71a48589a 100644 --- a/apps/sim/lib/billing/webhooks/outbox-handlers.ts +++ b/apps/sim/lib/billing/webhooks/outbox-handlers.ts @@ -3,6 +3,7 @@ import { db } from '@sim/db' import { member, subscription as subscriptionTable, user } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' +import { generateShortId } from '@sim/utils/id' import { and, eq } from 'drizzle-orm' import { isTeam } from '@/lib/billing/plan-helpers' import { getPlanByName } from '@/lib/billing/plans' @@ -10,21 +11,28 @@ import { requireStripeClient } from '@/lib/billing/stripe-client' import { resolveDefaultPaymentMethod } from '@/lib/billing/stripe-payment-method' import { hasPaidSubscriptionStatus } from '@/lib/billing/subscriptions/utils' import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + type CancelAtPeriodEndSyncPayload, + cancelAtPeriodEndSyncIdempotencyKey, + readRecordedSyncValue, + type SubscriptionSeatsSyncPayload, +} from '@/lib/billing/webhooks/subscription-sync' import type { OutboxHandler } from '@/lib/core/outbox/service' const logger = createLogger('BillingOutboxHandlers') -interface StripeSyncCancelAtPeriodEndPayload { - stripeSubscriptionId: string - /** The DB subscription row id — also our source-of-truth pointer. */ - subscriptionId: string - /** Optional: reason this was enqueued — e.g. 'member-joined-paid-org'. */ - reason?: string - /** Correlates Enterprise-issuance follow-up work for Admin progress/retry. */ - sourceOperationId?: string - operationId?: string - organizationId?: string - requestedBy?: { id: string | null; name: string; email: string | null } +/** + * Passes a DB→Stripe sync handler makes before throwing for an outbox retry: each pass re-reads + * the row after its Stripe write and goes again when the value moved in the meantime. + */ +const MAX_SYNC_ATTEMPTS = 2 + +/** + * A fresh key per Stripe write: the SDK reuses it across its own network retries of that call. + * A key derived from the pushed value would be replayed, unapplied, once that value comes back. + */ +function syncWriteIdempotencyKey(eventId: string): string { + return `outbox:${eventId}:${generateShortId()}` } interface StripeCancelSubscriptionImmediatelyPayload { @@ -62,12 +70,6 @@ async function recordAdminCancellationAudit(params: { }) } -interface StripeSyncSubscriptionSeatsPayload { - /** The DB subscription row id — the handler reads current seats from this row. */ - subscriptionId: string - reason?: string -} - interface StripeSyncCustomerContactPayload { /** The DB subscription row id — handler resolves current owner/contact at processing time. */ subscriptionId: string @@ -102,42 +104,85 @@ async function getSubscriptionSeatSyncState(subscriptionId: string) { return row ?? null } -const stripeSyncCancelAtPeriodEnd: OutboxHandler = async ( +/** + * The value this sync should push: the one recorded on its own event, re-read now, which every + * commit keeps current on each sync that can still run. The subscription row is not used for + * this, because the Stripe plugin can overwrite it with a stale webhook payload before the + * reconcile step restores it. An event enqueued before values were recorded falls back to the + * row. Null when the subscription no longer exists. + */ +async function readDesiredCancelAtPeriodEnd( + eventId: string, + subscriptionId: string +): Promise { + const [row] = await db + .select({ cancelAtPeriodEnd: subscriptionTable.cancelAtPeriodEnd }) + .from(subscriptionTable) + .where(eq(subscriptionTable.id, subscriptionId)) + .limit(1) + if (!row) return null + return (await readRecordedSyncValue(eventId))?.cancelAtPeriodEnd ?? Boolean(row.cancelAtPeriodEnd) +} + +/** + * Pushes the latest committed value (see `readDesiredCancelAtPeriodEnd`), never the claim-time + * payload: racing events for one subscription each converge on it. Stripe is read first and + * written only when it differs, and the value is re-read after the write so one committed while + * this event's request was in flight is pushed too, even when an earlier event's request lands + * in Stripe after a newer one. + */ +const stripeSyncCancelAtPeriodEnd: OutboxHandler = async ( payload, ctx ) => { await recordAdminCancellationAudit({ ...payload, timing: 'period_end' }) - // Read the DB value at processing time (not at enqueue time). This - // makes the handler idempotent across racing enqueues: multiple - // events for the same subscription all push whatever the DB - // currently says, converging on the last committed value. - const rows = await db - .select({ cancelAtPeriodEnd: subscriptionTable.cancelAtPeriodEnd }) - .from(subscriptionTable) - .where(eq(subscriptionTable.id, payload.subscriptionId)) - .limit(1) + const stripe = requireStripeClient() + + for (let attempt = 1; attempt <= MAX_SYNC_ATTEMPTS; attempt++) { + const desiredValue = await readDesiredCancelAtPeriodEnd(ctx.eventId, payload.subscriptionId) + if (desiredValue === null) { + logger.warn('Subscription not found when syncing cancel_at_period_end', { + eventId: ctx.eventId, + subscriptionId: payload.subscriptionId, + }) + return + } - if (rows.length === 0) { - logger.warn('Subscription not found when syncing cancel_at_period_end', { + const stripeSubscription = await stripe.subscriptions.retrieve(payload.stripeSubscriptionId) + const needsUpdate = stripeSubscription.cancel_at_period_end !== desiredValue + if (needsUpdate) { + await stripe.subscriptions.update( + payload.stripeSubscriptionId, + { cancel_at_period_end: desiredValue }, + { idempotencyKey: cancelAtPeriodEndSyncIdempotencyKey(ctx.eventId) } + ) + } + + const latestValue = await readDesiredCancelAtPeriodEnd(ctx.eventId, payload.subscriptionId) + if (latestValue !== desiredValue) { + logger.info('cancel_at_period_end changed during Stripe sync; retrying latest value', { + eventId: ctx.eventId, + subscriptionId: payload.subscriptionId, + stripeSubscriptionId: payload.stripeSubscriptionId, + attemptedValue: desiredValue, + latestValue, + attempt, + }) + continue + } + + logger.info('Synced cancel_at_period_end from DB to Stripe', { + eventId: ctx.eventId, + stripeSubscriptionId: payload.stripeSubscriptionId, subscriptionId: payload.subscriptionId, + desiredValue, + alreadySynced: !needsUpdate, + reason: payload.reason, }) return } - const desiredValue = Boolean(rows[0].cancelAtPeriodEnd) - const stripe = requireStripeClient() - await stripe.subscriptions.update( - payload.stripeSubscriptionId, - { cancel_at_period_end: desiredValue }, - { idempotencyKey: `outbox:${ctx.eventId}` } - ) - logger.info('Synced cancel_at_period_end from DB to Stripe', { - eventId: ctx.eventId, - stripeSubscriptionId: payload.stripeSubscriptionId, - subscriptionId: payload.subscriptionId, - desiredValue, - reason: payload.reason, - }) + throw new Error(`cancel_at_period_end changed while syncing ${payload.subscriptionId}`) } const stripeCancelSubscriptionImmediately: OutboxHandler< @@ -160,14 +205,17 @@ const stripeCancelSubscriptionImmediately: OutboxHandler< }) } -const stripeSyncSubscriptionSeats: OutboxHandler = async ( +/** + * Pushes the seat count recorded on its own event, re-read each pass (falling back to the row + * for an event enqueued before values were recorded), for the same reason as the cancel sync. + */ +const stripeSyncSubscriptionSeats: OutboxHandler = async ( payload, ctx ) => { const stripe = requireStripeClient() - const maxSyncAttempts = 2 - for (let attempt = 1; attempt <= maxSyncAttempts; attempt++) { + for (let attempt = 1; attempt <= MAX_SYNC_ATTEMPTS; attempt++) { const row = await getSubscriptionSeatSyncState(payload.subscriptionId) if (!row) { logger.warn('Subscription not found when syncing seats', { @@ -203,7 +251,7 @@ const stripeSyncSubscriptionSeats: OutboxHandler = async ( @@ -369,33 +419,31 @@ const stripeThresholdOverageInvoice: OutboxHandler = async ( - payload, - ctx -) => { +type CustomerContactState = + | { status: 'ready'; stripeCustomerId: string; email: string; name: string } + | { status: 'skipped'; reason: string; organizationId?: string } + +async function readCustomerContact(subscriptionId: string): Promise { const [subscriptionRow] = await db .select({ referenceId: subscriptionTable.referenceId, stripeCustomerId: subscriptionTable.stripeCustomerId, }) .from(subscriptionTable) - .where(eq(subscriptionTable.id, payload.subscriptionId)) + .where(eq(subscriptionTable.id, subscriptionId)) .limit(1) if (!subscriptionRow) { - logger.warn('Subscription not found when syncing Stripe customer contact', { - eventId: ctx.eventId, - subscriptionId: payload.subscriptionId, - }) - return + return { + status: 'skipped', + reason: 'Subscription not found when syncing Stripe customer contact', + } } - if (!subscriptionRow.stripeCustomerId) { - logger.warn('Subscription has no Stripe customer id when syncing contact', { - eventId: ctx.eventId, - subscriptionId: payload.subscriptionId, - }) - return + return { + status: 'skipped', + reason: 'Subscription has no Stripe customer id when syncing contact', + } } const [owner] = await db @@ -409,29 +457,85 @@ const stripeSyncCustomerContact: OutboxHandler .limit(1) if (!owner) { - logger.warn('Organization owner not found when syncing Stripe customer contact', { + return { + status: 'skipped', + reason: 'Organization owner not found when syncing Stripe customer contact', + organizationId: subscriptionRow.referenceId, + } + } + + return { + status: 'ready', + stripeCustomerId: subscriptionRow.stripeCustomerId, + email: owner.email, + name: owner.name, + } +} + +/** + * Pushes the organization owner's current contact, re-reading it after the write so an + * ownership change committed while this event's request was in flight is pushed too. + */ +const stripeSyncCustomerContact: OutboxHandler = async ( + payload, + ctx +) => { + const stripe = requireStripeClient() + + for (let attempt = 1; attempt <= MAX_SYNC_ATTEMPTS; attempt++) { + const contact = await readCustomerContact(payload.subscriptionId) + if (contact.status === 'skipped') { + logger.warn(contact.reason, { + eventId: ctx.eventId, + subscriptionId: payload.subscriptionId, + ...(contact.organizationId ? { organizationId: contact.organizationId } : {}), + }) + return + } + + const customer = await stripe.customers.retrieve(contact.stripeCustomerId) + if (customer.deleted) { + throw new Error(`Stripe customer ${contact.stripeCustomerId} is deleted`) + } + const needsUpdate = + customer.email !== contact.email || Boolean(contact.name && customer.name !== contact.name) + if (needsUpdate) { + await stripe.customers.update( + contact.stripeCustomerId, + { + email: contact.email, + ...(contact.name ? { name: contact.name } : {}), + }, + { idempotencyKey: syncWriteIdempotencyKey(ctx.eventId) } + ) + } + + const latest = await readCustomerContact(payload.subscriptionId) + if ( + latest.status !== 'ready' || + latest.stripeCustomerId !== contact.stripeCustomerId || + latest.email !== contact.email || + latest.name !== contact.name + ) { + logger.info('Stripe customer contact changed during sync; retrying latest value', { + eventId: ctx.eventId, + subscriptionId: payload.subscriptionId, + attempt, + }) + continue + } + + logger.info('Synced Stripe customer contact', { eventId: ctx.eventId, + stripeCustomerId: contact.stripeCustomerId, subscriptionId: payload.subscriptionId, - organizationId: subscriptionRow.referenceId, + alreadySynced: !needsUpdate, + reason: payload.reason, }) return } - const stripe = requireStripeClient() - await stripe.customers.update( - subscriptionRow.stripeCustomerId, - { - email: owner.email, - ...(owner.name ? { name: owner.name } : {}), - }, - { idempotencyKey: `outbox:${ctx.eventId}` } - ) - logger.info('Synced Stripe customer contact', { - eventId: ctx.eventId, - stripeCustomerId: subscriptionRow.stripeCustomerId, - subscriptionId: payload.subscriptionId, - reason: payload.reason, - }) + throw new Error(`Stripe customer contact changed while syncing ${payload.subscriptionId}`) } export const billingOutboxHandlers = { diff --git a/apps/sim/lib/billing/webhooks/stripe-sync-convergence.integration.ts b/apps/sim/lib/billing/webhooks/stripe-sync-convergence.integration.ts new file mode 100644 index 00000000000..6e7777e611a --- /dev/null +++ b/apps/sim/lib/billing/webhooks/stripe-sync-convergence.integration.ts @@ -0,0 +1,1447 @@ +/** + * DB → Stripe sync convergence against real PostgreSQL, the real outbox worker path, and the + * real Better Auth Stripe webhook endpoint (plugin write first, then Sim's callbacks), with an + * in-memory Stripe that applies requests in the order the test releases them. + */ + +import { stripe as stripePlugin } from '@better-auth/stripe' +import * as schema from '@sim/db/schema' +import { member, organization, outboxEvent, subscription, user } from '@sim/db/schema' +import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { withUtcTimestamps } from '@sim/db/timestamps' +import { envFlagsMock, resetEnvFlagsMock, setEnvFlags } from '@sim/testing/mocks/env-flags.mock' +import { + createInMemoryStripe, + type InMemoryStripe, + stripeClientMock, +} from '@sim/testing/mocks/stripe.mock' +import { generateId } from '@sim/utils/id' +import { type BetterAuthOptions, betterAuth } from 'better-auth' +import { createAuthMiddleware } from 'better-auth/api' +import { and, desc, eq, sql } from 'drizzle-orm' +import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js' +import { NextRequest } from 'next/server' +import postgres from 'postgres' +import type Stripe from 'stripe' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +const { ADMIN_API_KEY, restoreEnvironment } = vi.hoisted(() => { + const fixture: Record = { + ADMIN_API_KEY: 'integration-fixture-admin-key', + STRIPE_PRICE_TEAM_25_MO: 'price_team_pro_tier_month', + STRIPE_PRICE_TEAM_100_MO: 'price_team_max_tier_month', + } + const previous = Object.fromEntries(Object.keys(fixture).map((key) => [key, process.env[key]])) + Object.assign(process.env, fixture) + return { + ADMIN_API_KEY: fixture.ADMIN_API_KEY, + restoreEnvironment() { + for (const [key, value] of Object.entries(previous)) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } + }, + } +}) + +const database = vi.hoisted(() => ({ + current: undefined as PostgresJsDatabase | undefined, +})) + +vi.mock('@sim/db', () => ({ + get db() { + if (!database.current) throw new Error('Stripe sync test database is not initialized') + return database.current + }, +})) +vi.mock('@/lib/billing/stripe-client', () => stripeClientMock) +vi.mock('@/lib/core/config/env-flags', () => envFlagsMock) + +import { requestDashboardSubscriptionCancellation } from '@/lib/admin/subscription-lifecycle' +import { createSimAuthAdapter } from '@/lib/auth/sim-auth-adapter' +import { CREDIT_TIERS } from '@/lib/billing/constants' +import { + pauseProSubscriptionForOrgCoverage, + restoreUserProSubscription, +} from '@/lib/billing/organizations/membership' +import { ensureTeamOrganizationForAcceptance } from '@/lib/billing/organizations/provision-seat' +import { reconcileOrganizationSeats } from '@/lib/billing/organizations/seats' +import { isTeam } from '@/lib/billing/plan-helpers' +import { syncSeatsFromStripeQuantity } from '@/lib/billing/validation/seat-management' +import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers' +import { + enqueueCancelAtPeriodEndSync, + enqueueSubscriptionSeatsSync, + reconcileSubscriptionSyncFromStripe, + recordCustomerRestoreAfterHook, +} from '@/lib/billing/webhooks/subscription-sync' +import { enqueueOutboxEvent, processOutboxEventById } from '@/lib/core/outbox/service' +import { POST as requeueOutboxEvent } from '@/app/api/v1/admin/outbox/[id]/requeue/route' + +const schemaName = `stripe_sync_${generateId().replaceAll('-', '')}` +const connection = postgres( + readTestDatabaseUrl(), + withUtcTimestamps({ + max: 6, + prepare: false, + fetch_types: false, + connection: { search_path: schemaName }, + onnotice: () => {}, + }) +) +const testDatabase = drizzle(connection, { schema }) + +let stripe: InMemoryStripe + +/** Runs between the plugin's write and Sim's reconcile step, to place work in that window. */ +let beforeReconcile: (() => Promise) | undefined + +/** + * Sim's Better Auth Stripe wiring from `lib/auth/auth.ts`, reduced to the parts that touch the + * synced fields: the seat sync in `onSubscriptionUpdate`, the reconcile step in `onEvent`, and + * the restore step in the `after` hook. + */ +function createTestAuth() { + return betterAuth({ + baseURL: 'http://localhost:3000', + secret: 'isolated-integration-fixture-secret-not-a-real-credential', + database: (options: BetterAuthOptions) => createSimAuthAdapter(options, testDatabase), + emailAndPassword: { enabled: true }, + hooks: { after: createAuthMiddleware(recordCustomerRestoreAfterHook) }, + plugins: [ + stripePlugin({ + stripeClient: stripe.client, + stripeWebhookSecret: 'whsec_fixture', + subscription: { + enabled: true, + plans: [], + onSubscriptionUpdate: async ({ event, subscription: updated }) => { + const stripeSubscription = event.data.object as Stripe.Subscription + if (!isTeam(updated.plan)) return + await syncSeatsFromStripeQuantity( + updated.id, + updated.seats ?? null, + stripeSubscription.items?.data?.[0]?.quantity || 1 + ) + }, + }, + onEvent: async (event) => { + await beforeReconcile?.() + await reconcileSubscriptionSyncFromStripe(event) + }, + }), + ], + }) +} + +let auth: ReturnType + +async function deliver(event: Stripe.Event) { + const response = await auth.handler( + new Request('http://localhost:3000/api/auth/stripe/webhook', { + method: 'POST', + headers: { 'content-type': 'application/json', 'stripe-signature': 't=1,v1=fixture' }, + body: JSON.stringify(event), + }) + ) + expect(response.status).toBe(200) +} + +beforeAll(async () => { + await connection`CREATE SCHEMA ${connection(schemaName)}` + for (const table of [ + 'subscription', + 'outbox_event', + 'member', + 'user', + 'organization', + 'workspace', + 'permissions', + 'audit_log', + 'session', + 'account', + 'verification', + ]) { + await connection.unsafe(`CREATE TABLE "${table}" (LIKE public."${table}" INCLUDING ALL)`) + } + database.current = testDatabase +}) + +beforeEach(() => { + stripe = createInMemoryStripe() + stripeClientMock.requireStripeClient.mockReturnValue(stripe.client) + beforeReconcile = undefined + auth = createTestAuth() +}) + +afterAll(async () => { + resetEnvFlagsMock() + restoreEnvironment() + try { + await connection`DROP SCHEMA ${connection(schemaName)} CASCADE` + } finally { + await connection.end() + database.current = undefined + } +}) + +async function createUser(label: string) { + const id = generateId() + const now = new Date() + await testDatabase.insert(user).values({ + id, + name: label, + email: `${label}-${id}@example.com`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + return { id, email: `${label}-${id}@example.com`, name: label } +} + +async function createOrganizationWithPlan(plan: 'team', seats = 1) { + const organizationId = generateId() + await testDatabase + .insert(organization) + .values({ id: organizationId, name: 'Org', slug: organizationId }) + const subscriptionId = generateId() + const stripeSubscriptionId = `sub_${subscriptionId}` + const stripeCustomerId = `cus_${subscriptionId}` + await testDatabase.insert(subscription).values({ + id: subscriptionId, + plan, + referenceId: organizationId, + status: 'active', + seats, + stripeSubscriptionId, + stripeCustomerId, + cancelAtPeriodEnd: false, + }) + stripe.addSubscription({ id: stripeSubscriptionId, customer: stripeCustomerId, quantity: seats }) + return { organizationId, subscriptionId, stripeSubscriptionId, stripeCustomerId } +} + +async function addMember(organizationId: string, userId: string, role = 'member') { + await testDatabase + .insert(member) + .values({ id: generateId(), organizationId, userId, role, createdAt: new Date() }) +} + +/** A user on personal Pro, synced with Stripe, who belongs to a paid Team organization. */ +async function createProUserInPaidOrganization(existingUserId?: string) { + const proUser = existingUserId ? { id: existingUserId } : await createUser('pro') + const subscriptionId = generateId() + const stripeSubscriptionId = `sub_${subscriptionId}` + await testDatabase.insert(subscription).values({ + id: subscriptionId, + plan: 'pro', + referenceId: proUser.id, + status: 'active', + seats: 1, + stripeSubscriptionId, + stripeCustomerId: `cus_${subscriptionId}`, + cancelAtPeriodEnd: false, + }) + stripe.addSubscription({ id: stripeSubscriptionId, customer: `cus_${subscriptionId}` }) + const paidOrganization = await createOrganizationWithPlan('team') + await addMember(paidOrganization.organizationId, proUser.id) + return { userId: proUser.id, subscriptionId, stripeSubscriptionId, paidOrganization } +} + +async function leaveOrganization(userId: string, organizationId: string) { + await testDatabase + .delete(member) + .where(and(eq(member.userId, userId), eq(member.organizationId, organizationId))) +} + +/** + * The sync event enqueued last for a subscription. `created_at` is the enqueuing transaction's + * start time, compared at microsecond precision, so it orders events from different + * transactions. No test enqueues two of one type for one subscription in a single transaction, + * and a tie fails loudly rather than being broken arbitrarily: `outbox_event` has no per-insert + * sequence to break it with. + */ +async function latestOutboxEventId(eventType: string, subscriptionId: string) { + const [latest, previous] = await testDatabase + .select({ id: outboxEvent.id, createdAt: sql`${outboxEvent.createdAt}::text` }) + .from(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, eventType), + sql`${outboxEvent.payload} ->> 'subscriptionId' = ${subscriptionId}` + ) + ) + .orderBy(desc(outboxEvent.createdAt)) + .limit(2) + if (!latest) throw new Error(`No ${eventType} event for ${subscriptionId}`) + if (previous?.createdAt === latest.createdAt) { + throw new Error(`Two ${eventType} events for ${subscriptionId} share one enqueue time`) + } + return latest.id +} + +function processEvent(eventId: string) { + return processOutboxEventById(eventId, billingOutboxHandlers) +} + +async function makeDue(eventId: string) { + await testDatabase + .update(outboxEvent) + .set({ availableAt: new Date() }) + .where(eq(outboxEvent.id, eventId)) +} + +/** Runs the event's last attempt with Stripe unavailable, so it dead-letters without applying. */ +async function deadLetter(eventId: string) { + await testDatabase.update(outboxEvent).set({ maxAttempts: 1 }).where(eq(outboxEvent.id, eventId)) + stripe.failNextRequest('subscriptions.update') + await expect(processEvent(eventId)).resolves.toBe('dead_letter') +} + +async function requeueFromAdminApi(eventId: string) { + const response = await requeueOutboxEvent( + new NextRequest(`http://localhost:3000/api/v1/admin/outbox/${eventId}/requeue`, { + method: 'POST', + headers: { 'x-admin-key': ADMIN_API_KEY }, + }), + { params: Promise.resolve({ id: eventId }) } + ) + expect(response.status).toBe(200) +} + +/** Delivers a subscription update Stripe makes on its own, e.g. a renewal. */ +async function deliverUnrelatedUpdate(stripeSubscriptionId: string) { + stripe.updateOutsideSim(stripeSubscriptionId, { metadata: { renewedAt: generateId() } }) + await deliver(stripe.events.at(-1) as Stripe.Event) +} + +async function signUp(label: string) { + const email = `${label}-${generateId()}@example.com` + const response = await auth.api.signUpEmail({ + body: { email, password: 'integration-fixture-password', name: label }, + asResponse: true, + }) + expect(response.status).toBe(200) + const { user: created } = (await response.json()) as { user: { id: string } } + const cookie = response.headers + .getSetCookie() + .map((header) => header.split(';')[0]) + .join('; ') + return { userId: created.id, cookie } +} + +function restoreSubscription(cookie: string) { + return auth.handler( + new Request('http://localhost:3000/api/auth/subscription/restore', { + method: 'POST', + headers: { 'content-type': 'application/json', cookie, origin: 'http://localhost:3000' }, + body: '{}', + }) + ) +} + +async function cancelValuesOfRetryableSyncs(subscriptionId: string) { + const rows = await testDatabase + .select({ payload: outboxEvent.payload }) + .from(outboxEvent) + .where( + and( + eq(outboxEvent.eventType, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END), + sql`${outboxEvent.status} in ('pending', 'processing', 'dead_letter')`, + sql`${outboxEvent.payload} ->> 'subscriptionId' = ${subscriptionId}` + ) + ) + return rows.map((row) => (row.payload as { cancelAtPeriodEnd?: boolean }).cancelAtPeriodEnd) +} + +async function storedSubscription(subscriptionId: string) { + const [row] = await testDatabase + .select({ cancelAtPeriodEnd: subscription.cancelAtPeriodEnd, seats: subscription.seats }) + .from(subscription) + .where(eq(subscription.id, subscriptionId)) + if (!row) throw new Error(`Subscription ${subscriptionId} not found`) + return row +} + +/** Resolves once another backend is blocked on a lock, i.e. the racing transaction is parked. */ +type TestTransaction = Parameters[0]>[0] + +/** + * Starts a transaction that takes its locks in `holdLocks`, then parks until released and runs + * `finish`. `untilBlocking` resolves once another backend is waiting on one of its locks. + */ +function startParkedTransaction( + holdLocks: (tx: TestTransaction) => Promise, + finish: (tx: TestTransaction) => Promise = async () => {} +) { + let release: () => void = () => {} + const released = new Promise((resolve) => { + release = resolve + }) + let reportPid: (pid: number) => void = () => {} + const holderPid = new Promise((resolve) => { + reportPid = resolve + }) + const done = testDatabase.transaction(async (tx) => { + await holdLocks(tx) + const [row] = await tx.execute<{ pid: number }>(sql`select pg_backend_pid() as pid`) + reportPid(row.pid) + await released + await finish(tx) + }) + async function untilBlocking() { + const pid = await holderPid + for (let attempt = 0; attempt < 200; attempt++) { + const [row] = await connection<{ blocked: number }[]>` + select count(*)::int as blocked from pg_stat_activity + where ${pid}::int = any(pg_blocking_pids(pid))` + if (row.blocked > 0) return + await new Promise((resolve) => setImmediate(resolve)) + } + throw new Error('No transaction ever waited on the parked one') + } + return { done, release, untilBlocking } +} + +describe('cancel_at_period_end sync', () => { + it('pushes the latest value when an earlier sync lands in Stripe after a newer one', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + const gate = stripe.holdNextRequest('subscriptions.update') + const pausing = processEvent(pauseSync) + await gate.reached + + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + + gate.release() + await pausing + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + + for (const event of stripe.events) await deliver(event) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + }) + + it('keeps the latest committed value when the echo of an earlier sync arrives mid-sync', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + const stalePush = stripe.holdNextRequest('subscriptions.update') + const pausing = processEvent(pauseSync) + await stalePush.reached + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + + const correctingPush = stripe.holdNextRequest('subscriptions.update') + stalePush.release() + await correctingPush.reached + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + await deliver(stripe.events.at(-1) as Stripe.Event) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + + correctingPush.release() + await expect(pausing).resolves.toBe('completed') + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('retries after a request reached Stripe but failed on the client and the value then changed', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + await makeDue(pauseSync) + + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('keeps a change that has not reached Stripe when an unrelated subscription update arrives', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + const renewal = stripe.updateOutsideSim(pro.stripeSubscriptionId, { + metadata: { renewedAt: 'period-2' }, + }) + expect(renewal.cancel_at_period_end).toBe(false) + await deliver(stripe.events.at(-1) as Stripe.Event) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('keeps the live Stripe value when its webhooks arrive out of order', async () => { + const pro = await createProUserInPaidOrganization() + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + const [cancelled, renewed] = stripe.events + + await deliver(renewed) + await deliver(cancelled) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + }) + + it('treats clearing a scheduled cancel_at in Stripe as a change made in Stripe', async () => { + const pro = await createProUserInPaidOrganization() + const scheduledEnd = Math.floor(Date.now() / 1000) + 7 * 24 * 60 * 60 + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: scheduledEnd }) + await deliver(stripe.events.at(-1) as Stripe.Event) + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: '' }) + const restore = stripe.events.at(-1) as Stripe.Event + expect(Object.keys(restore.data.previous_attributes ?? {})).toEqual(['cancel_at']) + await deliver(restore) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(new Set(await cancelValuesOfRetryableSyncs(pro.subscriptionId))).toEqual( + new Set([false]) + ) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('records a later Stripe read even when every sync already carries its value', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + const restoreRead = stripe.holdNextRequest('subscriptions.retrieve') + const reconcilingRestore = deliver(stripe.events.at(-1) as Stripe.Event) + await restoreRead.reached + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + const cancelRead = stripe.holdNextRequest('subscriptions.retrieve') + const reconcilingCancel = deliver(stripe.events.at(-1) as Stripe.Event) + await cancelRead.reached + + cancelRead.release() + await reconcilingCancel + restoreRead.release() + await reconcilingRestore + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + await makeDue(pauseSync) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('keeps a pending value when an unrelated Stripe update moves the cancellation date', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + const periodEnd = Math.floor(Date.now() / 1000) + 30 * 24 * 60 * 60 + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: periodEnd }) + await deliver(stripe.events.at(-1) as Stripe.Event) + + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at: periodEnd + 335 * 24 * 60 * 60 }) + const moved = stripe.events.at(-1) as Stripe.Event + expect(Object.keys(moved.data.previous_attributes ?? {})).toEqual(['cancel_at']) + await deliver(moved) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('lets a change made in Stripe while a sync is pending win over the pending value', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + for (const event of stripe.events) await deliver(event) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + it('keeps a renewal made in Stripe after later updates while an earlier sync retries', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + await deliver(stripe.events.at(-1) as Stripe.Event) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + await deliver(stripe.events.at(-1) as Stripe.Event) + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + await makeDue(pauseSync) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('keeps a cancel-then-renew made in Stripe after a later update while a sync is pending', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + await deliver(stripe.events.at(-1) as Stripe.Event) + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + await deliver(stripe.events.at(-1) as Stripe.Event) + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('does not restore an older value while its slow sync is still running', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + const slowPush = stripe.holdNextRequest('subscriptions.update') + const pausing = processEvent(pauseSync) + await slowPush.reached + + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + + slowPush.release() + await expect(pausing).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('does not revive an older value when its dead-lettered sync is requeued', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await deadLetter(pauseSync) + + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + + await requeueFromAdminApi(pauseSync) + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('leaves the field alone while a sync enqueued by an older deploy is in flight', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await testDatabase.transaction(async (tx) => { + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: false }) + .where(eq(subscription.id, pro.subscriptionId)) + await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, { + stripeSubscriptionId: pro.stripeSubscriptionId, + subscriptionId: pro.subscriptionId, + reason: 'member-left-paid-org', + }) + }) + + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + }) + + it('restores a pending value without reading Stripe', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { metadata: { renewedAt: 'period-2' } }) + stripe.failNextRequest('subscriptions.retrieve') + await deliver(stripe.events.at(-1) as Stripe.Event) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + }) + + it('pushes the committed value when its sync runs between the plugin write and the reconcile', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + beforeReconcile = async () => { + beforeReconcile = undefined + await expect(processEvent(pauseSync)).resolves.toBe('completed') + } + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + }) + + it('orders concurrent reconciles of Stripe-side changes by when each read Stripe', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + const renewal = stripe.events.at(-1) as Stripe.Event + const firstRead = stripe.holdNextRequest('subscriptions.retrieve') + const secondRead = stripe.holdNextRequest('subscriptions.retrieve') + const reconcilingRenewal = deliver(renewal) + await firstRead.reached + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + const reconcilingCancel = deliver(stripe.events.at(-1) as Stripe.Event) + await secondRead.reached + + firstRead.release() + await reconcilingRenewal + secondRead.release() + await reconcilingCancel + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + await makeDue(pauseSync) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('keeps a value Sim committed after the reconcile read Stripe', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + const liveRead = stripe.holdNextRequest('subscriptions.retrieve') + const delivering = deliver(stripe.events.at(-1) as Stripe.Event) + await liveRead.reached + await testDatabase.transaction(async (tx) => { + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: true }) + .where(eq(subscription.id, pro.subscriptionId)) + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId: pro.stripeSubscriptionId, + subscriptionId: pro.subscriptionId, + cancelAtPeriodEnd: true, + reason: 'admin-cancel-at-period-end', + }) + }) + liveRead.release() + await delivering + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + const adminSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(adminSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('requeues with the pending value when the plugin has overwritten the row', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await deadLetter(pauseSync) + await testDatabase.transaction(async (tx) => { + await tx + .select({ id: subscription.id }) + .from(subscription) + .where(eq(subscription.id, pro.subscriptionId)) + .for('update') + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId: pro.stripeSubscriptionId, + subscriptionId: pro.subscriptionId, + cancelAtPeriodEnd: true, + reason: 'admin-cancel-at-period-end', + }) + }) + const adminSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + beforeReconcile = async () => { + beforeReconcile = undefined + await requeueFromAdminApi(pauseSync) + } + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + await expect(processEvent(adminSync)).resolves.toBe('completed') + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('keeps a pause committed after the reconcile read a customer restore from Stripe', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: false }) + const liveRead = stripe.holdNextRequest('subscriptions.retrieve') + const reconcilingRestore = deliver(stripe.events.at(-1) as Stripe.Event) + await liveRead.reached + await pauseProSubscriptionForOrgCoverage(pro.userId) + liveRead.release() + await reconcilingRestore + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + const latestSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(latestSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(true) + }) + + it('restores the personal Pro when its member leaves while the plugin has overwritten the row', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + + beforeReconcile = async () => { + beforeReconcile = undefined + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + } + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('does not revive an older value when a retry path resets its sync without re-committing', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await deadLetter(pauseSync) + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + const restoreSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await expect(processEvent(restoreSync)).resolves.toBe('completed') + + await testDatabase + .update(outboxEvent) + .set({ + status: 'pending', + attempts: 0, + lastError: null, + availableAt: new Date(), + lockedAt: null, + processedAt: null, + }) + .where(eq(outboxEvent.id, pauseSync)) + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it("treats a write from an older deploy's sync handler as Sim's own echo", async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await leaveOrganization(pro.userId, pro.paidOrganization.organizationId) + await restoreUserProSubscription(pro.userId) + + await stripe.client.subscriptions.update( + pro.stripeSubscriptionId, + { cancel_at_period_end: true }, + { idempotencyKey: `outbox:${pauseSync}` } + ) + await deliver(stripe.events.at(-1) as Stripe.Event) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + }) + + it("records a customer's restore through the real endpoint while an earlier sync is retrying", async () => { + const customer = await signUp('restorer') + const pro = await createProUserInPaidOrganization(customer.userId) + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + stripe.failNextUpdateAfterApplying('subscriptions') + await expect(processEvent(pauseSync)).resolves.toBe('pending') + + expect((await restoreSubscription(customer.cookie)).status).toBe(200) + const restoreEvent = stripe.events.at(-1) as Stripe.Event + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(new Set(await cancelValuesOfRetryableSyncs(pro.subscriptionId))).toEqual( + new Set([false]) + ) + + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + await makeDue(pauseSync) + await expect(processEvent(pauseSync)).resolves.toBe('completed') + await deliver(restoreEvent) + + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(stripe.subscription(pro.stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) + + it('records nothing when the restore endpoint refuses the request', async () => { + const customer = await signUp('not-cancelling') + const pro = await createProUserInPaidOrganization(customer.userId) + + expect((await restoreSubscription(customer.cookie)).status).toBe(400) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + expect(await cancelValuesOfRetryableSyncs(pro.subscriptionId)).toEqual([]) + }) +}) + +describe('Team activation', () => { + it('records the cleared cancellation when a cancel is committed while it activates Team', async () => { + const owner = await createUser('owner') + const subscriptionId = generateId() + const stripeSubscriptionId = `sub_${subscriptionId}` + await testDatabase.insert(subscription).values({ + id: subscriptionId, + plan: 'team', + referenceId: owner.id, + status: 'active', + seats: 1, + stripeSubscriptionId, + stripeCustomerId: `cus_${subscriptionId}`, + cancelAtPeriodEnd: false, + }) + stripe.addSubscription({ id: stripeSubscriptionId, customer: `cus_${subscriptionId}` }) + + const cancelling = startParkedTransaction(async (tx) => { + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: true }) + .where(eq(subscription.id, subscriptionId)) + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId, + subscriptionId, + cancelAtPeriodEnd: true, + reason: 'admin-cancel-at-period-end', + }) + }) + const activating = testDatabase.transaction((tx) => + ensureTeamOrganizationForAcceptance({ + billingOwnerUserId: owner.id, + workspaceOrganizationId: null, + executor: tx, + workspaceIdsToAttach: [], + }) + ) + await cancelling.untilBlocking() + cancelling.release() + await cancelling.done + await expect(activating).resolves.toMatchObject({ success: true }) + expect((await storedSubscription(subscriptionId)).cancelAtPeriodEnd).toBe(false) + + await deliverUnrelatedUpdate(stripeSubscriptionId) + expect((await storedSubscription(subscriptionId)).cancelAtPeriodEnd).toBe(false) + const cancelSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + subscriptionId + ) + await expect(processEvent(cancelSync)).resolves.toBe('completed') + expect(stripe.subscription(stripeSubscriptionId).cancel_at_period_end).toBe(false) + }) +}) + +describe('operator retry', () => { + it('requeues a dead-lettered sync while a writer commits a new value for the subscription', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + const pauseSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + pro.subscriptionId + ) + await deadLetter(pauseSync) + + const writing = startParkedTransaction( + async (tx) => { + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: false }) + .where(eq(subscription.id, pro.subscriptionId)) + }, + async (tx) => { + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId: pro.stripeSubscriptionId, + subscriptionId: pro.subscriptionId, + cancelAtPeriodEnd: false, + reason: 'member-left-paid-org', + }) + } + ) + const requeuing = requeueFromAdminApi(pauseSync) + await writing.untilBlocking() + writing.release() + + await expect(Promise.all([writing.done, requeuing])).resolves.toBeDefined() + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(false) + }) + + it('retries a dead-lettered dashboard cancellation while a writer commits a new value', async () => { + const org = await createOrganizationWithPlan('team') + const operationId = generateId() + const actor = { id: null, name: 'Admin', email: null } + await requestDashboardSubscriptionCancellation({ + organizationId: org.organizationId, + operationId, + timing: 'period_end', + actor, + }) + const cancelSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END, + org.subscriptionId + ) + await deadLetter(cancelSync) + + const writing = startParkedTransaction( + async (tx) => { + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: false }) + .where(eq(subscription.id, org.subscriptionId)) + }, + async (tx) => { + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId: org.stripeSubscriptionId, + subscriptionId: org.subscriptionId, + cancelAtPeriodEnd: false, + reason: 'pro-to-team-conversion', + }) + } + ) + const retrying = requestDashboardSubscriptionCancellation({ + organizationId: org.organizationId, + operationId, + timing: 'period_end', + actor, + }) + await writing.untilBlocking() + writing.release() + + await expect(Promise.all([writing.done, retrying])).resolves.toBeDefined() + expect((await storedSubscription(org.subscriptionId)).cancelAtPeriodEnd).toBe(true) + }) +}) + +describe('customer contact sync', () => { + it('pushes the current owner when an earlier sync lands in Stripe after a newer one', async () => { + const [first, second, third] = await Promise.all([ + createUser('first-owner'), + createUser('second-owner'), + createUser('third-owner'), + ]) + const org = await createOrganizationWithPlan('team') + stripe.addCustomer({ id: org.stripeCustomerId, email: first.email, name: first.name }) + await addMember(org.organizationId, first.id, 'owner') + await addMember(org.organizationId, second.id) + await addMember(org.organizationId, third.id) + + async function transferOwnership(from: string, to: string) { + await testDatabase.transaction(async (tx) => { + await tx.update(member).set({ role: 'admin' }).where(eq(member.userId, from)) + await tx.update(member).set({ role: 'owner' }).where(eq(member.userId, to)) + await enqueueOutboxEvent(tx, OUTBOX_EVENT_TYPES.STRIPE_SYNC_CUSTOMER_CONTACT, { + subscriptionId: org.subscriptionId, + reason: 'ownership-transfer', + }) + }) + return latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_CUSTOMER_CONTACT, + org.subscriptionId + ) + } + + const firstSync = await transferOwnership(first.id, second.id) + const gate = stripe.holdNextRequest('customers.update') + const firstSyncRun = processEvent(firstSync) + await gate.reached + + const secondSync = await transferOwnership(second.id, third.id) + await expect(processEvent(secondSync)).resolves.toBe('completed') + gate.release() + await firstSyncRun + + expect(stripe.customer(org.stripeCustomerId)).toMatchObject({ + email: third.email, + name: third.name, + }) + }) +}) + +describe('Team seat sync', () => { + beforeAll(() => setEnvFlags({ isBillingEnabled: true })) + + it('keeps a seat change that has not reached Stripe when a stale subscription update arrives', async () => { + const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')]) + const org = await createOrganizationWithPlan('team', 1) + await addMember(org.organizationId, owner.id, 'owner') + await addMember(org.organizationId, joiner.id) + await reconcileOrganizationSeats({ organizationId: org.organizationId, reason: 'member-added' }) + expect((await storedSubscription(org.subscriptionId)).seats).toBe(2) + const seatSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + + stripe.updateOutsideSim(org.stripeSubscriptionId, { metadata: { renewedAt: 'period-2' } }) + await deliver(stripe.events.at(-1) as Stripe.Event) + + expect((await storedSubscription(org.subscriptionId)).seats).toBe(2) + await expect(processEvent(seatSync)).resolves.toBe('completed') + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(2) + }) + it('pushes the committed seats when its sync runs between the plugin write and the reconcile', async () => { + const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')]) + const org = await createOrganizationWithPlan('team', 1) + await addMember(org.organizationId, owner.id, 'owner') + await addMember(org.organizationId, joiner.id) + await reconcileOrganizationSeats({ organizationId: org.organizationId, reason: 'member-added' }) + const seatSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + + beforeReconcile = async () => { + beforeReconcile = undefined + await expect(processEvent(seatSync)).resolves.toBe('completed') + } + await deliverUnrelatedUpdate(org.stripeSubscriptionId) + + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(2) + }) + + it('pushes the latest plan when an earlier seat sync lands in Stripe after a newer one', async () => { + const [smallTeam, largeTeam] = CREDIT_TIERS.map((tier) => `team_${tier.credits}`) + const org = await createOrganizationWithPlan('team', 1) + await testDatabase + .update(subscription) + .set({ plan: smallTeam }) + .where(eq(subscription.id, org.subscriptionId)) + async function commitPlan(plan: string) { + await testDatabase.transaction(async (tx) => { + await tx.update(subscription).set({ plan }).where(eq(subscription.id, org.subscriptionId)) + await enqueueSubscriptionSeatsSync(tx, { + subscriptionId: org.subscriptionId, + seats: 1, + reason: 'plan-change', + }) + }) + return latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + } + + const smallSync = await commitPlan(smallTeam) + const stalePush = stripe.holdNextRequest('subscriptions.update') + const pushingSmall = processEvent(smallSync) + await stalePush.reached + const largeSync = await commitPlan(largeTeam) + await expect(processEvent(largeSync)).resolves.toBe('completed') + stalePush.release() + await expect(pushingSmall).resolves.toBe('completed') + + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].price.id).toBe( + 'price_team_max_tier_month' + ) + }) + + it('pushes a seat count that changes away and back while its sync retries', async () => { + const org = await createOrganizationWithPlan('team', 1) + async function commitSeats(seats: number) { + await testDatabase.transaction(async (tx) => { + await tx.update(subscription).set({ seats }).where(eq(subscription.id, org.subscriptionId)) + await enqueueSubscriptionSeatsSync(tx, { + subscriptionId: org.subscriptionId, + seats, + reason: 'member-change', + }) + }) + return latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + } + + const seatSync = await commitSeats(2) + const firstPush = stripe.holdNextRequest('subscriptions.update') + const syncing = processEvent(seatSync) + await firstPush.reached + await commitSeats(3) + const secondPush = stripe.holdNextRequest('subscriptions.update') + firstPush.release() + await secondPush.reached + await commitSeats(2) + secondPush.release() + await expect(syncing).resolves.toBe('pending') + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(3) + + await makeDue(seatSync) + await expect(processEvent(seatSync)).resolves.toBe('completed') + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(2) + }) + + it('drops a seat when a member leaves while the plugin has overwritten the row', async () => { + const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')]) + const org = await createOrganizationWithPlan('team', 1) + await addMember(org.organizationId, owner.id, 'owner') + await addMember(org.organizationId, joiner.id) + await reconcileOrganizationSeats({ organizationId: org.organizationId, reason: 'member-added' }) + const growSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + + beforeReconcile = async () => { + beforeReconcile = undefined + await leaveOrganization(joiner.id, org.organizationId) + await reconcileOrganizationSeats({ + organizationId: org.organizationId, + reason: 'member-removed', + }) + } + await deliverUnrelatedUpdate(org.stripeSubscriptionId) + + expect((await storedSubscription(org.subscriptionId)).seats).toBe(1) + await expect(processEvent(growSync)).resolves.toBe('completed') + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(1) + }) + + it('does not revive an older seat count when its dead-lettered sync is requeued', async () => { + const [owner, joiner] = await Promise.all([createUser('owner'), createUser('joiner')]) + const org = await createOrganizationWithPlan('team', 1) + await addMember(org.organizationId, owner.id, 'owner') + await addMember(org.organizationId, joiner.id) + await reconcileOrganizationSeats({ organizationId: org.organizationId, reason: 'member-added' }) + const growSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + await deadLetter(growSync) + + await leaveOrganization(joiner.id, org.organizationId) + await reconcileOrganizationSeats({ + organizationId: org.organizationId, + reason: 'member-removed', + }) + const shrinkSync = await latestOutboxEventId( + OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS, + org.subscriptionId + ) + await expect(processEvent(shrinkSync)).resolves.toBe('completed') + + await requeueFromAdminApi(growSync) + await deliverUnrelatedUpdate(org.stripeSubscriptionId) + expect((await storedSubscription(org.subscriptionId)).seats).toBe(1) + await expect(processEvent(growSync)).resolves.toBe('completed') + expect(stripe.subscription(org.stripeSubscriptionId).items.data[0].quantity).toBe(1) + }) +}) + +describe('webhook reconcile cost', () => { + interface QueryPlan { + 'Node Type': string + 'Relation Name'?: string + 'Shared Hit Blocks': number + 'Shared Read Blocks': number + 'Actual Rows': number + Plans?: QueryPlan[] + } + + /** EXPLAIN ANALYZE inside a rolled-back transaction, so a measured UPDATE changes nothing. */ + async function explainWithoutEffects(query: string, parameters: unknown[]) { + const rollback = new Error('rollback') + let plan: QueryPlan | undefined + await connection + .begin(async (sql) => { + const [explained] = await sql.unsafe( + `EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${query}`, + parameters as never[] + ) + plan = (explained['QUERY PLAN'] as { Plan: QueryPlan }[])[0].Plan + throw rollback + }) + .catch((error: unknown) => { + if (error !== rollback) throw error + }) + if (!plan) throw new Error(`No plan for ${query}`) + return plan + } + const planNodes = (plan: QueryPlan): QueryPlan[] => [ + plan, + ...(plan.Plans ?? []).flatMap(planNodes), + ] + + it('reads only in-flight syncs, however many have completed or dead-lettered', async () => { + const pro = await createProUserInPaidOrganization() + await pauseProSubscriptionForOrgCoverage(pro.userId) + for (const [status, count] of [ + ['completed', 20000], + ['dead_letter', 50], + ] as const) { + await connection` + INSERT INTO outbox_event (id, event_type, payload, status, available_at, created_at, processed_at) + SELECT ${generateId()} || ':' || n, ${OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END}, + json_build_object( + 'subscriptionId', ${pro.subscriptionId}::text, + 'cancelAtPeriodEnd', false, + 'committedAt', n + ), + ${status}, now(), now(), now() + FROM generate_series(1, ${count}::integer) AS n` + } + await connection`ANALYZE outbox_event` + + const issued: { query: string; parameters: unknown[] }[] = [] + const traced = postgres( + readTestDatabaseUrl(), + withUtcTimestamps({ + max: 2, + prepare: false, + fetch_types: false, + connection: { search_path: schemaName }, + onnotice: () => {}, + debug: (_connection: number, query: string, parameters: unknown[]) => { + issued.push({ query, parameters }) + }, + }) + ) + database.current = drizzle(traced, { schema }) + try { + await deliverUnrelatedUpdate(pro.stripeSubscriptionId) + stripe.updateOutsideSim(pro.stripeSubscriptionId, { cancel_at_period_end: true }) + await deliver(stripe.events.at(-1) as Stripe.Event) + } finally { + database.current = testDatabase + await traced.end() + } + expect((await storedSubscription(pro.subscriptionId)).cancelAtPeriodEnd).toBe(true) + expect(new Set(await cancelValuesOfRetryableSyncs(pro.subscriptionId))).toEqual(new Set([true])) + + const outboxQueries = issued.filter(({ query }) => /"outbox_event"/i.test(query)) + expect(outboxQueries.some(({ query }) => /^update/i.test(query))).toBe(true) + for (const { query, parameters } of outboxQueries) { + const plan = await explainWithoutEffects(query, parameters) + expect( + planNodes(plan).some( + (node) => node['Node Type'] === 'Seq Scan' && node['Relation Name'] === 'outbox_event' + ) + ).toBe(false) + expect(plan['Shared Hit Blocks'] + plan['Shared Read Blocks']).toBeLessThan(200) + if (/^select/i.test(query)) expect(plan['Actual Rows']).toBeLessThanOrEqual(1) + } + }) +}) diff --git a/apps/sim/lib/billing/webhooks/subscription-sync.ts b/apps/sim/lib/billing/webhooks/subscription-sync.ts new file mode 100644 index 00000000000..86d2f2d87f7 --- /dev/null +++ b/apps/sim/lib/billing/webhooks/subscription-sync.ts @@ -0,0 +1,549 @@ +import { db } from '@sim/db' +import { subscription } from '@sim/db/schema' +import { createLogger } from '@sim/logger' +import { generateShortId } from '@sim/utils/id' +import { toRecord } from '@sim/utils/object' +import { eq, sql } from 'drizzle-orm' +import type Stripe from 'stripe' +import { requireStripeClient } from '@/lib/billing/stripe-client' +import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events' +import { + enqueueOutboxEvent, + listInflightOutboxEvents, + patchRetryableOutboxEvents, + readOutboxEventPayload, +} from '@/lib/core/outbox/service' +import type { DbOrTx } from '@/lib/db/types' + +const logger = createLogger('BillingSubscriptionSync') + +/** + * Prefix of the idempotency key on every `cancel_at_period_end` write the sync handler sends. + * Stripe copies the key onto the resulting event's `request.idempotency_key`, which is how a + * webhook is recognized as the echo of Sim's own sync rather than a change made in Stripe. + */ +const CANCEL_AT_PERIOD_END_SYNC_KEY_PREFIX = 'outbox-sync-cancel-at-period-end:' +/** The key prefix the cancel-sync handler used before the one above; old pods still send it mid-rollout. */ +const LEGACY_SYNC_KEY_PREFIX = 'outbox:' + +const CANCEL_SYNC = OUTBOX_EVENT_TYPES.STRIPE_SYNC_CANCEL_AT_PERIOD_END +const SEATS_SYNC = OUTBOX_EVENT_TYPES.STRIPE_SYNC_SUBSCRIPTION_SEATS + +export interface CancelAtPeriodEndSyncPayload { + stripeSubscriptionId: string + /** The DB subscription row id; the handler pushes this row's current value. */ + subscriptionId: string + /** The latest committed value. Absent on events enqueued before it was recorded. */ + cancelAtPeriodEnd?: boolean + /** When `cancelAtPeriodEnd` was committed; see `withCommittedAt`. */ + committedAt?: number + /** Reason this was enqueued, e.g. 'joined-paid-org'. */ + reason?: string + /** Correlates Enterprise-issuance follow-up work for Admin progress/retry. */ + sourceOperationId?: string + operationId?: string + organizationId?: string + requestedBy?: { id: string | null; name: string; email: string | null } +} + +export interface SubscriptionSeatsSyncPayload { + /** The DB subscription row id; the handler pushes this row's current plan and seats. */ + subscriptionId: string + /** The latest committed seat count. Absent on events enqueued before it was recorded. */ + seats?: number + /** When `seats` was committed; see `withCommittedAt`. */ + committedAt?: number + reason?: string +} + +type SubscriptionSyncEventType = typeof CANCEL_SYNC | typeof SEATS_SYNC +type SyncIntentFields = { cancelAtPeriodEnd: boolean } | { seats: number } + +export function isSubscriptionSyncEventType( + eventType: string +): eventType is SubscriptionSyncEventType { + return eventType === CANCEL_SYNC || eventType === SEATS_SYNC +} + +function subscriptionSubject(subscriptionId: string) { + return { payloadKey: 'subscriptionId', payloadValue: subscriptionId } +} + +/** + * Stamps `fields` with the database clock in microseconds since the epoch. Read while the + * caller holds the subscription row lock, so a later stamp is a later committed value. The + * transaction start time (`created_at`) cannot order them: a transaction that began earlier + * can take the row lock later. + */ +async function withCommittedAt( + tx: DbOrTx, + fields: T +): Promise { + return { ...fields, committedAt: await readDatabaseClock(tx) } +} + +/** The clock `committedAt` is stamped from, in microseconds since the epoch. */ +async function readDatabaseClock(executor: DbOrTx): Promise { + const [row] = await executor.execute<{ now: string }>( + sql`select (extract(epoch from clock_timestamp()) * 1000000)::bigint::text as "now"` + ) + if (!row) throw new Error('Database clock read returned no row') + return Number(row.now) +} + +/** + * Records `fields` as the subscription's latest committed value for `eventType` and writes it + * onto every event of that type that can still run: pending, processing, or dead-lettered (each + * operator retry path resets dead letters to pending). No event that can run again ever carries + * an older value for the webhook reconcile to restore. The caller must hold the subscription row + * lock (`FOR UPDATE`, or the `UPDATE` itself), per the lock order on + * {@link lockSubscriptionForSyncRetry}. + * + * A Sim commit is stamped with the clock under that lock and rewrites every such event. A value + * taken from Stripe passes `observedAt`, the clock read just before Stripe was read, so it orders + * by when it was observed: it rewrites (value and stamp) only the events last stamped before that + * observation, including ones already holding the value, so a slower reconcile of an earlier + * Stripe read cannot outrank a later one and never overwrites a newer commit. + */ +async function commitIntent( + tx: DbOrTx, + eventType: SubscriptionSyncEventType, + subscriptionId: string, + fields: T, + observedAt?: number +): Promise { + const committed = + observedAt === undefined + ? await withCommittedAt(tx, fields) + : { ...fields, committedAt: observedAt } + await patchRetryableOutboxEvents( + tx, + eventType, + subscriptionSubject(subscriptionId), + committed, + observedAt === undefined ? undefined : 'committedAt' + ) + return committed +} + +/** + * Enqueue the Stripe sync for a `cancelAtPeriodEnd` value written in this transaction. The + * caller must hold the subscription row lock. Once every in-flight event for the subscription + * completes, Stripe holds the last committed value. + */ +export async function enqueueCancelAtPeriodEndSync( + tx: DbOrTx, + payload: Omit & { + cancelAtPeriodEnd: boolean + } +): Promise { + const committed = await commitIntent(tx, CANCEL_SYNC, payload.subscriptionId, { + cancelAtPeriodEnd: payload.cancelAtPeriodEnd, + }) + const intent: CancelAtPeriodEndSyncPayload = { ...payload, ...committed } + return enqueueOutboxEvent(tx, CANCEL_SYNC, intent) +} + +/** + * Enqueue the Stripe sync for a Team subscription's plan and `seats` as written in this + * transaction. The caller must hold the subscription row lock. + */ +export async function enqueueSubscriptionSeatsSync( + tx: DbOrTx, + payload: Omit & { seats: number } +): Promise { + const committed = await commitIntent(tx, SEATS_SYNC, payload.subscriptionId, { + seats: payload.seats, + }) + const intent: SubscriptionSeatsSyncPayload = { ...payload, ...committed } + return enqueueOutboxEvent(tx, SEATS_SYNC, intent) +} + +/** + * Records a `cancelAtPeriodEnd` value written in this transaction without enqueuing a sync, for a + * writer whose value an existing sync will push or Stripe already holds. It is written onto every + * sync that can still run (including one just reset to pending), so none keeps an older value. + * The caller must hold the subscription row lock. + */ +export async function recordCancelAtPeriodEnd( + tx: DbOrTx, + subscriptionId: string, + cancelAtPeriodEnd: boolean +): Promise { + await commitIntent(tx, CANCEL_SYNC, subscriptionId, { cancelAtPeriodEnd }) +} + +/** + * Takes the subscription row lock for an operator retry of one of its sync events; call it + * before touching the event, then {@link recommitSubscriptionSync} after resetting it. + * + * Lock order for every writer of a subscription's synced fields and their outbox events: + * organization mutation lock (where taken) → subscription row → outbox rows. Committing a value + * rewrites the subscription's retryable sync events, dead letters included, so a retry that + * locked a dead-lettered event before the subscription would deadlock against any concurrent + * writer. + */ +export async function lockSubscriptionForSyncRetry( + tx: DbOrTx, + subscriptionId: string +): Promise { + await tx + .select({ id: subscription.id }) + .from(subscription) + .where(eq(subscription.id, subscriptionId)) + .for('update') + .limit(1) +} + +/** + * Re-commits the latest committed value onto the subscription's sync events that can still run, + * for a dead-lettered event that was just reset to `pending`: the retry then carries the latest + * value rather than the one it failed with. That is the newest in-flight value; the row is only a + * fallback when nothing in flight records one, because until the reconcile step runs the row can + * hold the Stripe plugin's stale webhook payload. The caller holds the lock from + * {@link lockSubscriptionForSyncRetry}, taken before the reset. + */ +export async function recommitSubscriptionSync( + tx: DbOrTx, + eventType: SubscriptionSyncEventType, + subscriptionId: string +): Promise { + const [current] = await tx + .select({ cancelAtPeriodEnd: subscription.cancelAtPeriodEnd, seats: subscription.seats }) + .from(subscription) + .where(eq(subscription.id, subscriptionId)) + .limit(1) + if (!current) return + + if (eventType === CANCEL_SYNC) { + await commitIntent(tx, eventType, subscriptionId, { + cancelAtPeriodEnd: await readCommittedCancelAtPeriodEnd( + tx, + subscriptionId, + Boolean(current.cancelAtPeriodEnd) + ), + }) + return + } + await commitIntent(tx, eventType, subscriptionId, { + seats: await readCommittedSeats(tx, subscriptionId, current.seats ?? 1), + }) +} + +/** + * Records a customer's restore through Better Auth's `/subscription/restore` as Sim's latest + * committed `cancelAtPeriodEnd`. That endpoint updates Stripe and then writes the row directly, + * so a cancel sync still in flight would otherwise carry the cancellation the customer just + * undid, and a webhook processed before the restore's own could restore it for that sync to + * push. Enqueuing re-pushes `false` even if such a sync already ran. + */ +async function commitCustomerRestoredSubscription(restored: unknown): Promise { + const stripeSubscription = toRecord(restored) + const stripeSubscriptionId = stripeSubscription.id + if ( + typeof stripeSubscriptionId !== 'string' || + stripeSubscription.cancel_at_period_end !== false + ) { + return + } + + await db.transaction(async (tx) => { + const [row] = await tx + .select({ id: subscription.id }) + .from(subscription) + .where(eq(subscription.stripeSubscriptionId, stripeSubscriptionId)) + .for('update') + .limit(1) + if (!row) return + + await tx + .update(subscription) + .set({ cancelAtPeriodEnd: false }) + .where(eq(subscription.id, row.id)) + await enqueueCancelAtPeriodEndSync(tx, { + stripeSubscriptionId, + subscriptionId: row.id, + cancelAtPeriodEnd: false, + reason: 'customer-restored', + }) + }) +} + +/** + * The Better Auth `after` hook step for `/subscription/restore`: records the restored value from + * the Stripe subscription the endpoint returned. A failure is logged rather than failing a + * restore that already reached Stripe and the row; the reconcile step then still treats the + * restore's own webhook as a change made in Stripe. + */ +export async function recordCustomerRestoreAfterHook(ctx: { + path: string + context: { returned?: unknown } +}): Promise { + if (ctx.path !== '/subscription/restore') return + try { + await commitCustomerRestoredSubscription(ctx.context.returned) + } catch (error) { + logger.error('Failed to record a restored subscription as the committed value', { error }) + } +} + +/** + * The value recorded on a sync event as of now, not as of its claim: every commit rewrites it on + * each sync that can still run. Undefined for an event enqueued before values were recorded. + */ +export async function readRecordedSyncValue( + eventId: string +): Promise<{ cancelAtPeriodEnd?: boolean; seats?: number } | undefined> { + const payload = toRecord(await readOutboxEventPayload(eventId)) + if (typeof payload.committedAt !== 'number') return undefined + return { + ...(typeof payload.cancelAtPeriodEnd === 'boolean' + ? { cancelAtPeriodEnd: payload.cancelAtPeriodEnd } + : {}), + ...(typeof payload.seats === 'number' ? { seats: payload.seats } : {}), + } +} + +/** + * The subscription's latest committed `cancelAtPeriodEnd`: the newest value an in-flight sync + * records, else `stored` (the row). Until the reconcile step runs, the row can hold the Stripe + * plugin's stale webhook payload, so a writer deciding whether a change is needed compares + * against this, never the row alone. The caller holds the subscription row lock. + */ +async function readCommittedCancelAtPeriodEnd( + tx: DbOrTx, + subscriptionId: string, + stored: boolean +): Promise { + const intent = (await readSyncIntents(tx, subscriptionId)).cancelAtPeriodEnd + return intent.status === 'value' ? intent.value : stored +} + +/** + * True when both the row and the latest committed `cancelAtPeriodEnd` already hold `desired`, so + * a writer has nothing to record. A writer that sees either one differ writes the row and + * commits: a redundant sync of the same value is harmless, while skipping on the committed value + * alone could let an older Stripe read, accepted after this writer, override it. + */ +export async function isCancelAtPeriodEndSettled( + tx: DbOrTx, + subscriptionId: string, + stored: boolean, + desired: boolean +): Promise { + return ( + stored === desired && + (await readCommittedCancelAtPeriodEnd(tx, subscriptionId, stored)) === desired + ) +} + +/** The seat-count counterpart of {@link readCommittedCancelAtPeriodEnd}. */ +export async function readCommittedSeats( + tx: DbOrTx, + subscriptionId: string, + stored: number +): Promise { + const intent = (await readSyncIntents(tx, subscriptionId)).seats + return intent.status === 'value' ? intent.value : stored +} + +/** A fresh key per Stripe write: the SDK reuses it across its own network retries of that call. */ +export function cancelAtPeriodEndSyncIdempotencyKey(eventId: string): string { + return `${CANCEL_AT_PERIOD_END_SYNC_KEY_PREFIX}${eventId}:${generateShortId()}` +} + +/** + * What a sync type's in-flight events say about its field: nothing in flight, the latest + * committed value, or `legacy` when an event predates recorded values (enqueued by an older + * deploy) so the committed value is unknown. + */ +type InflightIntent = + | { status: 'none' } + | { status: 'legacy' } + | { status: 'value'; value: T; committedAt: number } + +function latestIntent( + events: { eventType: string; payload: unknown }[], + eventType: SubscriptionSyncEventType, + readValue: (payload: Record) => T | undefined +): InflightIntent { + let latest: { committedAt: number; value: T } | undefined + for (const event of events) { + if (event.eventType !== eventType) continue + const payload = toRecord(event.payload) + const value = readValue(payload) + if (typeof payload.committedAt !== 'number' || value === undefined) return { status: 'legacy' } + if (!latest || payload.committedAt > latest.committedAt) { + latest = { committedAt: payload.committedAt, value } + } + } + return latest ? { status: 'value', ...latest } : { status: 'none' } +} + +/** + * One read of the subscription's in-flight syncs, through the status index and bounded by the + * in-flight backlog. Dead letters are not intents: they are failed syncs awaiting an operator, + * kept current by every commit so a retry pushes the latest value, but never a reason to override + * Stripe. + */ +async function readSyncIntents(executor: DbOrTx, subscriptionId: string) { + const events = await listInflightOutboxEvents( + executor, + [CANCEL_SYNC, SEATS_SYNC], + subscriptionSubject(subscriptionId) + ) + return { + cancelAtPeriodEnd: latestIntent(events, CANCEL_SYNC, (payload) => + typeof payload.cancelAtPeriodEnd === 'boolean' ? payload.cancelAtPeriodEnd : undefined + ), + seats: latestIntent(events, SEATS_SYNC, (payload) => + typeof payload.seats === 'number' ? payload.seats : undefined + ), + } +} + +/** + * True when the event records a cancellation change made in Stripe, not by Sim's sync. A + * `cancel_at` that was set or cleared counts too: Better Auth's restore clears `cancel_at` when it + * is set, and Stripe may then list only `cancel_at` among the previous attributes. A `cancel_at` + * that only moved (e.g. a billing-interval switch on a subscription already ending) is not a + * cancellation change. + */ +function isCancellationChangedInStripe(event: Stripe.Event): boolean { + const previousAttributes = toRecord(event.data.previous_attributes) + const scheduledOrCleared = + 'cancel_at' in previousAttributes && + (previousAttributes.cancel_at == null) !== (toRecord(event.data.object).cancel_at == null) + if (!('cancel_at_period_end' in previousAttributes) && !scheduledOrCleared) return false + const idempotencyKey = event.request?.idempotency_key + const issuedBySimSync = + idempotencyKey?.startsWith(CANCEL_AT_PERIOD_END_SYNC_KEY_PREFIX) || + idempotencyKey?.startsWith(LEGACY_SYNC_KEY_PREFIX) + return !issuedBySimSync +} + +type CancelAtPeriodEndSource = + | { source: 'unchanged' } + | { source: 'pending-sync'; value: boolean } + | { source: 'stripe' } + +/** + * A Stripe-side change wins over a pending value unless the pending value was committed after + * Stripe was read (`liveReadAt`, on the same database clock): that read predates Sim's newer + * commit, so applying it would roll the newer value back. + */ +function cancelAtPeriodEndSource( + intent: InflightIntent, + changedInStripe: boolean, + liveReadAt?: number +): CancelAtPeriodEndSource { + if (intent.status === 'legacy') return { source: 'unchanged' } + if ( + intent.status === 'value' && + (!changedInStripe || (liveReadAt !== undefined && intent.committedAt > liveReadAt)) + ) { + return { source: 'pending-sync', value: intent.value } + } + return { source: 'stripe' } +} + +/** + * Reconciles the Sim-owned subscription fields after the Better Auth Stripe plugin has copied a + * `customer.subscription.updated` payload into the row. The plugin writes `cancelAtPeriodEnd` + * and `seats` unconditionally, so a delayed, out-of-order, or unrelated event would otherwise + * overwrite a value Sim committed but has not yet pushed, and the pending sync would then push + * the overwritten value back to Stripe. + * + * Decided under the subscription row lock that every committing writer holds: + * - `cancelAtPeriodEnd`: while a cancel sync is in flight, its committed value wins over + * snapshots and over echoes of Sim's own writes. A change made in Stripe itself (customer + * portal, dashboard, Better Auth's cancel/restore endpoints), recognised by a non-Sim request + * changing `cancel_at_period_end` or setting or clearing `cancel_at`, wins and is committed + * onto every sync that can still run, unless Sim committed a newer value after Stripe was + * read. With no sync in flight Stripe wins, read live so out-of-order delivery cannot regress + * it. + * - Precedence across the two systems is arrival order, not wall-clock order: a Stripe-side + * change whose webhook is processed after a Sim commit wins even if the customer made it + * earlier. Stripe's `event.created` is not compared with the database clock, because skew + * between them could override a genuinely newer customer action. + * - `seats`: Team seats are Sim-owned; while a seat sync is in flight its committed value wins. + * - A field with an in-flight event from an older deploy is left as the plugin wrote it. + * + * Stripe is read only when its value decides, and never under the lock. Only the DB row and + * sync payloads are written, so this cannot trigger another webhook. + */ +export async function reconcileSubscriptionSyncFromStripe(event: Stripe.Event): Promise { + if (event.type !== 'customer.subscription.updated') return + const stripeSubscriptionId = event.data.object.id + + const [row] = await db + .select({ id: subscription.id }) + .from(subscription) + .where(eq(subscription.stripeSubscriptionId, stripeSubscriptionId)) + .limit(1) + if (!row) return + + const changedInStripe = isCancellationChangedInStripe(event) + let liveCancelAtPeriodEnd: boolean | undefined + let liveReadAt: number | undefined + + for (let pass = 1; pass <= 2; pass++) { + if (liveCancelAtPeriodEnd === undefined) { + const needsStripe = + pass > 1 || + cancelAtPeriodEndSource( + (await readSyncIntents(db, row.id)).cancelAtPeriodEnd, + changedInStripe + ).source === 'stripe' + if (needsStripe) { + liveReadAt = await readDatabaseClock(db) + const live = await requireStripeClient().subscriptions.retrieve(stripeSubscriptionId) + liveCancelAtPeriodEnd = Boolean(live.cancel_at_period_end) + } + } + + const reconciled = await db.transaction(async (tx) => { + const [current] = await tx + .select({ cancelAtPeriodEnd: subscription.cancelAtPeriodEnd, seats: subscription.seats }) + .from(subscription) + .where(eq(subscription.id, row.id)) + .for('update') + .limit(1) + if (!current) return true + + const intents = await readSyncIntents(tx, row.id) + const cancel = cancelAtPeriodEndSource(intents.cancelAtPeriodEnd, changedInStripe, liveReadAt) + let cancelAtPeriodEnd = Boolean(current.cancelAtPeriodEnd) + if (cancel.source === 'pending-sync') { + cancelAtPeriodEnd = cancel.value + } else if (cancel.source === 'stripe') { + if (liveCancelAtPeriodEnd === undefined) return false + cancelAtPeriodEnd = liveCancelAtPeriodEnd + await commitIntent(tx, CANCEL_SYNC, row.id, { cancelAtPeriodEnd }, liveReadAt) + } + const seats = intents.seats.status === 'value' ? intents.seats.value : current.seats + + if (Boolean(current.cancelAtPeriodEnd) === cancelAtPeriodEnd && current.seats === seats) { + return true + } + await tx + .update(subscription) + .set({ cancelAtPeriodEnd, seats }) + .where(eq(subscription.id, row.id)) + + logger.info('Reconciled Sim-owned subscription fields after a Stripe webhook', { + eventId: event.id, + subscriptionId: row.id, + stripeSubscriptionId, + cancelAtPeriodEnd: { + stored: current.cancelAtPeriodEnd, + reconciled: cancelAtPeriodEnd, + source: cancel.source, + }, + seats: { stored: current.seats, reconciled: seats }, + }) + return true + }) + if (reconciled) return + } +} diff --git a/apps/sim/lib/core/outbox/service.ts b/apps/sim/lib/core/outbox/service.ts index 2990714e39c..eed0495eafa 100644 --- a/apps/sim/lib/core/outbox/service.ts +++ b/apps/sim/lib/core/outbox/service.ts @@ -406,6 +406,90 @@ export async function findDeadLetteredEvents( .limit(DEAD_LETTER_SCAN_LIMIT) } +/** The event's current payload, read fresh rather than from the copy its handler was claimed with. */ +export async function readOutboxEventPayload(eventId: string): Promise { + const [row] = await db + .select({ payload: outboxEvent.payload }) + .from(outboxEvent) + .where(eq(outboxEvent.id, eventId)) + .limit(1) + return row?.payload +} + +/** Statuses of an event whose side effect may still run. */ +const INFLIGHT_OUTBOX_STATUSES = ['pending', 'processing'] as const +/** + * Statuses an event can still run from: in flight, or dead-lettered, which every operator retry + * path resets to `pending`. A `completed` event never runs again. + */ +const RETRYABLE_OUTBOX_STATUSES = [...INFLIGHT_OUTBOX_STATUSES, 'dead_letter'] as const + +/** Identifies the subject of an event by one scalar field of its JSON payload. */ +export interface OutboxPayloadSubject { + payloadKey: string + payloadValue: string +} + +function eventsForSubject( + eventTypes: readonly string[], + subject: OutboxPayloadSubject, + statuses: readonly string[] +) { + return and( + inArray(outboxEvent.eventType, [...eventTypes]), + inArray(outboxEvent.status, [...statuses]), + sql`${outboxEvent.payload} ->> ${subject.payloadKey} = ${subject.payloadValue}` + ) +} + +/** + * The `pending` or `processing` events of the given types for one subject. Pass the caller's + * transaction to read under its locks. + */ +export async function listInflightOutboxEvents( + executor: Pick, + eventTypes: readonly string[], + subject: OutboxPayloadSubject, + limit?: number +): Promise<{ id: string; eventType: string; payload: unknown }[]> { + const query = executor + .select({ id: outboxEvent.id, eventType: outboxEvent.eventType, payload: outboxEvent.payload }) + .from(outboxEvent) + .where(eventsForSubject(eventTypes, subject, INFLIGHT_OUTBOX_STATUSES)) + return limit === undefined ? query : query.limit(limit) +} + +/** + * Shallow-merges `patch` into the payload of every `pending`, `processing`, or `dead_letter` + * event of the type for one subject. With `onlyIfOlderThanPatch`, naming a numeric payload key + * that `patch` sets, an event whose own value for that key is already at least the patch's is + * left alone, so a stale writer never overwrites a newer one. One UPDATE; nothing is read into + * memory. Callers serialize writers for the subject with their domain lock. + */ +export async function patchRetryableOutboxEvents( + executor: Pick, + eventType: string, + subject: OutboxPayloadSubject, + patch: Record, + onlyIfOlderThanPatch?: string +): Promise { + const patched = await executor + .update(outboxEvent) + .set({ + payload: sql`(coalesce(${outboxEvent.payload}::jsonb, '{}'::jsonb) || ${JSON.stringify(patch)}::jsonb)::json`, + }) + .where( + and( + eventsForSubject([eventType], subject, RETRYABLE_OUTBOX_STATUSES), + onlyIfOlderThanPatch + ? sql`coalesce((${outboxEvent.payload} ->> ${onlyIfOlderThanPatch})::numeric, -1) < ${String(patch[onlyIfOlderThanPatch])}::numeric` + : undefined + ) + ) + .returning({ id: outboxEvent.id }) + return patched.length +} + /** * True when an event of the given type whose JSON payload has * `payload->>payloadKey === payloadValue` is still `pending` or `processing`. @@ -417,18 +501,8 @@ export async function hasInflightOutboxEvent( payloadKey: string, payloadValue: string ): Promise { - const [row] = await db - .select({ id: outboxEvent.id }) - .from(outboxEvent) - .where( - and( - eq(outboxEvent.eventType, eventType), - inArray(outboxEvent.status, ['pending', 'processing']), - sql`${outboxEvent.payload} ->> ${payloadKey} = ${payloadValue}` - ) - ) - .limit(1) - return Boolean(row) + const events = await listInflightOutboxEvents(db, [eventType], { payloadKey, payloadValue }, 1) + return events.length > 0 } /** diff --git a/packages/testing/src/mocks/billing-subscription-sync.mock.ts b/packages/testing/src/mocks/billing-subscription-sync.mock.ts new file mode 100644 index 00000000000..492ce5f4a49 --- /dev/null +++ b/packages/testing/src/mocks/billing-subscription-sync.mock.ts @@ -0,0 +1,68 @@ +import { vi } from 'vitest' + +/** + * Controllable mock functions for `@/lib/billing/webhooks/subscription-sync`. The enqueue + * functions resolve to a fixed event id; drive them with `mockResolvedValueOnce`. + * `mockIsSubscriptionSyncEventType` keeps the real logic; `mockReadCommittedSeats` and + * `mockIsCancelAtPeriodEndSettled` answer from the stored value, as when nothing is in flight. + * + * @example + * ```ts + * import { billingSubscriptionSyncMockFns } from '@sim/testing/mocks/billing-subscription-sync.mock' + * + * expect(billingSubscriptionSyncMockFns.mockEnqueueCancelAtPeriodEndSync).toHaveBeenCalledWith( + * expect.anything(), + * expect.objectContaining({ cancelAtPeriodEnd: true }) + * ) + * ``` + */ +export const billingSubscriptionSyncMockFns = { + mockEnqueueCancelAtPeriodEndSync: vi.fn(async () => 'cancel-at-period-end-sync-event'), + mockLockSubscriptionForSyncRetry: vi.fn(async () => undefined), + mockRecommitSubscriptionSync: vi.fn(async () => undefined), + mockRecordCancelAtPeriodEnd: vi.fn(async () => undefined), + mockIsSubscriptionSyncEventType: vi.fn( + (eventType: string) => + eventType === 'stripe.sync-cancel-at-period-end' || + eventType === 'stripe.sync-subscription-seats' + ), + mockEnqueueSubscriptionSeatsSync: vi.fn(async () => 'subscription-seats-sync-event'), + mockCancelAtPeriodEndSyncIdempotencyKey: vi.fn( + (eventId: string) => `outbox-sync-cancel-at-period-end:${eventId}:key` + ), + mockReconcileSubscriptionSyncFromStripe: vi.fn(async () => undefined), + mockRecordCustomerRestoreAfterHook: vi.fn(async () => undefined), + mockReadRecordedSyncValue: vi.fn(async () => undefined), + mockIsCancelAtPeriodEndSettled: vi.fn( + async (_tx: unknown, _subscriptionId: string, stored: boolean, desired: boolean) => + stored === desired + ), + mockReadCommittedSeats: vi.fn( + async (_tx: unknown, _subscriptionId: string, stored: number) => stored + ), +} + +/** + * Static mock module for `@/lib/billing/webhooks/subscription-sync`. + * + * @example + * ```ts + * vi.mock('@/lib/billing/webhooks/subscription-sync', () => billingSubscriptionSyncMock) + * ``` + */ +export const billingSubscriptionSyncMock = { + enqueueCancelAtPeriodEndSync: billingSubscriptionSyncMockFns.mockEnqueueCancelAtPeriodEndSync, + lockSubscriptionForSyncRetry: billingSubscriptionSyncMockFns.mockLockSubscriptionForSyncRetry, + recommitSubscriptionSync: billingSubscriptionSyncMockFns.mockRecommitSubscriptionSync, + recordCancelAtPeriodEnd: billingSubscriptionSyncMockFns.mockRecordCancelAtPeriodEnd, + isSubscriptionSyncEventType: billingSubscriptionSyncMockFns.mockIsSubscriptionSyncEventType, + enqueueSubscriptionSeatsSync: billingSubscriptionSyncMockFns.mockEnqueueSubscriptionSeatsSync, + cancelAtPeriodEndSyncIdempotencyKey: + billingSubscriptionSyncMockFns.mockCancelAtPeriodEndSyncIdempotencyKey, + reconcileSubscriptionSyncFromStripe: + billingSubscriptionSyncMockFns.mockReconcileSubscriptionSyncFromStripe, + recordCustomerRestoreAfterHook: billingSubscriptionSyncMockFns.mockRecordCustomerRestoreAfterHook, + readRecordedSyncValue: billingSubscriptionSyncMockFns.mockReadRecordedSyncValue, + isCancelAtPeriodEndSettled: billingSubscriptionSyncMockFns.mockIsCancelAtPeriodEndSettled, + readCommittedSeats: billingSubscriptionSyncMockFns.mockReadCommittedSeats, +} diff --git a/packages/testing/src/mocks/index.ts b/packages/testing/src/mocks/index.ts index 4332a76b106..f8dd5d0d97a 100644 --- a/packages/testing/src/mocks/index.ts +++ b/packages/testing/src/mocks/index.ts @@ -127,6 +127,10 @@ export { billingSubscriptionMock, billingSubscriptionMockFns, } from './billing-subscription.mock' +export { + billingSubscriptionSyncMock, + billingSubscriptionSyncMockFns, +} from './billing-subscription-sync.mock' export { billingSubscriptionUtilsMock, billingSubscriptionUtilsMockFns, @@ -746,7 +750,9 @@ export { storageServiceMockFns, } from './storage-service.mock' export { + createInMemoryStripe, createMockStripeEvent, + type InMemoryStripe, stripeClientMock, stripePaymentMethodMock, } from './stripe.mock' diff --git a/packages/testing/src/mocks/outbox-service.mock.ts b/packages/testing/src/mocks/outbox-service.mock.ts index 15f80918643..1744ecfb0bd 100644 --- a/packages/testing/src/mocks/outbox-service.mock.ts +++ b/packages/testing/src/mocks/outbox-service.mock.ts @@ -68,6 +68,9 @@ export const outboxServiceMockFns = { ) }), mockFindDeadLetteredEvents: vi.fn(), + mockListInflightOutboxEvents: vi.fn(), + mockReadOutboxEventPayload: vi.fn(), + mockPatchRetryableOutboxEvents: vi.fn(), mockHasInflightOutboxEvent: vi.fn(), mockHasDueOutboxWork: vi.fn(), mockProcessOutboxEvents: vi.fn(), @@ -97,6 +100,9 @@ export const outboxServiceMock = { outboxEventHasSourceOperationId: outboxServiceMockFns.mockOutboxEventHasSourceOperationId, outboxPayloadHasSourceOperationId: outboxServiceMockFns.mockOutboxPayloadHasSourceOperationId, findDeadLetteredEvents: outboxServiceMockFns.mockFindDeadLetteredEvents, + listInflightOutboxEvents: outboxServiceMockFns.mockListInflightOutboxEvents, + readOutboxEventPayload: outboxServiceMockFns.mockReadOutboxEventPayload, + patchRetryableOutboxEvents: outboxServiceMockFns.mockPatchRetryableOutboxEvents, hasInflightOutboxEvent: outboxServiceMockFns.mockHasInflightOutboxEvent, hasDueOutboxWork: outboxServiceMockFns.mockHasDueOutboxWork, processOutboxEvents: outboxServiceMockFns.mockProcessOutboxEvents, diff --git a/packages/testing/src/mocks/stripe.mock.ts b/packages/testing/src/mocks/stripe.mock.ts index 5babfd73b6f..6dfdd11cf2d 100644 --- a/packages/testing/src/mocks/stripe.mock.ts +++ b/packages/testing/src/mocks/stripe.mock.ts @@ -55,3 +55,320 @@ export function createMockStripeEvent( ...overrides, } as Stripe.Event } + +/** A Stripe subscription as the in-memory fake stores it: the fields billing code reads. */ +export interface InMemoryStripeSubscription { + id: string + object: 'subscription' + customer: string + status: Stripe.Subscription.Status + cancel_at_period_end: boolean + cancel_at: number | null + canceled_at: number | null + ended_at: number | null + trial_start: number | null + trial_end: number | null + schedule: string | null + metadata: Record + items: { + object: 'list' + data: Array<{ + id: string + quantity: number + current_period_start: number + current_period_end: number + price: { id: string; recurring: { interval: 'month' | 'year' } } + }> + } +} + +/** A Stripe customer as the in-memory fake stores it. */ +export interface InMemoryStripeCustomer { + id: string + object: 'customer' + email: string | null + name: string | null +} + +/** A request the fake parked on arrival, before Stripe processes it. */ +export interface InMemoryStripeRequestGate { + /** Resolves once the parked request has reached the fake. */ + reached: Promise + /** Lets the parked request proceed to Stripe's processing. */ + release(): void +} + +type UpdatableResource = 'subscriptions' | 'customers' +type StripeOperation = `${UpdatableResource}.${'retrieve' | 'update'}` + +interface SubscriptionUpdateParams { + cancel_at_period_end?: boolean + /** A Unix timestamp schedules the cancellation; `''` clears it. Only `cancel_at` changes. */ + cancel_at?: number | '' + metadata?: Record + items?: Array<{ id: string; quantity?: number; price?: string }> +} + +interface CustomerUpdateParams { + email?: string + name?: string +} + +/** + * An in-memory Stripe account for integration tests that need Stripe to behave like Stripe: + * `retrieve` returns current state, `update` applies params and emits the + * `customer.subscription.updated` event Stripe would send (with `previous_attributes` and the + * originating `request.idempotency_key`), and a reused idempotency key replays the first + * response, or rejects when its parameters differ. A test can park the next request on a gate to + * control the order requests land in (a parked retrieve answers with the state it arrived to), or make the next update apply and then fail on the + * client, as a dropped connection does. + * + * @example + * ```ts + * const stripe = createInMemoryStripe() + * stripeClientMock.requireStripeClient.mockReturnValue(stripe.client) + * const gate = stripe.holdNextRequest('subscriptions.update') + * ``` + */ +export function createInMemoryStripe() { + const subscriptions = new Map() + const customers = new Map() + const idempotentResults = new Map() + const gates = new Map void; released: Promise }>>() + const failuresAfterApply = new Map() + const failuresOnArrival = new Map() + const events: Stripe.Event[] = [] + let sequence = 0 + + function nextId(prefix: string) { + sequence += 1 + return `${prefix}_${sequence}` + } + + function requireSubscription(id: string) { + const subscription = subscriptions.get(id) + if (!subscription) throw new Error(`No such subscription: '${id}'`) + return subscription + } + + function requireCustomer(id: string) { + const customer = customers.get(id) + if (!customer) throw new Error(`No such customer: '${id}'`) + return customer + } + + function applySubscriptionUpdate( + id: string, + params: SubscriptionUpdateParams, + idempotencyKey: string | null + ) { + const current = requireSubscription(id) + const next = structuredClone(current) + const previousAttributes: Record = {} + if ( + params.cancel_at_period_end !== undefined && + params.cancel_at_period_end !== current.cancel_at_period_end + ) { + previousAttributes.cancel_at_period_end = current.cancel_at_period_end + next.cancel_at_period_end = params.cancel_at_period_end + } + if (params.cancel_at !== undefined) { + const cancelAt = params.cancel_at === '' ? null : params.cancel_at + if (cancelAt !== current.cancel_at) { + previousAttributes.cancel_at = current.cancel_at + next.cancel_at = cancelAt + } + } + if (params.metadata) { + const metadata = { ...current.metadata, ...params.metadata } + if (JSON.stringify(metadata) !== JSON.stringify(current.metadata)) { + previousAttributes.metadata = current.metadata + next.metadata = metadata + } + } + for (const item of params.items ?? []) { + const target = next.items.data.find((existing) => existing.id === item.id) + if (!target) throw new Error(`No such subscription item: '${item.id}'`) + if (item.quantity !== undefined) target.quantity = item.quantity + if (item.price !== undefined) target.price = { ...target.price, id: item.price } + } + if (JSON.stringify(next.items) !== JSON.stringify(current.items)) { + previousAttributes.items = current.items + } + subscriptions.set(id, next) + if (Object.keys(previousAttributes).length > 0) { + events.push({ + id: nextId('evt'), + object: 'event', + api_version: '2025-08-27.basil', + created: sequence, + livemode: false, + pending_webhooks: 1, + request: { id: nextId('req'), idempotency_key: idempotencyKey }, + type: 'customer.subscription.updated', + data: { + object: structuredClone(next) as unknown as Stripe.Subscription, + previous_attributes: previousAttributes as Partial, + }, + } as Stripe.Event) + } + return structuredClone(next) + } + + function rejectIfFailing(operation: StripeOperation) { + const failure = failuresOnArrival.get(operation)?.shift() + if (failure) throw failure + } + + async function waitAtGate(operation: StripeOperation) { + const gate = gates.get(operation)?.shift() + if (gate) { + gate.reached() + await gate.released + } + } + + async function retrieve(operation: StripeOperation, read: () => T): Promise { + rejectIfFailing(operation) + const snapshot = structuredClone(read()) + await waitAtGate(operation) + return snapshot + } + + async function update( + resource: UpdatableResource, + id: string, + params: unknown, + options: { idempotencyKey?: string } | undefined, + apply: () => T + ): Promise { + rejectIfFailing(`${resource}.update`) + await waitAtGate(`${resource}.update`) + + const idempotencyKey = options?.idempotencyKey + const fingerprint = JSON.stringify([resource, id, params]) + if (idempotencyKey) { + const previous = idempotentResults.get(idempotencyKey) + if (previous && previous.fingerprint !== fingerprint) { + throw Object.assign( + new Error( + 'Keys for idempotent requests can only be used with the same parameters they were first used with.' + ), + { type: 'StripeIdempotencyError' } + ) + } + if (previous) return structuredClone(previous.result) as T + } + + const result = apply() + if (idempotencyKey) idempotentResults.set(idempotencyKey, { fingerprint, result }) + const failure = failuresAfterApply.get(resource)?.shift() + if (failure) throw failure + return structuredClone(result) + } + + const client = { + subscriptions: { + retrieve: (id: string) => retrieve('subscriptions.retrieve', () => requireSubscription(id)), + update: ( + id: string, + params: SubscriptionUpdateParams, + options?: { idempotencyKey?: string } + ) => + update('subscriptions', id, params, options, () => + applySubscriptionUpdate(id, params, options?.idempotencyKey ?? null) + ), + }, + customers: { + retrieve: (id: string) => retrieve('customers.retrieve', () => requireCustomer(id)), + update: (id: string, params: CustomerUpdateParams, options?: { idempotencyKey?: string }) => + update('customers', id, params, options, () => { + const next = { ...requireCustomer(id), ...params } + customers.set(id, next) + return structuredClone(next) + }), + }, + webhooks: { + constructEventAsync: async (payload: string) => JSON.parse(payload) as Stripe.Event, + }, + } + + return { + /** Pass to `stripeClientMock.requireStripeClient` or to the Better Auth Stripe plugin. */ + client: client as unknown as Stripe, + /** Every `customer.subscription.updated` event Stripe emitted, in the order it applied them. */ + events, + addSubscription( + subscription: Pick & + Partial< + Pick + > & { + quantity?: number + priceId?: string + } + ) { + const now = Math.floor(Date.now() / 1000) + subscriptions.set(subscription.id, { + id: subscription.id, + object: 'subscription', + customer: subscription.customer, + status: subscription.status ?? 'active', + cancel_at_period_end: subscription.cancel_at_period_end ?? false, + cancel_at: subscription.cancel_at ?? null, + canceled_at: null, + ended_at: null, + trial_start: null, + trial_end: null, + schedule: null, + metadata: {}, + items: { + object: 'list', + data: [ + { + id: `si_${subscription.id}`, + quantity: subscription.quantity ?? 1, + current_period_start: now, + current_period_end: now + 30 * 24 * 60 * 60, + price: { + id: subscription.priceId ?? `price_${subscription.id}`, + recurring: { interval: 'month' }, + }, + }, + ], + }, + }) + }, + addCustomer(customer: Pick) { + customers.set(customer.id, { ...customer, object: 'customer' }) + }, + subscription: (id: string) => structuredClone(requireSubscription(id)), + customer: (id: string) => structuredClone(requireCustomer(id)), + /** Applies a change made outside Sim (dashboard, customer portal) and emits its event. */ + updateOutsideSim(id: string, params: SubscriptionUpdateParams) { + return applySubscriptionUpdate(id, params, null) + }, + /** Parks the next call to `operation` until the returned gate is released. */ + holdNextRequest(operation: StripeOperation): InMemoryStripeRequestGate { + let reached: () => void = () => {} + const arrival = new Promise((resolve) => { + reached = resolve + }) + let release: () => void = () => {} + const released = new Promise((resolve) => { + release = resolve + }) + gates.set(operation, [...(gates.get(operation) ?? []), { reached, released }]) + return { reached: arrival, release } + }, + /** Makes the next call to `operation` fail before Stripe processes it, as an outage does. */ + failNextRequest(operation: StripeOperation, error = new Error('Stripe is unavailable')) { + failuresOnArrival.set(operation, [...(failuresOnArrival.get(operation) ?? []), error]) + }, + /** Makes the next update to `resource` apply in Stripe, then fail on the client. */ + failNextUpdateAfterApplying(resource: UpdatableResource, error = new Error('socket hang up')) { + failuresAfterApply.set(resource, [...(failuresAfterApply.get(resource) ?? []), error]) + }, + } +} + +export type InMemoryStripe = ReturnType