Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 46 additions & 18 deletions apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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'

Expand All @@ -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) => {
Expand Down Expand Up @@ -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(
Expand Down
9 changes: 5 additions & 4 deletions apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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',
})
})
Expand Down
5 changes: 5 additions & 0 deletions apps/sim/lib/admin/subscription-lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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,
Expand Down Expand Up @@ -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',
Expand Down Expand Up @@ -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',
Expand All @@ -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',
Expand Down
50 changes: 36 additions & 14 deletions apps/sim/lib/admin/subscription-lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<string | null>`${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({
Expand All @@ -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')
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
7 changes: 7 additions & 0 deletions apps/sim/lib/auth/auth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1727,6 +1733,7 @@ export const auth = betterAuth({
case 'customer.subscription.created':
case 'customer.subscription.updated': {
await handleManualEnterpriseSubscription(event)
await reconcileSubscriptionSyncFromStripe(event)
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
await handleSubscriptionUsageUpdate(event)
break
}
Expand Down
18 changes: 17 additions & 1 deletion apps/sim/lib/billing/enterprise-provisioning.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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')
}
Expand All @@ -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
})

Expand Down
2 changes: 2 additions & 0 deletions apps/sim/lib/billing/organizations/lock-order.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand All @@ -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
Expand Down
Loading
Loading