diff --git a/bun.lock b/bun.lock index f5ecec2f9..5f80b89c7 100644 --- a/bun.lock +++ b/bun.lock @@ -253,6 +253,14 @@ "typescript-eslint": "^8.26.0", }, }, + "packages/stripe-sync": { + "name": "@autumn/stripe-sync", + "version": "1.0.0", + "dependencies": { + "@supabase/stripe-sync-engine": "^0.48.5", + "stripe": "catalog:", + }, + }, "scripts": { "name": "@autumn/scripts", "version": "1.0.0", @@ -282,6 +290,7 @@ "@anthropic-ai/sdk": "^0.32.1", "@autumn/ksuid": "workspace:*", "@autumn/shared": "workspace:*", + "@autumn/stripe-sync": "workspace:*", "@aws-sdk/client-s3": "^3.1017.0", "@aws-sdk/client-scheduler": "^3.1004.0", "@aws-sdk/client-sqs": "^3.958.0", @@ -603,6 +612,8 @@ "@autumn/shared": ["@autumn/shared@workspace:shared"], + "@autumn/stripe-sync": ["@autumn/stripe-sync@workspace:packages/stripe-sync"], + "@autumn/vite": ["@autumn/vite@workspace:vite"], "@aws-crypto/crc32": ["@aws-crypto/crc32@5.2.0", "", { "dependencies": { "@aws-crypto/util": "^5.2.0", "@aws-sdk/types": "^3.222.0", "tslib": "^2.6.2" } }, "sha512-nLbCWqQNgUiwwtFsen1AdzAtvuLRsQS8rYgMuxCrdKf9kOssamGLuPwyTY9wyYblNr9+1XM8v6zoDTPPSIeANg=="], @@ -1921,6 +1932,8 @@ "@supabase/storage-js": ["@supabase/storage-js@2.102.1", "", { "dependencies": { "iceberg-js": "^0.8.1", "tslib": "2.8.1" } }, "sha512-eCL9T4Xpe40nmKlkUJ7Zq/hk34db1xPiT0WL3Iv5MbJqHuCAe5TxhV8Rjqd6DNZrzjtfYObZtYl9jKJaHrivqw=="], + "@supabase/stripe-sync-engine": ["@supabase/stripe-sync-engine@0.48.5", "", { "dependencies": { "pg": "^8.20.0", "pg-node-migrations": "0.0.8", "yesql": "^7.0.0" }, "peerDependencies": { "stripe": "> 18" } }, "sha512-+LbtJH8n5Xiu289AL3FuWFdKXd0K7kDF0z4Lm+zMYoImWmOuGd3TgSx9gm/nv4nzLooOmIxGZh6LojoYBcJM+g=="], + "@supabase/supabase-js": ["@supabase/supabase-js@2.102.1", "", { "dependencies": { "@supabase/auth-js": "2.102.1", "@supabase/functions-js": "2.102.1", "@supabase/postgrest-js": "2.102.1", "@supabase/realtime-js": "2.102.1", "@supabase/storage-js": "2.102.1" } }, "sha512-bChxPVeLDnYN9M2d/u4fXsvylwSQG5grAl+HN8f+ZD9a9PuVU+Ru+xGmEsk+b9Iz3rJC9ZQnQUJYQ28fApdWYA=="], "@swc/helpers": ["@swc/helpers@0.5.15", "", { "dependencies": { "tslib": "^2.8.0" } }, "sha512-JQ5TuMi45Owi4/BIMAJBoSQoOJu12oOk/gADqlcUL9JEdHB8vyjUSsxqeNXnmXHjYKMi2WcYtezGEEhqUI/E2g=="], @@ -4369,6 +4382,8 @@ "pg-int8": ["pg-int8@1.0.1", "", {}, "sha512-WCtabS6t3c8SkpDBUlb1kjOs7l66xsGdKpIPZsg4wR+B3+u9UAum2odSsF9tnvxg80h4ZxLWMy4pRjOsFIqQpw=="], + "pg-node-migrations": ["pg-node-migrations@0.0.8", "", { "dependencies": { "pg": "^8.6.0", "sql-template-strings": "^2.2.2" }, "bin": { "pg-validate-migrations": "dist/bin/validate.js" } }, "sha512-44cMl9umOmCv0hzZyEcvjEq8Bm8u7mrzggZ06qXTJVSsMMB4j2OsjG+rSp+uzeKWyP2Vu0K9Ye2wKtjFUJwrdw=="], + "pg-pool": ["pg-pool@3.13.0", "", { "peerDependencies": { "pg": ">=8.0" } }, "sha512-gB+R+Xud1gLFuRD/QgOIgGOBE2KCQPaPwkzBBGC9oG69pHTkhQeIuejVIk3/cnDyX39av2AxomQiyPT13WKHQA=="], "pg-protocol": ["pg-protocol@1.13.0", "", {}, "sha512-zzdvXfS6v89r6v7OcFCHfHlyG/wvry1ALxZo4LqgUoy7W9xhBDMaqOuMiF3qEV45VqsN6rdlcehHrfDtlCPc8w=="], @@ -4867,6 +4882,8 @@ "sprintf-js": ["sprintf-js@1.0.3", "", {}, "sha512-D9cPgkvLlV3t3IzL0D0YLvGA9Ahk4PcvVwUbN0dSGr1aP0Nrt4AEnTUbuGvquEC0mA64Gqt1fzirlRs5ibXx8g=="], + "sql-template-strings": ["sql-template-strings@2.2.2", "", {}, "sha512-UXhXR2869FQaD+GMly8jAMCRZ94nU5KcrFetZfWEMd+LVVG6y0ExgHAhatEcKZ/wk8YcKPdi+hiD2wm75lq3/Q=="], + "sqs-consumer": ["sqs-consumer@6.0.2", "", { "dependencies": { "@aws-sdk/client-sqs": "^3.226.0", "debug": "^4.3.4" } }, "sha512-Y8ztFBc1VPj4q72j9Uji4XR6n2O88n9X83aOUaCvz9zZqjpGJaI4PBHBPnE8RDBgD0EQR8Dh+sQwH/+8mzr+Bg=="], "stack-utils": ["stack-utils@2.0.6", "", { "dependencies": { "escape-string-regexp": "^2.0.0" } }, "sha512-XlkWvfIm6RmsWtNJx+uqtKLS8eqFbxUg0ZzLXqY0caEy9l7hruX8IpiDnjsLavoBgqCCR71TqWO8MaXYheJ3RQ=="], @@ -5377,6 +5394,8 @@ "yauzl": ["yauzl@2.10.0", "", { "dependencies": { "buffer-crc32": "~0.2.3", "fd-slicer": "~1.1.0" } }, "sha512-p4a9I6X6nu6IhoGmBqAcbJy1mlC4j27vEPZX9F4L4/vZT3Lyq1VkFHw/V/PUcB9Buo+DG3iHkT0x3Qya58zc3g=="], + "yesql": ["yesql@7.0.0", "", {}, "sha512-sosfr7agy4ibLM7BvXBkM6BpBmKMGuBO8DUYQEuey+QqaqrgW+2bsSg6D050ocBYIz0PuHxUyehyzEztZTU4pw=="], + "yn": ["yn@3.1.1", "", {}, "sha512-Ux4ygGWsu2c7isFWe8Yu1YluJmqVhxqK2cLXNQA5AcC3QfbGNpM7fu0Y8b/z16pXLnFxZYvWhd3fhBY9DLmC6Q=="], "yocto-queue": ["yocto-queue@1.2.2", "", {}, "sha512-4LCcse/U2MHZ63HAJVE+v71o7yOdIe4cZ70Wpf8D/IyjDKYQLV5GD46B+hSTjJsvV5PztjvHoU580EftxjDZFQ=="], diff --git a/package.json b/package.json index 9f05e4477..0d571ce90 100644 --- a/package.json +++ b/package.json @@ -15,7 +15,8 @@ "packages/sdk", "packages/autumn-js", "packages/openapi", - "packages/ksuid" + "packages/ksuid", + "packages/stripe-sync" ], "catalog": { "stripe": "19.3.0-beta.1", @@ -107,7 +108,9 @@ "js:publish": "bun js:build && cd packages/autumn-js && npm publish", "js:publish-beta": "bun js:build && cd packages/autumn-js && npm publish --tag beta", "js:publish-dry": "bun js:build && cd packages/autumn-js && npm publish --dry-run", - "js:version": "git fetch --tags origin && git tag -l 'autumn-js-v*' --sort=-v:refname | head -1" + "js:version": "git fetch --tags origin && git tag -l 'autumn-js-v*' --sort=-v:refname | head -1", + "stripe-sync:migrate": "bun packages/stripe-sync/scripts/migrate.ts", + "stripe-sync:migrate:prod": "NODE_ENV=production infisical run --env=prod -- bun packages/stripe-sync/scripts/migrate.ts" }, "dependencies": { "@aws-sdk/client-sqs": "^3.985.0", diff --git a/packages/stripe-sync/package.json b/packages/stripe-sync/package.json new file mode 100644 index 000000000..5249fb210 --- /dev/null +++ b/packages/stripe-sync/package.json @@ -0,0 +1,20 @@ +{ + "name": "@autumn/stripe-sync", + "version": "1.0.0", + "description": "Stripe-to-Postgres sync engine wrapper for Autumn", + "type": "module", + "main": "./src/index.ts", + "types": "./src/index.ts", + "exports": { + ".": { + "types": "./src/index.ts", + "import": "./src/index.ts", + "default": "./src/index.ts" + } + }, + "private": true, + "dependencies": { + "@supabase/stripe-sync-engine": "^0.48.5", + "stripe": "catalog:" + } +} diff --git a/packages/stripe-sync/scripts/migrate.ts b/packages/stripe-sync/scripts/migrate.ts new file mode 100644 index 000000000..e80874c6d --- /dev/null +++ b/packages/stripe-sync/scripts/migrate.ts @@ -0,0 +1,26 @@ +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { config } from "dotenv"; +import { addAutumnColumns, runStripeSyncMigrations } from "../src/index.js"; + +const scriptDir = path.dirname(fileURLToPath(import.meta.url)); +const repoRoot = path.resolve(scriptDir, "../../.."); + +config({ path: path.join(repoRoot, "server/.env"), override: true }); + +const databaseUrl = process.env.STRIPE_SYNC_DATABASE_URL; + +if (!databaseUrl) { + console.error("STRIPE_SYNC_DATABASE_URL is required"); + process.exit(1); +} + +console.log("Running stripe sync migrations..."); +await runStripeSyncMigrations({ databaseUrl }); +console.log("Stripe sync migrations complete"); + +console.log( + "Adding Autumn columns (stripe_account_id, org_id) to synced tables...", +); +await addAutumnColumns({ databaseUrl }); +console.log("Autumn columns added"); diff --git a/packages/stripe-sync/src/addAccountIdColumn.ts b/packages/stripe-sync/src/addAccountIdColumn.ts new file mode 100644 index 000000000..95d09c968 --- /dev/null +++ b/packages/stripe-sync/src/addAccountIdColumn.ts @@ -0,0 +1,37 @@ +import pg from "pg"; +import { SYNCED_TABLES } from "./eventTypeToTable.js"; + +/** + * Adds `stripe_account_id` and `org_id` columns + indexes to all synced tables. + * Idempotent -- safe to run multiple times. + */ +export const addAutumnColumns = async ({ + databaseUrl, + schema = "stripe", +}: { + databaseUrl: string; + schema?: string; +}): Promise => { + const client = new pg.Client({ connectionString: databaseUrl }); + await client.connect(); + + try { + for (const table of SYNCED_TABLES) { + await client.query(` + ALTER TABLE "${schema}"."${table}" + ADD COLUMN IF NOT EXISTS stripe_account_id TEXT, + ADD COLUMN IF NOT EXISTS org_id TEXT + `); + await client.query(` + CREATE INDEX IF NOT EXISTS idx_${table}_stripe_account_id + ON "${schema}"."${table}" (stripe_account_id) + `); + await client.query(` + CREATE INDEX IF NOT EXISTS idx_${table}_org_id + ON "${schema}"."${table}" (org_id) + `); + } + } finally { + await client.end(); + } +}; diff --git a/packages/stripe-sync/src/eventTypeToTable.ts b/packages/stripe-sync/src/eventTypeToTable.ts new file mode 100644 index 000000000..6772935b2 --- /dev/null +++ b/packages/stripe-sync/src/eventTypeToTable.ts @@ -0,0 +1,61 @@ +/** + * Maps a Stripe event type to the sync DB table that stores the object. + * Returns undefined for event types that don't map to a known table. + */ +export const eventTypeToTable = ({ + eventType, +}: { + eventType: string; +}): string | undefined => { + if (eventType.startsWith("charge.dispute.")) return "disputes"; + if (eventType.startsWith("charge.")) return "charges"; + if (eventType.startsWith("checkout.session.")) return "checkout_sessions"; + if (eventType.startsWith("customer.subscription.")) return "subscriptions"; + if (eventType.startsWith("customer.tax_id.")) return "tax_ids"; + if (eventType.startsWith("customer.")) return "customers"; + if (eventType.startsWith("invoice.")) return "invoices"; + if (eventType.startsWith("product.")) return "products"; + if (eventType.startsWith("price.")) return "prices"; + if (eventType.startsWith("plan.")) return "plans"; + if (eventType.startsWith("setup_intent.")) return "setup_intents"; + if (eventType.startsWith("subscription_schedule.")) + return "subscription_schedules"; + if (eventType.startsWith("payment_method.")) return "payment_methods"; + if (eventType.startsWith("payment_intent.")) return "payment_intents"; + if (eventType.startsWith("credit_note.")) return "credit_notes"; + if (eventType.startsWith("radar.early_fraud_warning.")) + return "early_fraud_warnings"; + if (eventType.startsWith("refund.")) return "refunds"; + if (eventType.startsWith("review.")) return "reviews"; + if (eventType === "invoice_payment.paid") return "invoice_payments"; + + return undefined; +}; + +/** All tables in the stripe sync schema that store Stripe objects. */ +export const SYNCED_TABLES = [ + "charges", + "checkout_sessions", + "checkout_session_line_items", + "coupons", + "credit_notes", + "customers", + "disputes", + "early_fraud_warnings", + "events", + "invoices", + "invoice_payments", + "payment_intents", + "payment_methods", + "payouts", + "plans", + "prices", + "products", + "refunds", + "reviews", + "setup_intents", + "subscription_items", + "subscription_schedules", + "subscriptions", + "tax_ids", +] as const; diff --git a/packages/stripe-sync/src/index.ts b/packages/stripe-sync/src/index.ts new file mode 100644 index 000000000..dc23b3719 --- /dev/null +++ b/packages/stripe-sync/src/index.ts @@ -0,0 +1,12 @@ +export { addAutumnColumns } from "./addAccountIdColumn.js"; +export { eventTypeToTable, SYNCED_TABLES } from "./eventTypeToTable.js"; +export { + closeStripeSyncEngine, + getStripeSyncEngine, + processStripeSyncEvent, +} from "./initStripeSync.js"; +export { runStripeSyncMigrations } from "./runStripeSyncMigrations.js"; +export { + isSyncableEvent, + SYNCABLE_EVENT_PREFIXES, +} from "./syncableResources.js"; diff --git a/packages/stripe-sync/src/initStripeSync.ts b/packages/stripe-sync/src/initStripeSync.ts new file mode 100644 index 000000000..510d4e756 --- /dev/null +++ b/packages/stripe-sync/src/initStripeSync.ts @@ -0,0 +1,99 @@ +import { StripeSync } from "@supabase/stripe-sync-engine"; +import type Stripe from "stripe"; +import { eventTypeToTable } from "./eventTypeToTable.js"; + +const SCHEMA = "stripe"; + +let instance: StripeSync | null = null; +let initAttempted = false; + +/** + * Lazily creates a singleton StripeSync instance. + * Returns null if STRIPE_SYNC_DATABASE_URL is not configured or init fails. + */ +export const getStripeSyncEngine = (): StripeSync | null => { + if (instance) return instance; + if (initAttempted) return null; + + initAttempted = true; + + const databaseUrl = process.env.STRIPE_SYNC_DATABASE_URL; + const stripeSecretKey = + process.env.STRIPE_LIVE_SECRET_KEY || process.env.STRIPE_SANDBOX_SECRET_KEY; + + if (!databaseUrl || !stripeSecretKey) return null; + + try { + instance = new StripeSync({ + poolConfig: { + connectionString: databaseUrl, + max: 5, + keepAlive: true, + connectionTimeoutMillis: 5_000, + idleTimeoutMillis: 30_000, + }, + stripeSecretKey, + stripeWebhookSecret: "unused-processEvent-only", + schema: SCHEMA, + }); + } catch { + instance = null; + } + + return instance; +}; + +/** + * Upserts the Stripe event into the sync DB, then stamps the row + * with the originating Stripe account ID and org ID for multi-tenancy. + * Fully fail-open: any error is swallowed and returns silently. + */ +export const processStripeSyncEvent = async ({ + event, + stripeAccountId, + orgId, +}: { + event: Stripe.Event; + stripeAccountId?: string; + orgId?: string; +}): Promise => { + const engine = getStripeSyncEngine(); + if (!engine) return; + + try { + await engine.processEvent(event); + } catch { + return; + } + + const table = eventTypeToTable({ eventType: event.type }); + if (!table) return; + + const objectId = (event.data.object as { id?: string }).id; + if (!objectId) return; + + const accountId = stripeAccountId ?? event.account ?? null; + + if (!accountId && !orgId) return; + + try { + await engine.postgresClient.pool.query( + `UPDATE "${SCHEMA}"."${table}" SET stripe_account_id = COALESCE($1, stripe_account_id), org_id = COALESCE($2, org_id) WHERE id = $3`, + [accountId, orgId, objectId], + ); + } catch { + // Fail-open: metadata stamp is best-effort + } +}; + +/** Gracefully close the sync engine's PG pool (call on server shutdown). */ +export const closeStripeSyncEngine = async (): Promise => { + if (!instance) return; + try { + await instance.close(); + } catch { + // Best-effort cleanup + } + instance = null; + initAttempted = false; +}; diff --git a/packages/stripe-sync/src/runStripeSyncMigrations.ts b/packages/stripe-sync/src/runStripeSyncMigrations.ts new file mode 100644 index 000000000..565c7b95a --- /dev/null +++ b/packages/stripe-sync/src/runStripeSyncMigrations.ts @@ -0,0 +1,15 @@ +import { runMigrations } from "@supabase/stripe-sync-engine"; + +/** + * Run stripe-sync-engine migrations against the sync DB. + * Idempotent -- safe to run multiple times. + */ +export const runStripeSyncMigrations = async ({ + databaseUrl, + schema = "stripe", +}: { + databaseUrl: string; + schema?: string; +}): Promise => { + await runMigrations({ databaseUrl, schema }); +}; diff --git a/packages/stripe-sync/src/syncableResources.ts b/packages/stripe-sync/src/syncableResources.ts new file mode 100644 index 000000000..1317051e5 --- /dev/null +++ b/packages/stripe-sync/src/syncableResources.ts @@ -0,0 +1,16 @@ +export const SYNCABLE_EVENT_PREFIXES = [ + "customer.subscription.", + "payment_intent.", + "invoice.", + "customer.", + "product.", + "price.", +] as const; + +export const isSyncableEvent = ({ + eventType, +}: { + eventType: string; +}): boolean => { + return SYNCABLE_EVENT_PREFIXES.some((prefix) => eventType.startsWith(prefix)); +}; diff --git a/packages/stripe-sync/tsconfig.json b/packages/stripe-sync/tsconfig.json new file mode 100644 index 000000000..6efa41163 --- /dev/null +++ b/packages/stripe-sync/tsconfig.json @@ -0,0 +1,16 @@ +{ + "compilerOptions": { + "strict": true, + "composite": true, + "declaration": true, + "declarationMap": true, + "esModuleInterop": true, + "skipLibCheck": true, + "module": "Preserve", + "moduleResolution": "bundler", + "target": "ES2020", + "outDir": "./dist", + "rootDir": "." + }, + "include": ["src/**/*"] +} diff --git a/server/package.json b/server/package.json index 5cfedf8a2..a5b82626b 100644 --- a/server/package.json +++ b/server/package.json @@ -46,6 +46,7 @@ "@anthropic-ai/sdk": "^0.32.1", "@autumn/ksuid": "workspace:*", "@autumn/shared": "workspace:*", + "@autumn/stripe-sync": "workspace:*", "@aws-sdk/client-s3": "^3.1017.0", "@aws-sdk/client-scheduler": "^3.1004.0", "@aws-sdk/client-sqs": "^3.958.0", diff --git a/server/register.ts b/server/register.ts index 3a3699a77..df0fc8eec 100644 --- a/server/register.ts +++ b/server/register.ts @@ -1,29 +1,23 @@ import "dotenv/config"; import { loadLocalEnv } from "./src/utils/envUtils"; import Stripe from "stripe"; +import { + MAIN_STRIPE_EVENT_TYPES, + SYNC_STRIPE_EVENT_TYPES, +} from "./src/external/stripe/common/stripeConstants"; loadLocalEnv(); +const allEventTypes = [ + ...new Set([...MAIN_STRIPE_EVENT_TYPES, ...SYNC_STRIPE_EVENT_TYPES]), +] as Stripe.WebhookEndpointCreateParams.EnabledEvent[]; + const main = async () => { const stripe = new Stripe(process.env.STRIPE_SANDBOX_SECRET_KEY || ""); - const result = await stripe.webhookEndpoints.create({ url: `${process.env.STRIPE_WEBHOOK_URL}/webhooks/connect/sandbox`, - enabled_events: [ - "checkout.session.completed", - "customer.subscription.created", - "customer.subscription.updated", - "customer.subscription.deleted", - "customer.discount.deleted", - "invoice.paid", - "invoice.upcoming", - "invoice.created", - "invoice.finalized", - "invoice.updated", - "subscription_schedule.canceled", - "subscription_schedule.updated", - ], + enabled_events: allEventTypes, connect: true, }); diff --git a/server/src/external/stripe/common/stripeConstants.ts b/server/src/external/stripe/common/stripeConstants.ts new file mode 100644 index 000000000..9a9bf261c --- /dev/null +++ b/server/src/external/stripe/common/stripeConstants.ts @@ -0,0 +1,64 @@ +import type Stripe from "stripe"; + +type StripeEventType = Stripe.WebhookEndpointCreateParams.EnabledEvent; + +/** Events Autumn actively handles in its webhook handler. */ +export const MAIN_STRIPE_EVENT_TYPES: StripeEventType[] = [ + "checkout.session.completed", + "customer.subscription.created", + "customer.subscription.updated", + "customer.subscription.deleted", + "customer.discount.deleted", + "invoice.paid", + "invoice.upcoming", + "invoice.created", + "invoice.finalized", + "invoice.updated", + "subscription_schedule.canceled", + "subscription_schedule.updated", +]; + +/** Additional events needed to keep the stripe-sync DB up to date. */ +export const SYNC_STRIPE_EVENT_TYPES: StripeEventType[] = [ + // customers + "customer.created", + "customer.updated", + "customer.deleted", + + // subscriptions (extras beyond main) + "customer.subscription.paused", + "customer.subscription.resumed", + + // subscription schedules (extras beyond main) + "subscription_schedule.created", + "subscription_schedule.completed", + "subscription_schedule.released", + + // payment methods + "payment_method.attached", + "payment_method.detached", + "payment_method.updated", + + // products + "product.created", + "product.updated", + "product.deleted", + + // prices + "price.created", + "price.updated", + "price.deleted", + + // invoices (extras beyond main) + "invoice.deleted", + "invoice.payment_failed", + "invoice.payment_succeeded", + "invoice.voided", + "invoice.marked_uncollectible", + + // payment intents + "payment_intent.created", + "payment_intent.succeeded", + "payment_intent.payment_failed", + "payment_intent.canceled", +]; diff --git a/server/src/external/stripe/stripeOnboardingUtils.ts b/server/src/external/stripe/stripeOnboardingUtils.ts index 70e8b6f4e..5f6f1d43d 100644 --- a/server/src/external/stripe/stripeOnboardingUtils.ts +++ b/server/src/external/stripe/stripeOnboardingUtils.ts @@ -1,6 +1,10 @@ import { type AppEnv, ErrCode } from "@autumn/shared"; import Stripe from "stripe"; import RecaseError from "@/utils/errorUtils.js"; +import { + MAIN_STRIPE_EVENT_TYPES, + SYNC_STRIPE_EVENT_TYPES, +} from "./common/stripeConstants"; export const checkKeyValid = async (apiKey: string) => { const stripe = new Stripe(apiKey); @@ -29,20 +33,22 @@ export const createWebhookEndpoint = async ( const endpoint = await stripe.webhookEndpoints.create({ url: `${webhookBaseUrl}/webhooks/stripe/${orgId}/${env}`, - enabled_events: [ - "customer.subscription.created", - "customer.subscription.updated", - "customer.subscription.deleted", - "checkout.session.completed", - "invoice.paid", - "invoice.upcoming", - "invoice.created", - "invoice.finalized", - "invoice.updated", - "subscription_schedule.canceled", - "subscription_schedule.updated", - "customer.discount.deleted", - ], + enabled_events: [...MAIN_STRIPE_EVENT_TYPES, ...SYNC_STRIPE_EVENT_TYPES], + + // [ + // "customer.subscription.created", + // "customer.subscription.updated", + // "customer.subscription.deleted", + // "checkout.session.completed", + // "invoice.paid", + // "invoice.upcoming", + // "invoice.created", + // "invoice.finalized", + // "invoice.updated", + // "subscription_schedule.canceled", + // "subscription_schedule.updated", + // "customer.discount.deleted", + // ], }); return endpoint; diff --git a/server/src/external/stripe/stripeWebhookRouter.ts b/server/src/external/stripe/stripeWebhookRouter.ts index f1cbe08ac..9a07d41ce 100644 --- a/server/src/external/stripe/stripeWebhookRouter.ts +++ b/server/src/external/stripe/stripeWebhookRouter.ts @@ -4,6 +4,7 @@ import { handleStripeWebhookEvent } from "./handleStripeWebhookEvent.js"; import { stripeConnectSeederMiddleware } from "./webhookMiddlewares/stripeConnectSeederMiddleware.js"; import { stripeIdempotencyMiddleware } from "./webhookMiddlewares/stripeIdempotencyMiddleware.js"; import { stripeLegacySeederMiddleware } from "./webhookMiddlewares/stripeLegacySeederMiddleware.js"; +import { stripeSyncMiddleware } from "./webhookMiddlewares/stripeSyncMiddleware.js"; import { stripeToAutumnCustomerMiddleware } from "./webhookMiddlewares/stripeToAutumnCustomerMiddleware.js"; import type { StripeWebhookHonoEnv } from "./webhookMiddlewares/stripeWebhookContext.js"; import { stripeWebhookRefreshMiddleware } from "./webhookMiddlewares/stripeWebhookRefreshMiddleware.js"; @@ -15,6 +16,7 @@ stripeWebhookRouter.post( "/webhooks/stripe/:orgId/:env", stripeLegacySeederMiddleware, stripeWebhookRefreshMiddleware, + stripeSyncMiddleware, stripeToAutumnCustomerMiddleware, stripeLoggerMiddleware, stripeIdempotencyMiddleware, @@ -26,6 +28,7 @@ stripeWebhookRouter.post( "/webhooks/connect/:env", stripeConnectSeederMiddleware, stripeWebhookRefreshMiddleware, + stripeSyncMiddleware, stripeToAutumnCustomerMiddleware, stripeLoggerMiddleware, stripeIdempotencyMiddleware, diff --git a/server/src/external/stripe/webhookMiddlewares/stripeSyncMiddleware.ts b/server/src/external/stripe/webhookMiddlewares/stripeSyncMiddleware.ts new file mode 100644 index 000000000..d75fb66bd --- /dev/null +++ b/server/src/external/stripe/webhookMiddlewares/stripeSyncMiddleware.ts @@ -0,0 +1,57 @@ +import { isSyncableEvent, processStripeSyncEvent } from "@autumn/stripe-sync"; +import type { Context, Next } from "hono"; +import { isStripeSyncEnabled } from "@/internal/misc/stripeSync/stripeSyncStore.js"; +import type { + StripeWebhookContext, + StripeWebhookHonoEnv, +} from "./stripeWebhookContext.js"; + +/** + * Post-handler middleware that syncs Stripe events to the sync DB. + * Fire-and-forget -- errors are caught and logged, never propagated. + */ +export const stripeSyncMiddleware = async ( + c: Context, + next: Next, +) => { + await next(); + + const ctx = c.get("ctx") as StripeWebhookContext; + const { logger, org, stripeEvent } = ctx; + + if (!org || !stripeEvent) return; + if ( + process.env.NODE_ENV === "production" && + !isStripeSyncEnabled({ orgId: org.id }) + ) + return; + + if (!isSyncableEvent({ eventType: stripeEvent.type })) return; + + try { + const stripeAccountId = stripeEvent.account ?? undefined; + + void processStripeSyncEvent({ + event: stripeEvent, + stripeAccountId, + orgId: org.id, + }).catch((error) => { + logger.error(`Stripe sync failed for event ${stripeEvent.id}: ${error}`, { + error: { + message: error instanceof Error ? error.message : String(error), + }, + data: { + eventId: stripeEvent.id, + eventType: stripeEvent.type, + orgId: org.id, + }, + }); + }); + } catch (error) { + logger.error(`Stripe sync middleware error: ${error}`, { + error: { + message: error instanceof Error ? error.message : String(error), + }, + }); + } +}; diff --git a/server/src/init.ts b/server/src/init.ts index 57dd13765..37d441293 100644 --- a/server/src/init.ts +++ b/server/src/init.ts @@ -20,6 +20,8 @@ import { import "./internal/misc/requestBlocks/requestBlockStore.js"; import "./internal/misc/featureFlags/featureFlagStore.js"; import "./internal/misc/customerBlocks/customerBlockStore.js"; +import "./internal/misc/stripeSync/stripeSyncStore.js"; +import { closeStripeSyncEngine } from "@autumn/stripe-sync"; import { warmupRegionalRedis } from "./external/redis/initRedis.js"; import { createHonoApp } from "./initHono.js"; import { otelSdk } from "./instrumentation.js"; @@ -117,6 +119,7 @@ async function gracefulShutdown() { client.end(), clientCritical.end(), clientReplica?.end(), + closeStripeSyncEngine(), ]); console.log("Shutdown complete. Exiting process."); process.exit(0); diff --git a/server/src/internal/customers/cancel/handleCancelV2.ts b/server/src/internal/customers/cancel/handleCancelV2.ts index 930368e3f..61427da8b 100644 --- a/server/src/internal/customers/cancel/handleCancelV2.ts +++ b/server/src/internal/customers/cancel/handleCancelV2.ts @@ -8,10 +8,27 @@ import { executeBillingPlan } from "@/internal/billing/v2/execute/executeBilling import { evaluateStripeBillingPlan } from "@/internal/billing/v2/providers/stripe/actionBuilders/evaluateStripeBillingPlan"; import { logStripeBillingPlan } from "@/internal/billing/v2/providers/stripe/logs/logStripeBillingPlan"; import { logStripeBillingResult } from "@/internal/billing/v2/providers/stripe/logs/logStripeBillingResult"; +import { buildBillingLockKey } from "@/internal/billing/v2/utils/billingLock/buildBillingLockKey"; import { logAutumnBillingPlan } from "@/internal/billing/v2/utils/logs/logAutumnBillingPlan"; export const handleCancelV2 = createRoute({ - // body: CancelBodySchema, + lock: + process.env.NODE_ENV !== "development" + ? { + ttlMs: 120000, + errorMessage: + "Cancel already in progress for this customer, try again in a few seconds", + getKey: (c) => { + const ctx = c.get("ctx"); + if (!ctx.customerId) return null; + return buildBillingLockKey({ + orgId: ctx.org.id, + env: ctx.env, + customerId: ctx.customerId, + }); + }, + } + : undefined, handler: async (c) => { const ctx = c.get("ctx"); diff --git a/server/src/internal/misc/stripeSync/stripeSyncSchemas.ts b/server/src/internal/misc/stripeSync/stripeSyncSchemas.ts new file mode 100644 index 000000000..1a91f653f --- /dev/null +++ b/server/src/internal/misc/stripeSync/stripeSyncSchemas.ts @@ -0,0 +1,7 @@ +import { z } from "zod/v4"; + +export const StripeSyncConfigSchema = z.object({ + enabledOrgIds: z.array(z.string()).default([]), +}); + +export type StripeSyncConfig = z.infer; diff --git a/server/src/internal/misc/stripeSync/stripeSyncStore.ts b/server/src/internal/misc/stripeSync/stripeSyncStore.ts new file mode 100644 index 000000000..e711de5d0 --- /dev/null +++ b/server/src/internal/misc/stripeSync/stripeSyncStore.ts @@ -0,0 +1,20 @@ +import { registerEdgeConfig } from "@/internal/misc/edgeConfig/edgeConfigRegistry.js"; +import { createEdgeConfigStore } from "@/internal/misc/edgeConfig/edgeConfigStore.js"; +import { + type StripeSyncConfig, + StripeSyncConfigSchema, +} from "./stripeSyncSchemas.js"; + +const store = createEdgeConfigStore({ + s3Key: "admin/stripe-sync-config.json", + schema: StripeSyncConfigSchema, + defaultValue: () => ({ enabledOrgIds: [] }), + pollIntervalMs: 60_000, +}); + +registerEdgeConfig({ store }); + +/** Pure in-memory lookup -- zero I/O, sync. */ +export const isStripeSyncEnabled = ({ orgId }: { orgId: string }): boolean => { + return store.get().enabledOrgIds.includes(orgId); +}; diff --git a/server/tests/_temp/temp2.test.ts b/server/tests/_temp/temp2.test.ts index e846fc0c9..edced424c 100644 --- a/server/tests/_temp/temp2.test.ts +++ b/server/tests/_temp/temp2.test.ts @@ -1,67 +1,106 @@ -import { expect, test } from "bun:test"; +import { test } from "bun:test"; +import type { AttachParamsV1Input } from "@autumn/shared"; +import { + applySubscriptionDiscount, + createPercentCoupon, + getStripeSubscription, +} from "@tests/integration/billing/utils/discounts/discountTestUtils"; +import { TestFeature } from "@tests/setup/v2Features.js"; import { items } from "@tests/utils/fixtures/items"; import { products } from "@tests/utils/fixtures/products"; import { initScenario, s } from "@tests/utils/testInitUtils/initScenario"; import chalk from "chalk"; -import { CusService } from "@/internal/customers/CusService"; -test(`${chalk.yellowBright("check where default payment method lives after attach")}`, async () => { - const customerId = "pm-location-check"; +/** + * Repro for Mintlify multi-entity discount bug. + * + * Production sequence (from Axiom logs): + * 1. Entity 1 attaches pro (upgrade from Hobby) — creates subscription + * 2. Discount applied to subscription + * 3. Entity 1 cancels pro (end_of_cycle) — creates subscription schedule + * 4. Entity 1 uncancels pro — schedule modified + * 5. Entity 2 attaches pro — subscription update succeeds but schedule + * creation/update fails: "Discount di_xxx has exceeded its maximum + * number of applications and cannot be reused" + * + * The error is in stripeDiscountsToPhaseDiscounts which passes + * { discount: "di_xxx" } to schedule phases — but the discount is bound + * to the subscription, not the schedule. + */ +test(`${chalk.yellowBright("bug repro: entity 2 attach after cancel+uncancel with discount on sub")}`, async () => { + const customerId = "multi-ent-cancel-uncancel"; const messagesItem = items.monthlyMessages({ includedUsage: 100 }); + const pro = products.pro({ - id: "pro-pm-check", + id: "pro", items: [messagesItem], }); - const { ctx } = await initScenario({ + // Step 1: Entity 1 attaches pro + const { autumnV2_2, autumnV1, entities } = await initScenario({ customerId, setup: [ s.customer({ paymentMethod: "success" }), s.products({ list: [pro] }), + s.entities({ count: 2, featureId: TestFeature.Users }), + ], + actions: [ + s.billing.attach({ productId: pro.id, entityIndex: 0 }), ], - actions: [s.attach({ productId: pro.id })], }); - const fullCustomer = await CusService.getFull({ - ctx, - idOrInternalId: customerId, + // Step 2: Apply a discount to entity 1's subscription + const { stripeCli, subscription } = await getStripeSubscription({ + customerId, }); - const stripeCustomerId = fullCustomer.processor?.id; - expect(stripeCustomerId).toBeDefined(); - - const subs = await ctx.stripeCli.subscriptions.list({ - customer: stripeCustomerId!, - status: "all", + const coupon = await createPercentCoupon({ + stripeCli, + percentOff: 10, }); - const subscription = subs.data.find( - (sub) => sub.status === "active" || sub.status === "trialing", - ); - expect(subscription).toBeDefined(); + await applySubscriptionDiscount({ + stripeCli, + subscriptionId: subscription.id, + couponIds: [coupon.id], + }); - const subDefaultPm = subscription!.default_payment_method; - console.log(`Subscription default_payment_method: ${JSON.stringify(subDefaultPm)}`); + // Step 3: Entity 1 cancels pro (end_of_cycle) — creates subscription schedule + console.log("Canceling entity 1 pro (end of cycle)..."); + await autumnV2_2.subscriptions.update({ + customer_id: customerId, + entity_id: entities[0].id, + plan_id: pro.id, + cancel_action: "cancel_end_of_cycle", + }); - const stripeCustomer = await ctx.stripeCli.customers.retrieve(stripeCustomerId!); - const customerDefaultPm = "deleted" in stripeCustomer - ? null - : stripeCustomer.invoice_settings?.default_payment_method; - console.log(`Customer invoice_settings.default_payment_method: ${JSON.stringify(customerDefaultPm)}`); + await new Promise((resolve) => setTimeout(resolve, 3000)); - const customerDefaultSource = "deleted" in stripeCustomer - ? null - : stripeCustomer.default_source; - console.log(`Customer default_source: ${JSON.stringify(customerDefaultSource)}`); + // Step 4: Entity 1 uncancels pro — schedule released/modified + console.log("Uncanceling entity 1 pro..."); + await autumnV2_2.subscriptions.update({ + customer_id: customerId, + entity_id: entities[0].id, + plan_id: pro.id, + cancel_action: "uncancel", + }); - if (subDefaultPm) { - console.log(chalk.green("Payment method is set at the SUBSCRIPTION level")); - } else if (customerDefaultPm) { - console.log(chalk.green("Payment method is set at the CUSTOMER level (invoice_settings)")); - } else if (customerDefaultSource) { - console.log(chalk.green("Payment method is set at the CUSTOMER level (default_source)")); - } else { - console.log(chalk.red("No default payment method found on subscription or customer")); + await new Promise((resolve) => setTimeout(resolve, 3000)); + + // Step 5: Entity 2 attaches pro + console.log("Attaching pro to entity 2..."); + try { + const result = await autumnV2_2.billing.attach({ + customer_id: customerId, + entity_id: entities[1].id, + plan_id: pro.id, + redirect_mode: "if_required", + }); + console.log("Entity 2 result:", JSON.stringify(result, null, 2)); + console.log(chalk.red("BUG NOT REPRODUCED — entity 2 attach succeeded")); + } catch (error: any) { + console.log(chalk.green("BUG REPRODUCED — entity 2 attach failed:")); + console.log("Error:", JSON.stringify(error, null, 2)); } }); diff --git a/server/tests/integration/billing/attach/new-plan/schedule-addon-attach.test.ts b/server/tests/integration/billing/attach/new-plan/schedule-addon-attach.test.ts new file mode 100644 index 000000000..9df42def7 --- /dev/null +++ b/server/tests/integration/billing/attach/new-plan/schedule-addon-attach.test.ts @@ -0,0 +1,113 @@ +import { test } from "bun:test"; +import type { ApiCustomerV3 } from "@autumn/shared"; +import { + applySubscriptionDiscount, + createPercentCoupon, + getStripeSubscription, +} from "@tests/integration/billing/utils/discounts/discountTestUtils"; +import { expectCustomerFeatureCorrect } from "@tests/integration/billing/utils/expectCustomerFeatureCorrect"; +import { expectCustomerInvoiceCorrect } from "@tests/integration/billing/utils/expectCustomerInvoiceCorrect"; +import { + expectProductCanceling, + expectProductScheduled, +} from "@tests/integration/billing/utils/expectCustomerProductCorrect"; +import { TestFeature } from "@tests/setup/v2Features"; +import { items } from "@tests/utils/fixtures/items.js"; +import { products } from "@tests/utils/fixtures/products.js"; +import { initScenario, s } from "@tests/utils/testInitUtils/initScenario.js"; +import chalk from "chalk"; + +/** + * Reproduces: attaching a one-off add-on to a subscription that already + * has a schedule (from a prior downgrade) fails with Stripe error + * "You cannot migrate a subscription that is already attached to a schedule". + * + * The root cause is that setupStripeBillingContext skips fetching the schedule + * when targetCustomerProduct is undefined (new attachment, no transition). + * buildStripeSubscriptionScheduleAction then sees hasSchedule=false and emits + * type:"create" instead of type:"update", which Stripe rejects. + */ +test.concurrent(`${chalk.yellowBright("bug-repro: attach one-off add-on after scheduled downgrade hits 'subscription already attached to schedule'")}`, async () => { + const customerId = "repro-schedule-addon-attach"; + + const messagesItem = items.monthlyMessages({ includedUsage: 100 }); + + const premium = products.premium({ + id: "premium", + items: [messagesItem], + }); + + const pro = products.pro({ + id: "pro", + items: [messagesItem], + }); + + const topupAddon = products.base({ + id: "topup-addon", + items: [ + items.oneOffWords({ + includedUsage: 0, + billingUnits: 100, + price: 25, + }), + ], + }); + + const { autumnV1 } = await initScenario({ + customerId, + setup: [ + s.customer({ paymentMethod: "success" }), + s.products({ list: [premium, pro, topupAddon] }), + ], + actions: [s.billing.attach({ productId: premium.id })], + }); + + // Apply a 20% discount to the subscription (matches production scenario) + const { stripeCli, subscription: subBefore } = await getStripeSubscription({ + customerId, + }); + + const coupon = await createPercentCoupon({ + stripeCli, + percentOff: 20, + }); + + await applySubscriptionDiscount({ + stripeCli, + subscriptionId: subBefore.id, + couponIds: [coupon.id], + }); + + // Schedule downgrade: Premium -> Pro (creates a subscription schedule) + await autumnV1.billing.attach({ + customer_id: customerId, + product_id: pro.id, + redirect_mode: "if_required", + }); + + // Verify the downgrade was scheduled correctly + const customer = await autumnV1.customers.get(customerId); + await expectProductCanceling({ customer, productId: premium.id }); + await expectProductScheduled({ customer, productId: pro.id }); + + await autumnV1.billing.attach({ + customer_id: customerId, + product_id: topupAddon.id, + redirect_mode: "if_required", + options: [{ feature_id: TestFeature.Words, quantity: 100 }], + }); + + const customerAfter = await autumnV1.customers.get(customerId); + expectCustomerFeatureCorrect({ + customer: customerAfter, + featureId: TestFeature.Words, + balance: 100, + usage: 0, + }); + + await expectCustomerInvoiceCorrect({ + customer: customerAfter, + count: 2, + latestTotal: 25, + }); +});