diff --git a/bun.lock b/bun.lock index 7c423c2d..3284febd 100644 --- a/bun.lock +++ b/bun.lock @@ -96,12 +96,14 @@ "@polar-sh/sdk": "0.49.0", "ai": "7.0.49", "better-auth": "^1.6.23", + "buffer": "6.0.3", "convex": "^1.42.1", "files-sdk": "2.2.3", }, "devDependencies": { "@baseblocks/tsconfig": "workspace:*", "@types/bun": "1.3.14", + "convex-test": "0.0.55", }, "peerDependencies": { "typescript": ">=5 <8", @@ -1388,6 +1390,8 @@ "bare-url": ["bare-url@2.4.7", "", { "dependencies": { "bare-path": "^3.0.0" } }, "sha512-o8CRCiJtib+ycO3mE4A5UChtGX4dDP2XxsWVu9P+Zc3H8tcmKwNVEDoDTXmwN+uuMhfKeT7/i7Y26xS8W7ohoA=="], + "base64-js": ["base64-js@1.5.1", "", {}, "sha512-AKpaYlHn8t4SVbOHCy+b5+KKgvR4vrsD8vbvrbiQJps7fKDTkjkDry6ji0rUJjC0kzbNePLwzxq8iypo41qeWA=="], + "baseline-browser-mapping": ["baseline-browser-mapping@2.11.12", "", { "bin": { "baseline-browser-mapping": "dist/cli.cjs" } }, "sha512-r7WnVImvVCeFpf2DOXfy41aPWzeNg3H/A2X4dKmy1QL0MSyyk/e7z8ihJ3N6Nn2PsdhkVlqnEfnUE4a05P2aTA=="], "basic-ftp": ["basic-ftp@5.3.1", "", {}, "sha512-bopVNp6ugyA150DDuZfPFdt1KZ5a94ZDiwX4hMgZDzF+GttD80lEy8kj98kbyhLXnPvhtIo93mdnLIjpCAeeOw=="], @@ -1408,6 +1412,8 @@ "braces": ["braces@3.0.3", "", { "dependencies": { "fill-range": "^7.1.1" } }, "sha512-yQbXgO/OSZVD2IsiLlro+7Hf6Q18EJrKSEsdoMzKePKXct3gvD8oLcOQdIzGupr5Fj+EDe8gO/lxc1BzfMpxvA=="], + "buffer": ["buffer@6.0.3", "", { "dependencies": { "base64-js": "^1.3.1", "ieee754": "^1.2.1" } }, "sha512-FTiCpNxtwiZZHEZbcbTIcZjERVICn9yq/pDFkTl95/AxzD1naBctN7YO68riM/gLSDY7sdrMby8hofADYuuqOA=="], + "buffer-crc32": ["buffer-crc32@0.2.13", "", {}, "sha512-VO9Ht/+p3SN7SKWqcrgEzjGbRSJYTx+Q1pTQC0wrWqHx0vpJraQ6GtHx8tvcg1rlK1byhU5gccxgOgj7B0TDkQ=="], "bun-types": ["bun-types@1.3.14", "", { "dependencies": { "@types/node": "*" } }, "sha512-4N0ig0fEomHt5R0KCFWjovxow98rIoRwKolrYdCcknNwMekCXRnWEUvgu5soYV8QXtVsrUD8B95MBOZGPvr6KQ=="], @@ -1484,6 +1490,8 @@ "convex-helpers": ["convex-helpers@0.1.121", "", { "peerDependencies": { "@standard-schema/spec": "^1.0.0", "convex": "^1.43.0", "hono": "^4.0.5", "react": "^17.0.2 || ^18.0.0 || ^19.0.0", "typescript": "^5.5 || ^6.0.0", "zod": "^3.25.0 || ^4.0.0" }, "optionalPeers": ["@standard-schema/spec", "hono", "react", "typescript", "zod"], "bin": { "convex-helpers": "bin.cjs" } }, "sha512-qLrAr8p3dGI59xVa4485ARaGPuBRMjKOXC56FELKVnJ5AoIHFOkjPSKR8phE3KZUxyi6eZtb54Z/MYcS+pxoAQ=="], + "convex-test": ["convex-test@0.0.55", "", { "peerDependencies": { "convex": "^1.43.0" } }, "sha512-ablJ6vR3WUpyTaYrrRQf8yxInGrhHyqBEQsCFzloda3G0MDEu2RMKNGUkvlBXYWPiiEP0UIKPS7gZuLbqGnvQw=="], + "cookie": ["cookie@0.7.2", "", {}, "sha512-yki5XnKuf750l50uGTllt6kKILY4nQ1eNIQatoXEByZ5dWgnKqbnqmTrBE5B4N7lrMJKQ2ytWMiTO2o0v6Ew/w=="], "cookie-signature": ["cookie-signature@1.2.2", "", {}, "sha512-D76uU73ulSXrD1UXF4KE2TMxVVwhsnCgfAyTg9k8P6KGZjlXKrOLe4dJQKI3Bxi5wjesZoFXJWElNWBjPZMbhg=="], @@ -1750,6 +1758,8 @@ "icu-minify": ["icu-minify@4.13.6", "", { "dependencies": { "@formatjs/icu-messageformat-parser": "^3.4.0" } }, "sha512-iYZGCJZ+kX6o7GrxpVe2sOSdW86AvEqh8RQBvWeBd9jqmuABsMc2B6xongACfItLOogyIWH6GuBslNHr79OU8Q=="], + "ieee754": ["ieee754@1.2.1", "", {}, "sha512-dcyqhDvX1C46lXZcVqCpK+FtMRQVdIMN6/Df5js2zouUsqG7I6sFxitIC+7KYK29KdXOLHdu9zL4sFnoVQnqaA=="], + "image-ssim": ["image-ssim@0.2.0", "", {}, "sha512-W7+sO6/yhxy83L0G7xR8YAc5Z5QFtYEXXRV6EaE8tuYBZJnA3gVgp3q7X7muhLZVodeb9UfvjSbwt9VJwjIYAg=="], "immediate": ["immediate@3.0.6", "", {}, "sha512-XXOFtyqDjNDAQxVfYxuF7g9Il/IbWmmlQg2MYKOH8ExIT1qg6xc4zyS3HaEEATgs1btfzxq15ciUiY7gjSXRGQ=="], diff --git a/packages/backend/convex/_generated/api.d.ts b/packages/backend/convex/_generated/api.d.ts index a7a16971..444b4600 100644 --- a/packages/backend/convex/_generated/api.d.ts +++ b/packages/backend/convex/_generated/api.d.ts @@ -16,7 +16,8 @@ import type * as billing_checkoutIntent from "../billing/checkoutIntent.js"; import type * as billing_polar from "../billing/polar.js"; import type * as billingModel from "../billingModel.js"; import type * as billingRetention from "../billingRetention.js"; -import type * as billingWebhooks from "../billingWebhooks.js"; +import type * as billing_webhook_model from "../billing_webhook_model.js"; +import type * as billing_webhooks from "../billing_webhooks.js"; import type * as crons from "../crons.js"; import type * as deploymentPreflight from "../deploymentPreflight.js"; import type * as draftRestore from "../draftRestore.js"; @@ -36,7 +37,6 @@ import type * as libraries from "../libraries.js"; import type * as model_aiCredits from "../model/aiCredits.js"; import type * as model_aiWorkspaceBounds from "../model/aiWorkspaceBounds.js"; import type * as model_aiWorkspaceFingerprint from "../model/aiWorkspaceFingerprint.js"; -import type * as model_billingEventOrdering from "../model/billingEventOrdering.js"; import type * as model_billingRetention from "../model/billingRetention.js"; import type * as model_contentObjects from "../model/contentObjects.js"; import type * as model_draft from "../model/draft.js"; @@ -47,7 +47,6 @@ import type * as model_libraryAccess from "../model/libraryAccess.js"; import type * as model_pageDeletion from "../model/pageDeletion.js"; import type * as model_pageDocuments from "../model/pageDocuments.js"; import type * as model_pageHierarchy from "../model/pageHierarchy.js"; -import type * as model_polarOrderAmounts from "../model/polarOrderAmounts.js"; import type * as model_publishedRelease from "../model/publishedRelease.js"; import type * as model_releaseChangeDetails from "../model/releaseChangeDetails.js"; import type * as model_releaseChanges from "../model/releaseChanges.js"; @@ -104,7 +103,8 @@ declare const fullApi: ApiFromModules<{ "billing/polar": typeof billing_polar; billingModel: typeof billingModel; billingRetention: typeof billingRetention; - billingWebhooks: typeof billingWebhooks; + billing_webhook_model: typeof billing_webhook_model; + billing_webhooks: typeof billing_webhooks; crons: typeof crons; deploymentPreflight: typeof deploymentPreflight; draftRestore: typeof draftRestore; @@ -124,7 +124,6 @@ declare const fullApi: ApiFromModules<{ "model/aiCredits": typeof model_aiCredits; "model/aiWorkspaceBounds": typeof model_aiWorkspaceBounds; "model/aiWorkspaceFingerprint": typeof model_aiWorkspaceFingerprint; - "model/billingEventOrdering": typeof model_billingEventOrdering; "model/billingRetention": typeof model_billingRetention; "model/contentObjects": typeof model_contentObjects; "model/draft": typeof model_draft; @@ -135,7 +134,6 @@ declare const fullApi: ApiFromModules<{ "model/pageDeletion": typeof model_pageDeletion; "model/pageDocuments": typeof model_pageDocuments; "model/pageHierarchy": typeof model_pageHierarchy; - "model/polarOrderAmounts": typeof model_polarOrderAmounts; "model/publishedRelease": typeof model_publishedRelease; "model/releaseChangeDetails": typeof model_releaseChangeDetails; "model/releaseChanges": typeof model_releaseChanges; diff --git a/packages/backend/convex/_generated/server.d.ts b/packages/backend/convex/_generated/server.d.ts index 9731e330..3fe4bed1 100644 --- a/packages/backend/convex/_generated/server.d.ts +++ b/packages/backend/convex/_generated/server.d.ts @@ -25,7 +25,12 @@ import type { DataModel } from "./dataModel.js"; * Typesafe environment variables declared in `convex.config.ts`. */ type Env = { + readonly BASEBLOCKS_BILLING_ENVIRONMENT: "sandbox" | "production"; + readonly BASEBLOCKS_PAST_DUE_GRACE_DAYS: string | undefined; readonly INTEGRATIONS_ENABLED: "true" | "false"; + readonly POLAR_ACCESS_TOKEN: string; + readonly POLAR_ALLOW_PRODUCTION: "true" | "false"; + readonly POLAR_WEBHOOK_SECRET: string; }; /** diff --git a/packages/backend/convex/billing-webhook-model.test.ts b/packages/backend/convex/billing-webhook-model.test.ts new file mode 100644 index 00000000..d1f3ca78 --- /dev/null +++ b/packages/backend/convex/billing-webhook-model.test.ts @@ -0,0 +1,227 @@ +import { describe, expect, test } from "bun:test"; +import { convexTest } from "convex-test"; +import type { QueryCtx } from "./_generated/server"; +import { + applyPolarBillingEvent, + type PolarBillingEventCommand, +} from "./billing_webhook_model"; +import schema from "./schema"; + +const organizationId = "organization-alpha"; +const modules = { "./_generated/api.ts": async () => ({}) }; +const paidAt = Date.parse("2026-08-12T21:22:19.594Z"); + +function paidCreditOrder( + overrides: Partial = {}, +): PolarBillingEventCommand { + return { + providerEnvironment: "production", + deliveryId: "delivery-order-paid-1", + eventType: "order.paid", + eventOccurredAt: paidAt, + event: { + kind: "order", + organizationId, + providerOrderId: "polar-order-1", + providerCheckoutId: "polar-checkout-1", + providerCustomerId: "polar-customer-1", + providerProductId: "polar-credit-product", + state: "paid", + subtotalAmountMinor: 500n, + discountAmountMinor: 0n, + taxAmountMinor: 83n, + grossAmountMinor: 500n, + netAmountMinor: 417n, + refundedGrossAmountMinor: 0n, + currency: "usd", + billingReason: "purchase", + providerModifiedAt: paidAt, + }, + ...overrides, + }; +} + +async function seedWorkspace( + t: ReturnType, + kind: "credits" | "plus" = "credits", +) { + await t.mutation(async (ctx) => { + const now = Date.now(); + await ctx.db.insert("workspaceProfiles", { + organizationId, + intent: kind === "plus" ? "work" : "personal", + source: "onboarding", + schemaVersion: 1, + createdAt: now, + updatedAt: now, + }); + await ctx.db.insert("billingCatalogItems", { + providerEnvironment: "production", + sku: kind === "plus" ? "plus_monthly" : "ai_credit_top_up", + kind: kind === "plus" ? "plus" : "aiCreditPack", + providerProductId: + kind === "plus" ? "polar-plus-product" : "polar-credit-product", + providerPriceId: kind === "plus" ? "polar-plus-price" : undefined, + planKey: kind === "plus" ? "plus" : undefined, + recurringInterval: kind === "plus" ? "month" : undefined, + priceAmountMinor: kind === "plus" ? 1_000n : 500n, + currency: "usd", + creditUnits: kind === "plus" ? 2_000_000n : undefined, + configurationVersion: + kind === "plus" ? "polar-plus-v1" : "polar-credit-v1", + active: true, + createdAt: now, + updatedAt: now, + }); + }); +} + +async function readState(ctx: QueryCtx) { + const [account, orders, entitlement] = await Promise.all([ + ctx.db + .query("aiCreditAccounts") + .withIndex("by_organization", (query) => + query.eq("organizationId", organizationId), + ) + .unique(), + ctx.db + .query("billingOrders") + .withIndex("by_organization_created", (query) => + query.eq("organizationId", organizationId), + ) + .collect(), + ctx.db + .query("workspaceEntitlements") + .withIndex("by_organization", (query) => + query.eq("organizationId", organizationId), + ) + .unique(), + ]); + return { account, orders, entitlement }; +} + +describe("Polar billing event interface", () => { + test("grants a paid order exactly once across both forms of redelivery", async () => { + const t = convexTest(schema, modules); + await seedWorkspace(t); + + const applied = await t.mutation(async (ctx) => + applyPolarBillingEvent(ctx, paidCreditOrder()), + ); + const repeatedDelivery = await t.mutation(async (ctx) => + applyPolarBillingEvent(ctx, paidCreditOrder()), + ); + await t.mutation(async (ctx) => + applyPolarBillingEvent( + ctx, + paidCreditOrder({ deliveryId: "delivery-provider-redelivery" }), + ), + ); + const state = await t.query(readState); + + expect(applied.outcome).toBe("applied"); + expect(repeatedDelivery.outcome).toBe("duplicate"); + expect(state.account?.availablePrepaidUnits).toBe(5_000_000n); + expect(state.orders).toHaveLength(1); + }); + + test("refunds monotonically and rejects an older paid state", async () => { + const t = convexTest(schema, modules); + await seedWorkspace(t); + await t.mutation(async (ctx) => + applyPolarBillingEvent(ctx, paidCreditOrder()), + ); + const refundedAt = Date.parse("2026-08-14T00:00:00Z"); + const order = paidCreditOrder().event; + if (order.kind !== "order") throw new Error("Expected an order fixture"); + await t.mutation(async (ctx) => + applyPolarBillingEvent( + ctx, + paidCreditOrder({ + deliveryId: "delivery-order-refunded", + eventType: "order.refunded", + eventOccurredAt: refundedAt, + event: { + ...order, + state: "refunded", + refundedGrossAmountMinor: 500n, + providerModifiedAt: refundedAt, + }, + }), + ), + ); + const stale = await t.mutation(async (ctx) => + applyPolarBillingEvent( + ctx, + paidCreditOrder({ deliveryId: "delivery-stale-paid" }), + ), + ); + const state = await t.query(readState); + + expect(stale.outcome).toBe("ignored"); + expect(state.account?.availablePrepaidUnits).toBe(0n); + expect(state.orders).toHaveLength(1); + }); + + test("reconciles a newer authoritative order amount upward", async () => { + const t = convexTest(schema, modules); + await seedWorkspace(t); + await t.mutation(async (ctx) => + applyPolarBillingEvent(ctx, paidCreditOrder()), + ); + const order = paidCreditOrder().event; + if (order.kind !== "order") throw new Error("Expected an order fixture"); + const updatedAt = Date.parse("2026-08-13T00:00:00Z"); + await t.mutation(async (ctx) => + applyPolarBillingEvent( + ctx, + paidCreditOrder({ + deliveryId: "delivery-order-updated-amount", + eventType: "order.updated", + eventOccurredAt: updatedAt, + event: { + ...order, + grossAmountMinor: 700n, + providerModifiedAt: updatedAt, + }, + }), + ), + ); + + const state = await t.query(readState); + expect(state.account?.availablePrepaidUnits).toBe(7_000_000n); + expect(state.orders).toHaveLength(1); + }); + + test("atomically enables an active subscription entitlement", async () => { + const t = convexTest(schema, modules); + await seedWorkspace(t, "plus"); + + const result = await t.mutation(async (ctx) => + applyPolarBillingEvent(ctx, { + providerEnvironment: "production", + deliveryId: "delivery-subscription-active-1", + eventType: "subscription.active", + eventOccurredAt: paidAt, + event: { + kind: "subscription", + organizationId, + providerSubscriptionId: "polar-subscription-1", + providerCustomerId: "polar-customer-1", + providerProductId: "polar-plus-product", + providerStatus: "active", + seatQuantity: 2, + cancelAtPeriodEnd: false, + currentPeriodStart: Date.parse("2026-08-01T00:00:00.000Z"), + currentPeriodEnd: Date.parse("2026-09-01T00:00:00.000Z"), + providerModifiedAt: paidAt, + }, + }), + ); + const state = await t.query(readState); + + expect(result.outcome).toBe("applied"); + expect(state.entitlement?.plusEnabled).toBe(true); + expect(state.entitlement?.paidSeatCapacity).toBe(2); + }); +}); diff --git a/packages/backend/convex/billing-webhooks.test.ts b/packages/backend/convex/billing-webhooks.test.ts new file mode 100644 index 00000000..f966f136 --- /dev/null +++ b/packages/backend/convex/billing-webhooks.test.ts @@ -0,0 +1,168 @@ +import { createHmac } from "node:crypto"; +import { afterEach, describe, expect, test } from "bun:test"; +import { handlePolarWebhook } from "./billing_webhooks"; + +type RegisteredHttpAction = { + _handler: (ctx: unknown, request: Request) => Promise; +}; + +const originalEnvironment = process.env.BASEBLOCKS_BILLING_ENVIRONMENT; +const originalSecret = process.env.POLAR_WEBHOOK_SECRET; + +afterEach(() => { + process.env.BASEBLOCKS_BILLING_ENVIRONMENT = originalEnvironment; + process.env.POLAR_WEBHOOK_SECRET = originalSecret; +}); + +function signedRequest(secret: string, body: string, valid = true) { + const deliveryId = "delivery-official-verifier-contract"; + const timestamp = Math.floor(Date.now() / 1_000); + const signature = createHmac("sha256", valid ? secret : "wrong-secret") + .update(`${deliveryId}.${timestamp}.${body}`) + .digest("base64"); + return new Request("https://example.test/billing/webhooks/polar", { + method: "POST", + body, + headers: { + "content-type": "application/json", + "webhook-id": deliveryId, + "webhook-timestamp": String(timestamp), + "webhook-signature": `v1,${signature}`, + }, + }); +} + +function validOrderPaidBody() { + const occurredAt = "2026-08-12T21:22:19.594Z"; + return JSON.stringify({ + type: "order.paid", + timestamp: occurredAt, + data: { + id: "polar-order-1", + created_at: occurredAt, + modified_at: occurredAt, + status: "paid", + paid: true, + subtotal_amount: 500, + discount_amount: 0, + net_amount: 417, + tax_amount: 83, + total_amount: 500, + applied_balance_amount: 0, + due_amount: 0, + refunded_amount: 0, + refunded_tax_amount: 0, + currency: "usd", + billing_reason: "purchase", + billing_name: null, + billing_address: null, + invoice_number: "TEST-1", + is_invoice_generated: true, + receipt_number: "TEST-1", + customer_id: "customer-1", + product_id: "product-1", + discount_id: null, + subscription_id: null, + checkout_id: "checkout-1", + metadata: { baseblocks_workspace_id: "workspace-1" }, + platform_fee_amount: 0, + platform_fee_currency: null, + customer: { + id: "customer-1", + created_at: occurredAt, + modified_at: occurredAt, + metadata: {}, + external_id: "workspace-1", + email_verified: true, + type: "organization", + name: "Test workspace", + billing_name: null, + billing_address: null, + tax_id: null, + organization_id: "polar-organization-1", + deleted_at: null, + avatar_url: null, + }, + product: { + metadata: {}, + id: "product-1", + created_at: occurredAt, + modified_at: occurredAt, + trial_interval: null, + trial_interval_count: null, + name: "AI credits", + description: null, + visibility: "public", + recurring_interval: null, + recurring_interval_count: null, + meter_interval: null, + meter_interval_count: null, + is_recurring: false, + is_archived: false, + organization_id: "polar-organization-1", + }, + discount: null, + subscription: null, + items: [ + { + created_at: occurredAt, + modified_at: occurredAt, + id: "item-1", + label: "AI credits", + amount: 417, + tax_amount: 83, + proration: false, + product_price_id: "price-1", + }, + ], + description: "AI credits", + refundable_amount: 417, + refundable_tax_amount: 83, + }, + }); +} + +describe("Polar webhook HTTP boundary", () => { + test("verifies and applies a real SDK-shaped order with the complete secret", async () => { + const secret = "whsec_literal-secret-contract"; + const mutations: unknown[] = []; + process.env.BASEBLOCKS_BILLING_ENVIRONMENT = "production"; + process.env.POLAR_WEBHOOK_SECRET = secret; + + const response = await ( + handlePolarWebhook as unknown as RegisteredHttpAction + )._handler( + { + runMutation: async (_reference: unknown, command: unknown) => { + mutations.push(command); + return { outcome: "applied" }; + }, + }, + signedRequest(secret, validOrderPaidBody()), + ); + + expect(response.status).toBe(202); + expect(mutations).toEqual([ + expect.objectContaining({ + eventType: "order.paid", + event: expect.objectContaining({ + organizationId: "workspace-1", + providerOrderId: "polar-order-1", + grossAmountMinor: 500n, + }), + }), + ]); + }); + + test("rejects an invalid signature before processing the payload", async () => { + const secret = "whsec_literal-secret-contract"; + process.env.BASEBLOCKS_BILLING_ENVIRONMENT = "production"; + process.env.POLAR_WEBHOOK_SECRET = secret; + + const response = await ( + handlePolarWebhook as unknown as RegisteredHttpAction + )._handler({}, signedRequest(secret, validOrderPaidBody(), false)); + + expect(response.status).toBe(403); + }); +}); diff --git a/packages/backend/convex/billing/polar.test.ts b/packages/backend/convex/billing/polar.test.ts index cff84bd7..dab6f325 100644 --- a/packages/backend/convex/billing/polar.test.ts +++ b/packages/backend/convex/billing/polar.test.ts @@ -1,4 +1,3 @@ -import { createHmac } from "node:crypto"; import { describe, expect, test } from "bun:test"; import { billingOperationMetadata, @@ -8,7 +7,6 @@ import { normalizeSubscriptionLifecycle, PolarApiError, resolvePolarOrganizationCustomer, - verifyPolarWebhook, type PolarBillingProvider, type PolarSubscription, } from "./polar"; @@ -445,45 +443,3 @@ describe("subscription normalization", () => { }); }); }); - -describe("Polar Standard Webhooks verification", () => { - const signingKey = "polar-webhook-secret"; - const secret = `whsec_${Buffer.from(signingKey).toString("base64")}`; - const body = JSON.stringify({ - type: "subscription.updated", - timestamp: "2026-08-09T12:00:00Z", - data: { id: "sub_1" }, - }); - const timestamp = 1_786_276_800; - const deliveryId = "evt_delivery_1"; - const signature = createHmac("sha256", signingKey) - .update(`${deliveryId}.${timestamp}.${body}`) - .digest("base64"); - const headers = { - "webhook-id": deliveryId, - "webhook-timestamp": String(timestamp), - "webhook-signature": `v1,${signature}`, - }; - - test("accepts the unmodified raw body and returns replay identity", async () => { - const verified = await verifyPolarWebhook(body, headers, secret, { - now: timestamp, - }); - expect(verified?.deliveryId).toBe(deliveryId); - expect(verified?.payload.type).toBe("subscription.updated"); - }); - - test("rejects tampering, stale delivery, and missing headers", async () => { - expect( - await verifyPolarWebhook(`${body} `, headers, secret, { now: timestamp }), - ).toBeNull(); - expect( - await verifyPolarWebhook(body, headers, secret, { - now: timestamp + 301, - }), - ).toBeNull(); - expect( - await verifyPolarWebhook(body, {}, secret, { now: timestamp }), - ).toBeNull(); - }); -}); diff --git a/packages/backend/convex/billing/polar.ts b/packages/backend/convex/billing/polar.ts index 6259f98f..5efba85b 100644 --- a/packages/backend/convex/billing/polar.ts +++ b/packages/backend/convex/billing/polar.ts @@ -676,116 +676,3 @@ export function parsePolarSubscription(value: unknown): PolarSubscription { pendingUpdate: data.pending_update ?? null, }; } - -export type PolarWebhookHeaders = - | Headers - | Record; - -export type VerifiedPolarWebhook = Readonly<{ - deliveryId: string; - timestamp: number; - payload: JsonObject; -}>; - -export type VerifyPolarWebhookOptions = Readonly<{ - now?: number; - toleranceSeconds?: number; -}>; - -/** - * Verifies Polar's Standard Webhooks signature against the unmodified body. - * The caller must persist `deliveryId` uniquely to reject valid replays. - */ -export async function verifyPolarWebhook( - rawBody: string, - headers: PolarWebhookHeaders, - webhookSecret: string | undefined, - options: VerifyPolarWebhookOptions = {}, -): Promise { - if (!webhookSecret) return null; - const deliveryId = header(headers, "webhook-id"); - const timestampHeader = header(headers, "webhook-timestamp"); - const signatureHeader = header(headers, "webhook-signature"); - if (!deliveryId || !timestampHeader || !signatureHeader) return null; - - const timestamp = Number(timestampHeader); - if (!Number.isSafeInteger(timestamp)) return null; - const now = options.now ?? Math.floor(Date.now() / 1000); - const tolerance = options.toleranceSeconds ?? 300; - if (tolerance < 0 || Math.abs(now - timestamp) > tolerance) return null; - - const signedContent = `${deliveryId}.${timestampHeader}.${rawBody}`; - const signingKey = decodeWebhookSecret(webhookSecret); - if (!signingKey) return null; - const expected = await hmacSha256(signingKey, signedContent); - const candidates = signatureHeader - .split(" ") - .map((part) => part.split(",", 2)) - .filter(([version, signature]) => version === "v1" && !!signature) - .map(([, signature]) => signature as string); - if (!candidates.some((candidate) => constantTimeEqual(candidate, expected))) { - return null; - } - - let payload: unknown; - try { - payload = JSON.parse(rawBody); - } catch { - return null; - } - return { deliveryId, timestamp, payload: object(payload, "webhook payload") }; -} - -function decodeWebhookSecret(secret: string): ArrayBuffer | null { - const encoded = secret.startsWith("whsec_") ? secret.slice(6) : secret; - try { - const binary = atob(encoded); - const bytes = new Uint8Array(new ArrayBuffer(binary.length)); - for (let index = 0; index < binary.length; index += 1) { - bytes[index] = binary.charCodeAt(index); - } - return bytes.buffer; - } catch { - return null; - } -} - -async function hmacSha256( - signingKey: ArrayBuffer, - content: string, -): Promise { - const encoder = new TextEncoder(); - const key = await crypto.subtle.importKey( - "raw", - signingKey, - { name: "HMAC", hash: "SHA-256" }, - false, - ["sign"], - ); - const digest = await crypto.subtle.sign("HMAC", key, encoder.encode(content)); - return bytesToBase64(new Uint8Array(digest)); -} - -function bytesToBase64(bytes: Uint8Array): string { - let binary = ""; - for (const byte of bytes) binary += String.fromCharCode(byte); - return btoa(binary); -} - -function constantTimeEqual(left: string, right: string): boolean { - if (left.length !== right.length) return false; - let difference = 0; - for (let index = 0; index < left.length; index += 1) { - difference |= left.charCodeAt(index) ^ right.charCodeAt(index); - } - return difference === 0; -} - -function header(headers: PolarWebhookHeaders, name: string): string | null { - if (headers instanceof Headers) return headers.get(name); - const entry = Object.entries(headers).find( - ([key]) => key.toLowerCase() === name, - )?.[1]; - if (typeof entry === "string") return entry; - return entry?.[0] ?? null; -} diff --git a/packages/backend/convex/billingDeletion.test.ts b/packages/backend/convex/billingDeletion.test.ts index 85fffde5..b96b67d8 100644 --- a/packages/backend/convex/billingDeletion.test.ts +++ b/packages/backend/convex/billingDeletion.test.ts @@ -1,8 +1,5 @@ import { describe, expect, test } from "bun:test"; -import { - processWebhook, - terminateSubscriptionForWorkspaceDeletion, -} from "./billingModel"; +import { terminateSubscriptionForWorkspaceDeletion } from "./billingModel"; type RegisteredFunction = { _handler: (ctx: unknown, args: unknown) => Promise; @@ -58,113 +55,4 @@ describe("workspace billing deletion", () => { }), }); }); - - test("ignores late provider events after the workspace is gone", async () => { - const patches: Array> = []; - const event = { - _id: "event-1", - status: "pending", - attemptCount: 0, - organizationId: "organization-deleted", - eventType: "subscription.updated", - payload: { data: { id: "subscription-1" } }, - }; - await (processWebhook as unknown as RegisteredFunction)._handler( - { - db: { - get: async () => event, - patch: async (_id: string, value: Record) => { - patches.push(value); - }, - }, - runQuery: async () => null, - }, - { eventId: "event-1" }, - ); - - expect(patches).toContainEqual( - expect.objectContaining({ - status: "ignored", - failureCode: "WORKSPACE_DELETED", - }), - ); - }); - - test("ignores unsupported events before looking up workspace metadata", async () => { - const patches: Array> = []; - let organizationLookups = 0; - const event = { - _id: "event-unsupported", - status: "failed", - attemptCount: 74, - eventType: "checkout.created", - organizationId: "checkout-contract-probe", - payload: { data: { id: "checkout-1" } }, - }; - await (processWebhook as unknown as RegisteredFunction)._handler( - { - db: { - get: async () => event, - patch: async (_id: string, value: Record) => { - patches.push(value); - }, - }, - runQuery: async () => { - organizationLookups += 1; - throw new Error("Unsupported events must not query organizations"); - }, - }, - { eventId: "event-unsupported" }, - ); - - expect(organizationLookups).toBe(0); - expect(patches).toContainEqual( - expect.objectContaining({ - status: "ignored", - failureCode: "EVENT_NOT_APPLICABLE", - attemptCount: 75, - }), - ); - }); - - test("dead-letters a permanent failure after a bounded number of attempts", async () => { - const patches: Array> = []; - let scheduledRetries = 0; - const event = { - _id: "event-poison", - status: "failed", - attemptCount: 7, - eventType: "subscription.updated", - organizationId: "organization-invalid", - payload: { data: { id: "subscription-1" } }, - }; - await (processWebhook as unknown as RegisteredFunction)._handler( - { - db: { - get: async () => event, - patch: async (_id: string, value: Record) => { - patches.push(value); - }, - }, - runQuery: async () => { - throw new Error("Permanent organization lookup failure"); - }, - scheduler: { - runAfter: async () => { - scheduledRetries += 1; - }, - }, - }, - { eventId: "event-poison" }, - ); - - expect(scheduledRetries).toBe(0); - expect(patches).toContainEqual( - expect.objectContaining({ - status: "deadLettered", - attemptCount: 8, - nextAttemptAt: undefined, - }), - ); - }); }); diff --git a/packages/backend/convex/billingModel.ts b/packages/backend/convex/billingModel.ts index 214b4249..23d0094f 100644 --- a/packages/backend/convex/billingModel.ts +++ b/packages/backend/convex/billingModel.ts @@ -6,98 +6,7 @@ import { checkoutAttemptShouldReplay, newCheckoutIntentDocument, } from "./billing/checkoutIntent"; -import { - aiTopUpAmountToCreditUnits, - moneyAmountMinorToCreditUnits, -} from "@baseblocks/domain"; -import { components, internal } from "./_generated/api"; -import { - internalMutation, - internalQuery, - type MutationCtx, -} from "./_generated/server"; -import type { Doc, Id } from "./_generated/dataModel"; -import { - grantAiCredits, - reconcileAiCreditGrantUpward, - replaceUnusedIncludedCreditLots, - revokeAiCreditGrant, -} from "./model/aiCredits"; -import { parsePolarOrderAmounts } from "./model/polarOrderAmounts"; -import { shouldApplyProviderUpdate } from "./model/billingEventOrdering"; -import { - normalizeSubscriptionLifecycle, - parsePolarSubscription, -} from "./billing/polar"; - -type JsonObject = Record; -type WebhookKind = "subscription" | "order"; - -const MAX_WEBHOOK_ATTEMPTS = 8; - -function object(value: unknown): JsonObject | null { - return value && typeof value === "object" && !Array.isArray(value) - ? (value as JsonObject) - : null; -} - -function string(value: unknown): string | undefined { - return typeof value === "string" && value.length > 0 ? value : undefined; -} - -function number(value: unknown): number | undefined { - return typeof value === "number" && Number.isFinite(value) - ? value - : undefined; -} - -function timestamp(value: unknown, fallback: number): number { - if (typeof value !== "string") return fallback; - const parsed = Date.parse(value); - return Number.isFinite(parsed) ? parsed : fallback; -} - -function findWorkspaceId(data: JsonObject): string | undefined { - const metadata = object(data.metadata); - const customer = object(data.customer); - return ( - string(metadata?.baseblocks_workspace_id) ?? - string(data.external_customer_id) ?? - string(customer?.external_id) - ); -} - -function webhookKind(eventType: string): WebhookKind | null { - if (eventType.startsWith("subscription.")) return "subscription"; - if (eventType.startsWith("order.")) return "order"; - return null; -} - -function normalizedOrderState(eventType: string, data: JsonObject) { - if (eventType === "order.refunded") { - const amounts = parsePolarOrderAmounts(data); - const amount = amounts.grossMinor; - const refunded = amounts.refundedGrossMinor; - return refunded >= amount - ? ("refunded" as const) - : ("partiallyRefunded" as const); - } - if (eventType === "order.paid" || data.status === "paid") - return "paid" as const; - if (data.status === "failed") return "failed" as const; - return "pending" as const; -} - -function pastDueGraceEnabled(data: { - normalizedStatus: string; - pastDueAt?: number; - now: number; -}) { - if (data.normalizedStatus !== "grace") return false; - const days = Number(process.env.BASEBLOCKS_PAST_DUE_GRACE_DAYS ?? "0"); - if (!Number.isSafeInteger(days) || days <= 0 || !data.pastDueAt) return false; - return data.now <= data.pastDueAt + days * 86_400_000; -} +import { internalMutation, internalQuery } from "./_generated/server"; export const configureCatalogItem = internalMutation({ args: { @@ -593,518 +502,3 @@ export const failSeatSyncOperation = internalMutation({ return null; }, }); - -export const ingestWebhook = internalMutation({ - args: { - providerEnvironment: v.union(v.literal("sandbox"), v.literal("production")), - deliveryId: v.string(), - eventType: v.string(), - eventOccurredAt: v.number(), - providerModifiedAt: v.optional(v.number()), - resourceType: v.optional(v.string()), - resourceId: v.optional(v.string()), - organizationId: v.optional(v.string()), - providerCustomerId: v.optional(v.string()), - providerSubscriptionId: v.optional(v.string()), - providerOrderId: v.optional(v.string()), - payloadHash: v.string(), - rawPayload: v.string(), - payload: v.any(), - }, - handler: async (ctx, args) => { - const existing = await ctx.db - .query("billingWebhookEvents") - .withIndex("by_environment_delivery", (q) => - q - .eq("providerEnvironment", args.providerEnvironment) - .eq("deliveryId", args.deliveryId), - ) - .unique(); - if (existing) { - if (existing.payloadHash !== args.payloadHash) { - throw new ConvexError({ - code: "WEBHOOK_REPLAY_CONFLICT", - message: "Webhook delivery ID was replayed with a different payload", - }); - } - return { eventId: existing._id, duplicate: true }; - } - const now = Date.now(); - const eventId = await ctx.db.insert("billingWebhookEvents", { - provider: "polar", - ...args, - status: "pending", - attemptCount: 0, - receivedAt: now, - updatedAt: now, - }); - return { eventId, duplicate: false }; - }, -}); - -async function upsertSubscription( - ctx: MutationCtx, - event: Doc<"billingWebhookEvents">, - data: JsonObject, - organizationId: string, -) { - const subscription = parsePolarSubscription(data); - const lifecycle = normalizeSubscriptionLifecycle(subscription); - const providerModifiedAt = timestamp( - subscription.modifiedAt, - event.providerModifiedAt ?? event.eventOccurredAt, - ); - const existing = await ctx.db - .query("billingSubscriptions") - .withIndex("by_provider_subscription", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerSubscriptionId", subscription.id), - ) - .unique(); - if ( - existing && - !shouldApplyProviderUpdate(existing.providerModifiedAt, providerModifiedAt) - ) - return existing; - const catalog = await ctx.db - .query("billingCatalogItems") - .withIndex("by_environment_product", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerProductId", subscription.productId), - ) - .first(); - if (catalog?.kind !== "plus" || !catalog.recurringInterval) { - return null; - } - await upsertCustomerFromWebhook( - ctx, - event.providerEnvironment, - organizationId, - subscription.customerId, - ); - const pendingUpdate = object(subscription.pendingUpdate); - const now = Date.now(); - const value = { - organizationId, - providerEnvironment: event.providerEnvironment, - providerSubscriptionId: subscription.id, - providerCustomerId: subscription.customerId, - providerProductId: subscription.productId, - providerPriceId: catalog.providerPriceId, - planKey: "plus" as const, - recurringInterval: catalog.recurringInterval, - providerStatus: subscription.status, - normalizedStatus: lifecycle.state, - seatQuantity: Math.max(1, subscription.seats ?? 1), - pendingSeatQuantity: number(pendingUpdate?.seats), - pendingProductId: string(pendingUpdate?.product_id), - cancelAtPeriodEnd: subscription.cancelAtPeriodEnd, - currentPeriodStart: timestamp(subscription.currentPeriodStart, now), - currentPeriodEnd: timestamp(subscription.currentPeriodEnd, now), - pastDueAt: subscription.pastDueAt - ? timestamp(subscription.pastDueAt, now) - : undefined, - canceledAt: subscription.canceledAt - ? timestamp(subscription.canceledAt, now) - : undefined, - endedAt: subscription.endedAt - ? timestamp(subscription.endedAt, now) - : undefined, - providerModifiedAt, - latestEventOccurredAt: event.eventOccurredAt, - latestWebhookDeliveryId: event.deliveryId, - updatedAt: now, - }; - const subscriptionId = existing?._id; - if (existing) await ctx.db.patch(existing._id, value); - else { - const inserted = await ctx.db.insert("billingSubscriptions", { - ...value, - createdAt: now, - }); - const created = await ctx.db.get(inserted); - if (!created) throw new Error("Subscription insert failed"); - await deriveEntitlement(ctx, event._id, created); - return created; - } - const updated = subscriptionId ? await ctx.db.get(subscriptionId) : null; - if (updated) await deriveEntitlement(ctx, event._id, updated); - return updated; -} - -async function upsertCustomerFromWebhook( - ctx: MutationCtx, - providerEnvironment: "sandbox" | "production", - organizationId: string, - providerCustomerId: string, -) { - const existing = await ctx.db - .query("billingCustomers") - .withIndex("by_organization_environment", (q) => - q - .eq("organizationId", organizationId) - .eq("providerEnvironment", providerEnvironment), - ) - .unique(); - const now = Date.now(); - const value = { - organizationId, - provider: "polar" as const, - providerEnvironment, - providerCustomerId, - externalCustomerId: organizationId, - lastSyncedAt: now, - lastErrorCode: undefined, - updatedAt: now, - }; - if (existing) await ctx.db.patch(existing._id, value); - else await ctx.db.insert("billingCustomers", { ...value, createdAt: now }); -} - -async function deriveEntitlement( - ctx: MutationCtx, - eventId: Id<"billingWebhookEvents">, - subscription: Doc<"billingSubscriptions">, -) { - const latestSeatSnapshot = await ctx.db - .query("billingSeatSnapshots") - .withIndex("by_organization_observed", (q) => - q.eq("organizationId", subscription.organizationId), - ) - .order("desc") - .first(); - const now = Date.now(); - const plusEnabled = - subscription.normalizedStatus === "entitled" || - pastDueGraceEnabled({ - normalizedStatus: subscription.normalizedStatus, - pastDueAt: subscription.pastDueAt, - now, - }); - const existing = await ctx.db - .query("workspaceEntitlements") - .withIndex("by_organization", (q) => - q.eq("organizationId", subscription.organizationId), - ) - .unique(); - const value = { - organizationId: subscription.organizationId, - plan: plusEnabled ? ("plus" as const) : ("free" as const), - subscriptionStatus: subscription.normalizedStatus, - statusReason: `polar:${subscription.providerStatus}`, - plusEnabled, - paidSeatCapacity: subscription.seatQuantity, - billableSeatCount: Math.max(1, latestSeatSnapshot?.billableSeatCount ?? 1), - sourceSubscriptionId: subscription._id, - sourceEventId: eventId, - effectiveFrom: subscription.currentPeriodStart, - effectiveThrough: - subscription.cancelAtPeriodEnd || !plusEnabled - ? subscription.currentPeriodEnd - : undefined, - policyVersion: "polar-entitlements-v1", - derivedAt: now, - updatedAt: now, - }; - if (existing) await ctx.db.patch(existing._id, value); - else await ctx.db.insert("workspaceEntitlements", value); -} - -async function processOrder( - ctx: MutationCtx, - event: Doc<"billingWebhookEvents">, - data: JsonObject, - organizationId: string, -) { - const providerOrderId = string(data.id); - const providerProductId = - string(data.product_id) ?? string(object(data.product)?.id); - const providerCustomerId = - string(data.customer_id) ?? string(object(data.customer)?.id); - if (!providerOrderId || !providerProductId || !providerCustomerId) { - throw new Error("Polar order webhook omitted required identifiers"); - } - const catalog = await ctx.db - .query("billingCatalogItems") - .withIndex("by_environment_product", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerProductId", providerProductId), - ) - .first(); - if (!catalog) return null; - const amounts = parsePolarOrderAmounts(data); - const state = normalizedOrderState(event.eventType, data); - // Polar can emit a negative subscription_update order for a downgrade or - // proration credit. Keep that signed order for reconciliation, but never - // treat it as a paid credit grant. - const grossAmountMinor = amounts.grossMinor; - const refundedGrossAmountMinor = amounts.refundedGrossMinor; - const existing = await ctx.db - .query("billingOrders") - .withIndex("by_provider_order", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerOrderId", providerOrderId), - ) - .unique(); - const now = Date.now(); - const providerModifiedAt = timestamp(data.modified_at, event.eventOccurredAt); - if ( - existing && - !shouldApplyProviderUpdate(existing.providerModifiedAt, providerModifiedAt) - ) { - return existing._id; - } - const providerSubscriptionId = string(data.subscription_id); - const billingReason = string(data.billing_reason); - const value = { - organizationId, - providerEnvironment: event.providerEnvironment, - providerOrderId, - providerCheckoutId: string(data.checkout_id), - providerCustomerId, - providerSubscriptionId, - providerProductId, - kind: - catalog.kind === "aiCreditPack" - ? ("prepaid" as const) - : ("subscription" as const), - state, - subtotalAmountMinor: amounts.subtotalMinor, - discountAmountMinor: amounts.discountMinor, - taxAmountMinor: amounts.taxMinor, - grossAmountMinor: amounts.grossMinor, - netAmountMinor: amounts.netMinor, - refundedGrossAmountMinor: amounts.refundedGrossMinor, - currency: string(data.currency) ?? catalog.currency, - billingReason, - providerModifiedAt, - updatedAt: now, - }; - let orderId = existing?._id; - if (existing) await ctx.db.patch(existing._id, value); - else { - orderId = await ctx.db.insert("billingOrders", { - ...value, - createdAt: now, - }); - } - await upsertCustomerFromWebhook( - ctx, - event.providerEnvironment, - organizationId, - providerCustomerId, - ); - const providerCheckoutId = string(data.checkout_id); - if (providerCheckoutId) { - const intent = await ctx.db - .query("billingCheckoutIntents") - .withIndex("by_provider_checkout", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerCheckoutId", providerCheckoutId), - ) - .unique(); - if (intent && intent.status !== "completed") { - await ctx.db.patch(intent._id, { status: "completed", updatedAt: now }); - } - } - if ( - state === "paid" && - grossAmountMinor > 0n && - (catalog.kind === "aiCreditPack" || - (catalog.creditUnits !== undefined && catalog.creditUnits > 0n)) - ) { - if (catalog.kind === "aiCreditPack") { - const grantedUnits = aiTopUpAmountToCreditUnits(grossAmountMinor); - const grant = { - organizationId, - bucket: "prepaid", - sourceKind: "purchase", - sourceRef: `prepaid:${providerOrderId}`, - units: grantedUnits, - policyVersion: catalog.configurationVersion, - billingEventId: event._id, - now, - } as const; - const lotId = await grantAiCredits(ctx, grant); - await reconcileAiCreditGrantUpward(ctx, grant); - if (orderId) await ctx.db.patch(orderId, { creditLotId: lotId }); - } else if (providerSubscriptionId) { - const recurringUnits = catalog.creditUnits; - if (!recurringUnits || recurringUnits <= 0n) { - throw new Error("Plus catalog item has no included AI credits"); - } - const subscription = await ctx.db - .query("billingSubscriptions") - .withIndex("by_provider_subscription", (q) => - q - .eq("providerEnvironment", event.providerEnvironment) - .eq("providerSubscriptionId", providerSubscriptionId), - ) - .unique(); - if (subscription?.normalizedStatus !== "entitled") { - throw new Error( - "Paid subscription order is waiting for authoritative subscription state", - ); - } - if (billingReason === "subscription_update") { - await replaceUnusedIncludedCreditLots(ctx, { - organizationId, - replacementRef: providerOrderId, - preserveSourceRef: `included:${providerSubscriptionId}:${subscription.currentPeriodStart}`, - policyVersion: catalog.configurationVersion, - billingEventId: event._id, - now, - }); - } - await grantAiCredits(ctx, { - organizationId, - bucket: "included", - sourceKind: "recurring", - sourceRef: `included:${providerSubscriptionId}:${subscription.currentPeriodStart}`, - units: recurringUnits, - periodStart: subscription.currentPeriodStart, - periodEnd: subscription.currentPeriodEnd, - expiresAt: subscription.currentPeriodEnd, - policyVersion: catalog.configurationVersion, - billingEventId: event._id, - now, - }); - } - } - if ( - (state === "refunded" || state === "partiallyRefunded") && - catalog.kind === "aiCreditPack" && - grossAmountMinor > 0n - ) { - const targetRevokedUnits = moneyAmountMinorToCreditUnits( - refundedGrossAmountMinor, - ); - await revokeAiCreditGrant(ctx, { - organizationId, - sourceRef: `prepaid:${providerOrderId}`, - targetRevokedUnits, - policyVersion: catalog.configurationVersion, - billingEventId: event._id, - now, - }); - } - return orderId; -} - -export const processWebhook = internalMutation({ - args: { eventId: v.id("billingWebhookEvents") }, - handler: async (ctx, { eventId }) => { - const event = await ctx.db.get(eventId); - if ( - !event || - event.status === "processed" || - event.status === "ignored" || - event.status === "deadLettered" - ) { - return null; - } - const now = Date.now(); - const kind = webhookKind(event.eventType); - if (!kind) { - await ctx.db.patch(event._id, { - status: "ignored", - attemptCount: event.attemptCount + 1, - failureCode: "EVENT_NOT_APPLICABLE", - failureMessage: undefined, - processedAt: now, - nextAttemptAt: undefined, - updatedAt: now, - }); - return null; - } - if ( - event.status === "failed" && - event.nextAttemptAt !== undefined && - event.nextAttemptAt > now - ) { - return null; - } - try { - const payload = object(event.payload); - const data = object(payload?.data); - if (!data) throw new Error("Polar webhook omitted data"); - const organizationId = event.organizationId ?? findWorkspaceId(data); - if (!organizationId) { - await ctx.db.patch(event._id, { - status: "ignored", - attemptCount: event.attemptCount + 1, - failureCode: "WORKSPACE_ID_MISSING", - processedAt: now, - nextAttemptAt: undefined, - updatedAt: now, - }); - return null; - } - const organization = await ctx.runQuery( - components.betterAuth.adapter.findOne, - { - model: "organization", - where: [{ field: "_id", operator: "eq", value: organizationId }], - }, - ); - if (!organization) { - await ctx.db.patch(event._id, { - status: "ignored", - organizationId, - attemptCount: event.attemptCount + 1, - failureCode: "WORKSPACE_DELETED", - processedAt: now, - nextAttemptAt: undefined, - updatedAt: now, - }); - return null; - } - if (kind === "subscription") { - await upsertSubscription(ctx, event, data, organizationId); - } else { - await processOrder(ctx, event, data, organizationId); - } - await ctx.db.patch(event._id, { - status: "processed", - organizationId, - attemptCount: event.attemptCount + 1, - processedAt: Date.now(), - nextAttemptAt: undefined, - failureCode: undefined, - failureMessage: undefined, - updatedAt: Date.now(), - }); - } catch (error) { - const attemptCount = event.attemptCount + 1; - const retryDelay = Math.min( - 60 * 60_000, - 60_000 * 2 ** Math.min(event.attemptCount, 6), - ); - const terminal = attemptCount >= MAX_WEBHOOK_ATTEMPTS; - await ctx.db.patch(event._id, { - status: terminal ? "deadLettered" : "failed", - attemptCount, - nextAttemptAt: terminal ? undefined : Date.now() + retryDelay, - failureCode: "PROCESSING_FAILED", - failureMessage: - error instanceof Error - ? error.message.slice(0, 1_000) - : "Unknown error", - updatedAt: Date.now(), - }); - if (!terminal) { - await ctx.scheduler.runAfter( - retryDelay, - internal.billingModel.processWebhook, - { eventId: event._id }, - ); - } - } - return null; - }, -}); diff --git a/packages/backend/convex/billingWebhooks.ts b/packages/backend/convex/billingWebhooks.ts deleted file mode 100644 index 9f689b15..00000000 --- a/packages/backend/convex/billingWebhooks.ts +++ /dev/null @@ -1,115 +0,0 @@ -import { makeFunctionReference } from "convex/server"; -import type { Id } from "./_generated/dataModel"; -import { httpAction } from "./_generated/server"; -import { - polarEnvironmentFromEnvironment, - verifyPolarWebhook, -} from "./billing/polar"; - -type JsonObject = Record; -type ProviderEnvironment = "sandbox" | "production"; - -const ingestWebhook = makeFunctionReference< - "mutation", - { - providerEnvironment: ProviderEnvironment; - deliveryId: string; - eventType: string; - eventOccurredAt: number; - providerModifiedAt?: number; - resourceType?: string; - resourceId?: string; - organizationId?: string; - providerCustomerId?: string; - providerSubscriptionId?: string; - providerOrderId?: string; - payloadHash: string; - rawPayload: string; - payload: unknown; - }, - { eventId: Id<"billingWebhookEvents">; duplicate: boolean } ->("billingModel:ingestWebhook"); -const processWebhook = makeFunctionReference< - "mutation", - { eventId: Id<"billingWebhookEvents"> }, - null ->("billingModel:processWebhook"); - -function object(value: unknown): JsonObject | undefined { - return value && typeof value === "object" && !Array.isArray(value) - ? (value as JsonObject) - : undefined; -} - -function string(value: unknown): string | undefined { - return typeof value === "string" && value.length > 0 ? value : undefined; -} - -function milliseconds(value: unknown): number | undefined { - if (typeof value !== "string") return undefined; - const parsed = Date.parse(value); - return Number.isFinite(parsed) ? parsed : undefined; -} - -async function sha256Hex(value: string): Promise { - const digest = await crypto.subtle.digest( - "SHA-256", - new TextEncoder().encode(value), - ); - return Array.from(new Uint8Array(digest), (byte) => - byte.toString(16).padStart(2, "0"), - ).join(""); -} - -export const handlePolarWebhook = httpAction(async (ctx, request) => { - let providerEnvironment: ProviderEnvironment; - try { - providerEnvironment = polarEnvironmentFromEnvironment(); - } catch { - return new Response("Billing webhook is not configured", { status: 503 }); - } - const rawPayload = await request.text(); - if (rawPayload.length === 0 || rawPayload.length > 2_000_000) { - return new Response("Invalid webhook payload", { status: 413 }); - } - const verified = await verifyPolarWebhook( - rawPayload, - request.headers, - process.env.POLAR_WEBHOOK_SECRET, - ); - if (!verified) - return new Response("Invalid webhook signature", { status: 401 }); - - const eventType = string(verified.payload.type); - const data = object(verified.payload.data); - if (!eventType || !data) - return new Response("Malformed webhook payload", { status: 400 }); - const metadata = object(data.metadata); - const customer = object(data.customer); - const subscription = object(data.subscription); - const resourceType = eventType.split(".", 1)[0]; - const event = await ctx.runMutation(ingestWebhook, { - providerEnvironment, - deliveryId: verified.deliveryId, - eventType, - eventOccurredAt: - milliseconds(verified.payload.timestamp) ?? verified.timestamp * 1_000, - providerModifiedAt: milliseconds(data.modified_at), - resourceType, - resourceId: string(data.id), - organizationId: - string(metadata?.baseblocks_workspace_id) ?? - string(data.external_customer_id) ?? - string(customer?.external_id), - providerCustomerId: string(data.customer_id) ?? string(customer?.id), - providerSubscriptionId: - string(data.subscription_id) ?? string(subscription?.id), - providerOrderId: resourceType === "order" ? string(data.id) : undefined, - payloadHash: await sha256Hex(rawPayload), - rawPayload, - payload: verified.payload, - }); - if (!event.duplicate) - await ctx.scheduler.runAfter(0, processWebhook, { eventId: event.eventId }); - return new Response(null, { status: 202 }); -}); diff --git a/packages/backend/convex/billing_webhook_model.ts b/packages/backend/convex/billing_webhook_model.ts new file mode 100644 index 00000000..8e723747 --- /dev/null +++ b/packages/backend/convex/billing_webhook_model.ts @@ -0,0 +1,540 @@ +import { + aiTopUpAmountToCreditUnits, + moneyAmountMinorToCreditUnits, +} from "@baseblocks/domain"; +import { type Infer, v } from "convex/values"; +import { internalMutation, type MutationCtx } from "./_generated/server"; +import { normalizeSubscriptionLifecycle } from "./billing/polar"; +import { + grantAiCredits, + reconcileAiCreditGrantUpward, + replaceUnusedIncludedCreditLots, + revokeAiCreditGrant, +} from "./model/aiCredits"; + +const providerEnvironment = v.union( + v.literal("sandbox"), + v.literal("production"), +); +const orderEvent = v.object({ + kind: v.literal("order"), + organizationId: v.string(), + providerOrderId: v.string(), + providerCheckoutId: v.optional(v.string()), + providerCustomerId: v.string(), + providerSubscriptionId: v.optional(v.string()), + providerProductId: v.string(), + state: v.union( + v.literal("pending"), + v.literal("paid"), + v.literal("refunded"), + v.literal("partiallyRefunded"), + v.literal("failed"), + ), + subtotalAmountMinor: v.int64(), + discountAmountMinor: v.int64(), + taxAmountMinor: v.int64(), + grossAmountMinor: v.int64(), + netAmountMinor: v.int64(), + refundedGrossAmountMinor: v.int64(), + currency: v.string(), + billingReason: v.optional(v.string()), + providerModifiedAt: v.number(), +}); +const subscriptionEvent = v.object({ + kind: v.literal("subscription"), + organizationId: v.string(), + providerSubscriptionId: v.string(), + providerCustomerId: v.string(), + providerProductId: v.string(), + providerStatus: v.string(), + seatQuantity: v.number(), + pendingSeatQuantity: v.optional(v.number()), + pendingProductId: v.optional(v.string()), + cancelAtPeriodEnd: v.boolean(), + pauseAtPeriodEnd: v.optional(v.boolean()), + currentPeriodStart: v.number(), + currentPeriodEnd: v.number(), + pastDueAt: v.optional(v.number()), + canceledAt: v.optional(v.number()), + endedAt: v.optional(v.number()), + providerModifiedAt: v.number(), +}); +const commandValidator = { + providerEnvironment, + deliveryId: v.string(), + eventType: v.union( + v.literal("order.created"), + v.literal("order.updated"), + v.literal("order.paid"), + v.literal("order.refunded"), + v.literal("subscription.active"), + v.literal("subscription.canceled"), + v.literal("subscription.created"), + v.literal("subscription.past_due"), + v.literal("subscription.revoked"), + v.literal("subscription.uncanceled"), + v.literal("subscription.updated"), + ), + eventOccurredAt: v.number(), + event: v.union(orderEvent, subscriptionEvent), +}; + +type PolarOrderEvent = Infer; +type PolarSubscriptionEvent = Infer; +export type PolarBillingEventCommand = Readonly<{ + [Key in keyof typeof commandValidator]: Infer<(typeof commandValidator)[Key]>; +}>; + +function assertEventTypeMatchesResource(command: PolarBillingEventCommand) { + if ( + (command.event.kind === "order" && + !command.eventType.startsWith("order.")) || + (command.event.kind === "subscription" && + !command.eventType.startsWith("subscription.")) + ) { + throw new Error("Polar event type does not match its billing resource"); + } +} + +function pastDueGraceEnabled(data: { + normalizedStatus: string; + pastDueAt?: number; + now: number; +}) { + if (data.normalizedStatus !== "grace") return false; + const days = Number(process.env.BASEBLOCKS_PAST_DUE_GRACE_DAYS ?? "0"); + return ( + Number.isSafeInteger(days) && + days > 0 && + data.pastDueAt !== undefined && + data.now <= data.pastDueAt + days * 86_400_000 + ); +} + +function resourceId(event: PolarOrderEvent | PolarSubscriptionEvent) { + return event.kind === "order" + ? event.providerOrderId + : event.providerSubscriptionId; +} + +async function upsertCustomer( + ctx: MutationCtx, + command: PolarBillingEventCommand, +) { + const existing = await ctx.db + .query("billingCustomers") + .withIndex("by_organization_environment", (query) => + query + .eq("organizationId", command.event.organizationId) + .eq("providerEnvironment", command.providerEnvironment), + ) + .unique(); + const now = Date.now(); + const value = { + organizationId: command.event.organizationId, + provider: "polar" as const, + providerEnvironment: command.providerEnvironment, + providerCustomerId: command.event.providerCustomerId, + externalCustomerId: command.event.organizationId, + lastSyncedAt: now, + lastErrorCode: undefined, + updatedAt: now, + }; + if (existing) await ctx.db.patch(existing._id, value); + else await ctx.db.insert("billingCustomers", { ...value, createdAt: now }); +} + +export async function applyPolarBillingEvent( + ctx: MutationCtx, + command: PolarBillingEventCommand, +) { + assertEventTypeMatchesResource(command); + const existingDelivery = await ctx.db + .query("billingWebhookEvents") + .withIndex("by_environment_delivery", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("deliveryId", command.deliveryId), + ) + .unique(); + if (existingDelivery) { + if ( + existingDelivery.eventType !== command.eventType || + existingDelivery.resourceId !== resourceId(command.event) + ) { + throw new Error("Polar delivery ID was reused for another event"); + } + return { + outcome: "duplicate" as const, + eventType: command.eventType, + resourceId: resourceId(command.event), + }; + } + + const workspace = await ctx.db + .query("workspaceProfiles") + .withIndex("by_organization", (query) => + query.eq("organizationId", command.event.organizationId), + ) + .unique(); + if (!workspace) { + const now = Date.now(); + await ctx.db.insert("billingWebhookEvents", { + providerEnvironment: command.providerEnvironment, + deliveryId: command.deliveryId, + eventType: command.eventType, + eventOccurredAt: command.eventOccurredAt, + resourceId: resourceId(command.event), + organizationId: command.event.organizationId, + outcome: "ignored", + ignoredReason: "workspace_not_found", + receivedAt: now, + processedAt: now, + }); + return { + outcome: "ignored" as const, + eventType: command.eventType, + resourceId: resourceId(command.event), + }; + } + + const catalog = await ctx.db + .query("billingCatalogItems") + .withIndex("by_environment_product", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("providerProductId", command.event.providerProductId), + ) + .first(); + if (!catalog) { + throw new Error("Polar event references an unknown billing product"); + } + + const now = Date.now(); + const billingEventId = await ctx.db.insert("billingWebhookEvents", { + providerEnvironment: command.providerEnvironment, + deliveryId: command.deliveryId, + eventType: command.eventType, + eventOccurredAt: command.eventOccurredAt, + resourceId: resourceId(command.event), + organizationId: command.event.organizationId, + outcome: "applied", + receivedAt: now, + processedAt: now, + }); + + if (command.event.kind === "subscription") { + if (catalog.kind !== "plus" || !catalog.recurringInterval) { + throw new Error("Polar subscription references a non-recurring product"); + } + const lifecycle = normalizeSubscriptionLifecycle({ + status: command.event.providerStatus, + cancelAtPeriodEnd: command.event.cancelAtPeriodEnd, + pauseAtPeriodEnd: command.event.pauseAtPeriodEnd ?? false, + currentPeriodEnd: new Date(command.event.currentPeriodEnd).toISOString(), + endedAt: command.event.endedAt + ? new Date(command.event.endedAt).toISOString() + : null, + }); + const existingSubscription = await ctx.db + .query("billingSubscriptions") + .withIndex("by_provider_subscription", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("providerSubscriptionId", command.event.providerSubscriptionId!), + ) + .unique(); + if ( + existingSubscription && + command.event.providerModifiedAt < existingSubscription.providerModifiedAt + ) { + await ctx.db.patch(billingEventId, { + outcome: "ignored", + ignoredReason: "stale_provider_update", + }); + return { + outcome: "ignored" as const, + eventType: command.eventType, + resourceId: command.event.providerSubscriptionId, + }; + } + await upsertCustomer(ctx, command); + const subscriptionValue = { + organizationId: command.event.organizationId, + providerEnvironment: command.providerEnvironment, + providerSubscriptionId: command.event.providerSubscriptionId, + providerCustomerId: command.event.providerCustomerId, + providerProductId: command.event.providerProductId, + providerPriceId: catalog.providerPriceId, + planKey: "plus" as const, + recurringInterval: catalog.recurringInterval, + providerStatus: command.event.providerStatus, + normalizedStatus: lifecycle.state, + seatQuantity: Math.max(1, command.event.seatQuantity), + pendingSeatQuantity: command.event.pendingSeatQuantity, + pendingProductId: command.event.pendingProductId, + cancelAtPeriodEnd: command.event.cancelAtPeriodEnd, + currentPeriodStart: command.event.currentPeriodStart, + currentPeriodEnd: command.event.currentPeriodEnd, + pastDueAt: command.event.pastDueAt, + canceledAt: command.event.canceledAt, + endedAt: command.event.endedAt, + providerModifiedAt: command.event.providerModifiedAt, + latestEventOccurredAt: command.eventOccurredAt, + latestWebhookDeliveryId: command.deliveryId, + updatedAt: now, + }; + let subscriptionId = existingSubscription?._id; + if ( + existingSubscription && + command.event.providerModifiedAt >= + existingSubscription.providerModifiedAt + ) { + await ctx.db.patch(existingSubscription._id, subscriptionValue); + } else if (!existingSubscription) { + subscriptionId = await ctx.db.insert("billingSubscriptions", { + ...subscriptionValue, + createdAt: now, + }); + } + if (!subscriptionId) throw new Error("Polar subscription was not stored"); + const latestSeatSnapshot = await ctx.db + .query("billingSeatSnapshots") + .withIndex("by_organization_observed", (query) => + query.eq("organizationId", command.event.organizationId), + ) + .order("desc") + .first(); + const entitlement = await ctx.db + .query("workspaceEntitlements") + .withIndex("by_organization", (query) => + query.eq("organizationId", command.event.organizationId), + ) + .unique(); + const plusEnabled = + lifecycle.state === "entitled" || + pastDueGraceEnabled({ + normalizedStatus: lifecycle.state, + pastDueAt: command.event.pastDueAt, + now, + }); + const entitlementValue = { + organizationId: command.event.organizationId, + plan: plusEnabled ? ("plus" as const) : ("free" as const), + subscriptionStatus: lifecycle.state, + statusReason: `polar:${command.event.providerStatus}`, + plusEnabled, + paidSeatCapacity: Math.max(1, command.event.seatQuantity), + billableSeatCount: Math.max( + 1, + latestSeatSnapshot?.billableSeatCount ?? 1, + ), + sourceSubscriptionId: subscriptionId, + sourceEventId: billingEventId, + effectiveFrom: command.event.currentPeriodStart, + effectiveThrough: + command.event.cancelAtPeriodEnd || !plusEnabled + ? command.event.currentPeriodEnd + : undefined, + policyVersion: "polar-entitlements-v1", + derivedAt: now, + updatedAt: now, + }; + if (entitlement) await ctx.db.patch(entitlement._id, entitlementValue); + else { + await ctx.db.insert("workspaceEntitlements", entitlementValue); + } + return { + outcome: "applied" as const, + eventType: command.eventType, + resourceId: command.event.providerSubscriptionId, + }; + } + + const orderEvent = command.event; + const existingOrder = await ctx.db + .query("billingOrders") + .withIndex("by_provider_order", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("providerOrderId", orderEvent.providerOrderId), + ) + .unique(); + if ( + existingOrder && + orderEvent.providerModifiedAt < existingOrder.providerModifiedAt + ) { + await ctx.db.patch(billingEventId, { + outcome: "ignored", + ignoredReason: "stale_provider_update", + }); + return { + outcome: "ignored" as const, + eventType: command.eventType, + resourceId: orderEvent.providerOrderId, + }; + } + await upsertCustomer(ctx, command); + if ( + orderEvent.providerCheckoutId && + (orderEvent.state === "paid" || + orderEvent.state === "refunded" || + orderEvent.state === "partiallyRefunded") + ) { + const checkoutIntent = await ctx.db + .query("billingCheckoutIntents") + .withIndex("by_provider_checkout", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("providerCheckoutId", orderEvent.providerCheckoutId!), + ) + .unique(); + if (checkoutIntent && checkoutIntent.status !== "completed") { + await ctx.db.patch(checkoutIntent._id, { + status: "completed", + updatedAt: now, + }); + } + } + const orderValue = { + organizationId: command.event.organizationId, + providerEnvironment: command.providerEnvironment, + providerOrderId: command.event.providerOrderId, + providerCheckoutId: command.event.providerCheckoutId, + providerCustomerId: command.event.providerCustomerId, + providerSubscriptionId: command.event.providerSubscriptionId, + providerProductId: command.event.providerProductId, + kind: + catalog.kind === "aiCreditPack" + ? ("prepaid" as const) + : ("subscription" as const), + state: command.event.state, + subtotalAmountMinor: command.event.subtotalAmountMinor, + discountAmountMinor: command.event.discountAmountMinor, + taxAmountMinor: command.event.taxAmountMinor, + grossAmountMinor: command.event.grossAmountMinor, + netAmountMinor: command.event.netAmountMinor, + refundedGrossAmountMinor: command.event.refundedGrossAmountMinor, + currency: command.event.currency, + billingReason: command.event.billingReason, + providerModifiedAt: command.event.providerModifiedAt, + updatedAt: now, + }; + const orderId = existingOrder?._id; + if ( + existingOrder && + command.event.providerModifiedAt >= existingOrder.providerModifiedAt + ) { + await ctx.db.patch(existingOrder._id, orderValue); + } + const persistedOrderId = + orderId ?? + (await ctx.db.insert("billingOrders", { + ...orderValue, + createdAt: now, + })); + + if ( + catalog.kind === "aiCreditPack" && + command.event.state === "paid" && + command.event.grossAmountMinor > 0n + ) { + const sourceRef = `prepaid:${command.event.providerOrderId}`; + const lotId = await grantAiCredits(ctx, { + organizationId: command.event.organizationId, + bucket: "prepaid", + sourceKind: "purchase", + sourceRef, + units: aiTopUpAmountToCreditUnits(command.event.grossAmountMinor), + policyVersion: catalog.configurationVersion, + billingEventId, + now, + }); + await reconcileAiCreditGrantUpward(ctx, { + organizationId: command.event.organizationId, + bucket: "prepaid", + sourceRef, + units: aiTopUpAmountToCreditUnits(command.event.grossAmountMinor), + policyVersion: catalog.configurationVersion, + billingEventId, + now, + }); + await ctx.db.patch(persistedOrderId, { creditLotId: lotId }); + } + if ( + catalog.kind === "plus" && + command.event.state === "paid" && + command.event.providerSubscriptionId && + catalog.creditUnits !== undefined && + catalog.creditUnits > 0n + ) { + const subscription = await ctx.db + .query("billingSubscriptions") + .withIndex("by_provider_subscription", (query) => + query + .eq("providerEnvironment", command.providerEnvironment) + .eq("providerSubscriptionId", command.event.providerSubscriptionId!), + ) + .unique(); + if (subscription?.normalizedStatus !== "entitled") { + throw new Error( + "Paid subscription order is waiting for authoritative subscription state", + ); + } + const sourceRef = `included:${command.event.providerSubscriptionId}:${subscription.currentPeriodStart}`; + if (command.event.billingReason === "subscription_update") { + await replaceUnusedIncludedCreditLots(ctx, { + organizationId: command.event.organizationId, + replacementRef: command.event.providerOrderId, + preserveSourceRef: sourceRef, + policyVersion: catalog.configurationVersion, + billingEventId, + now, + }); + } + const lotId = await grantAiCredits(ctx, { + organizationId: command.event.organizationId, + bucket: "included", + sourceKind: "recurring", + sourceRef, + units: catalog.creditUnits, + periodStart: subscription.currentPeriodStart, + periodEnd: subscription.currentPeriodEnd, + expiresAt: subscription.currentPeriodEnd, + policyVersion: catalog.configurationVersion, + billingEventId, + now, + }); + await ctx.db.patch(persistedOrderId, { creditLotId: lotId }); + } + if ( + catalog.kind === "aiCreditPack" && + (command.event.state === "refunded" || + command.event.state === "partiallyRefunded") && + command.event.grossAmountMinor > 0n + ) { + await revokeAiCreditGrant(ctx, { + organizationId: command.event.organizationId, + sourceRef: `prepaid:${command.event.providerOrderId}`, + targetRevokedUnits: moneyAmountMinorToCreditUnits( + command.event.refundedGrossAmountMinor, + ), + policyVersion: catalog.configurationVersion, + billingEventId, + now, + }); + } + + return { + outcome: "applied" as const, + eventType: command.eventType, + resourceId: resourceId(command.event), + }; +} + +/** The only durable entry point for verified Polar billing events. */ +export const apply = internalMutation({ + args: commandValidator, + handler: applyPolarBillingEvent, +}); diff --git a/packages/backend/convex/billing_webhooks.ts b/packages/backend/convex/billing_webhooks.ts new file mode 100644 index 00000000..bd1a1abd --- /dev/null +++ b/packages/backend/convex/billing_webhooks.ts @@ -0,0 +1,216 @@ +import { + validateEvent, + WebhookVerificationError, +} from "@polar-sh/sdk/webhooks"; +import { Buffer as BufferPolyfill } from "buffer"; +import { makeFunctionReference } from "convex/server"; +import { httpAction } from "./_generated/server"; +import { polarEnvironmentFromEnvironment } from "./billing/polar"; +import type { PolarBillingEventCommand } from "./billing_webhook_model"; + +// Polar's official Convex component installs this polyfill because Convex +// HTTP actions use the default runtime and Polar's verifier uses Buffer. +globalThis.Buffer = BufferPolyfill; + +type ValidatedPolarEvent = ReturnType; + +const applyBillingEvent = makeFunctionReference< + "mutation", + PolarBillingEventCommand, + { + outcome: "applied" | "duplicate" | "ignored"; + eventType: string; + resourceId: string; + } +>("billing_webhook_model:apply"); + +function requiredWorkspaceId( + metadata: Record, + externalCustomerId: string | null | undefined, +) { + const metadataId = metadata.baseblocks_workspace_id; + const organizationId = + typeof metadataId === "string" && metadataId.length > 0 + ? metadataId + : externalCustomerId; + if (!organizationId) { + throw new Error("Polar billing event has no BaseBlocks workspace ID"); + } + return organizationId; +} + +function time(value: Date | null | undefined, fallback: number) { + return value?.getTime() ?? fallback; +} + +function integer(value: number, field: string) { + if (!Number.isSafeInteger(value)) { + throw new Error(`Polar ${field} is not a safe integer`); + } + return BigInt(value); +} + +function errorSummary(error: unknown) { + return error instanceof Error + ? { name: error.name, message: error.message.slice(0, 500) } + : { name: "UnknownError", message: "Non-Error value thrown" }; +} + +function polarWebhookToCommand( + verified: ValidatedPolarEvent, + deliveryId: string, + providerEnvironment: "sandbox" | "production", +): PolarBillingEventCommand | null { + const eventOccurredAt = verified.timestamp.getTime(); + switch (verified.type) { + case "order.created": + case "order.updated": + case "order.paid": + case "order.refunded": { + const data = verified.data; + if (!data.productId) { + throw new Error("Polar order has no product ID"); + } + const gross = integer(data.totalAmount, "order total"); + const refundedGross = + integer(data.refundedAmount, "refunded amount") + + integer(data.refundedTaxAmount, "refunded tax amount"); + const state = + verified.type === "order.refunded" + ? refundedGross >= gross + ? ("refunded" as const) + : ("partiallyRefunded" as const) + : verified.type === "order.paid" || data.status === "paid" + ? ("paid" as const) + : data.status === "failed" + ? ("failed" as const) + : ("pending" as const); + return { + providerEnvironment, + deliveryId, + eventType: verified.type, + eventOccurredAt, + event: { + kind: "order", + organizationId: requiredWorkspaceId( + data.metadata, + data.customer.externalId, + ), + providerOrderId: data.id, + providerCheckoutId: data.checkoutId ?? undefined, + providerCustomerId: data.customerId, + providerSubscriptionId: data.subscriptionId ?? undefined, + providerProductId: data.productId, + state, + subtotalAmountMinor: integer(data.subtotalAmount, "subtotal"), + discountAmountMinor: integer(data.discountAmount, "discount"), + taxAmountMinor: integer(data.taxAmount, "tax"), + grossAmountMinor: gross, + netAmountMinor: integer(data.netAmount, "net amount"), + refundedGrossAmountMinor: refundedGross, + currency: data.currency, + billingReason: data.billingReason, + providerModifiedAt: time(data.modifiedAt, eventOccurredAt), + }, + }; + } + case "subscription.active": + case "subscription.canceled": + case "subscription.created": + case "subscription.past_due": + case "subscription.revoked": + case "subscription.uncanceled": + case "subscription.updated": { + const data = verified.data; + return { + providerEnvironment, + deliveryId, + eventType: verified.type, + eventOccurredAt, + event: { + kind: "subscription", + organizationId: requiredWorkspaceId( + data.metadata, + data.customer.externalId, + ), + providerSubscriptionId: data.id, + providerCustomerId: data.customerId, + providerProductId: data.productId, + providerStatus: data.status, + seatQuantity: data.seats ?? 1, + pendingSeatQuantity: data.pendingUpdate?.seats ?? undefined, + pendingProductId: data.pendingUpdate?.productId ?? undefined, + cancelAtPeriodEnd: data.cancelAtPeriodEnd, + pauseAtPeriodEnd: data.pauseAtPeriodEnd, + currentPeriodStart: data.currentPeriodStart.getTime(), + currentPeriodEnd: data.currentPeriodEnd.getTime(), + pastDueAt: data.pastDueAt?.getTime(), + canceledAt: data.canceledAt?.getTime(), + endedAt: data.endedAt?.getTime(), + providerModifiedAt: time(data.modifiedAt, eventOccurredAt), + }, + }; + } + default: + return null; + } +} + +export const handlePolarWebhook = httpAction(async (ctx, request) => { + let providerEnvironment: "sandbox" | "production"; + const webhookSecret = process.env.POLAR_WEBHOOK_SECRET; + try { + providerEnvironment = polarEnvironmentFromEnvironment(); + if (!webhookSecret) throw new Error("Missing POLAR_WEBHOOK_SECRET"); + } catch { + return new Response("Billing webhook is not configured", { status: 503 }); + } + + const rawBody = await request.text(); + if (rawBody.length === 0 || rawBody.length > 2_000_000) { + return new Response("Invalid webhook payload", { status: 413 }); + } + + let verified: ValidatedPolarEvent; + try { + verified = validateEvent( + rawBody, + Object.fromEntries(request.headers.entries()), + webhookSecret, + ); + } catch (error) { + if (error instanceof WebhookVerificationError) { + return new Response("Invalid webhook signature", { status: 403 }); + } + // biome-ignore lint/suspicious/noConsole: Convex captures production function errors in its dashboard logs. + console.error( + "Polar webhook payload validation failed", + errorSummary(error), + ); + return new Response("Invalid Polar webhook payload", { status: 400 }); + } + + const deliveryId = request.headers.get("webhook-id"); + if (!deliveryId) { + return new Response("Missing Polar delivery ID", { status: 400 }); + } + + try { + const command = polarWebhookToCommand( + verified, + deliveryId, + providerEnvironment, + ); + if (command) await ctx.runMutation(applyBillingEvent, command); + return new Response(null, { status: 202 }); + } catch (error) { + // biome-ignore lint/suspicious/noConsole: Convex captures production function errors in its dashboard logs. + console.error( + "Polar billing event application failed", + errorSummary(error), + ); + return new Response("Polar billing event could not be applied", { + status: 500, + }); + } +}); diff --git a/packages/backend/convex/convex.config.ts b/packages/backend/convex/convex.config.ts index 0b4209d6..8b6ebb1b 100644 --- a/packages/backend/convex/convex.config.ts +++ b/packages/backend/convex/convex.config.ts @@ -7,6 +7,15 @@ import workpool from "@convex-dev/workpool/convex.config.js"; const app = defineApp({ env: { INTEGRATIONS_ENABLED: v.union(v.literal("true"), v.literal("false")), + BASEBLOCKS_BILLING_ENVIRONMENT: v.optional( + v.union(v.literal("sandbox"), v.literal("production")), + ), + POLAR_ACCESS_TOKEN: v.optional(v.string()), + POLAR_WEBHOOK_SECRET: v.optional(v.string()), + POLAR_ALLOW_PRODUCTION: v.optional( + v.union(v.literal("true"), v.literal("false")), + ), + BASEBLOCKS_PAST_DUE_GRACE_DAYS: v.optional(v.string()), }, }); app.use(betterAuth); diff --git a/packages/backend/convex/http.ts b/packages/backend/convex/http.ts index 01584fd7..2d1f042e 100644 --- a/packages/backend/convex/http.ts +++ b/packages/backend/convex/http.ts @@ -1,7 +1,7 @@ import { httpRouter } from "convex/server"; import { authComponent, createAuth } from "./auth"; import { handleNangoWebhook } from "./integrationWebhooks"; -import { handlePolarWebhook } from "./billingWebhooks"; +import { handlePolarWebhook } from "./billing_webhooks"; const http = httpRouter(); diff --git a/packages/backend/convex/model/billingEventOrdering.test.ts b/packages/backend/convex/model/billingEventOrdering.test.ts deleted file mode 100644 index d9f6875c..00000000 --- a/packages/backend/convex/model/billingEventOrdering.test.ts +++ /dev/null @@ -1,14 +0,0 @@ -import { describe, expect, test } from "bun:test"; -import { shouldApplyProviderUpdate } from "./billingEventOrdering"; - -describe("billing provider event ordering", () => { - test("accepts the first, newer, and idempotent same-version event", () => { - expect(shouldApplyProviderUpdate(undefined, 100)).toBe(true); - expect(shouldApplyProviderUpdate(100, 101)).toBe(true); - expect(shouldApplyProviderUpdate(100, 100)).toBe(true); - }); - - test("rejects an out-of-order provider event", () => { - expect(shouldApplyProviderUpdate(200, 199)).toBe(false); - }); -}); diff --git a/packages/backend/convex/model/billingEventOrdering.ts b/packages/backend/convex/model/billingEventOrdering.ts deleted file mode 100644 index 0b2d3692..00000000 --- a/packages/backend/convex/model/billingEventOrdering.ts +++ /dev/null @@ -1,8 +0,0 @@ -export function shouldApplyProviderUpdate( - currentModifiedAt: number | undefined, - incomingModifiedAt: number, -): boolean { - return ( - currentModifiedAt === undefined || incomingModifiedAt >= currentModifiedAt - ); -} diff --git a/packages/backend/convex/model/polarOrderAmounts.test.ts b/packages/backend/convex/model/polarOrderAmounts.test.ts deleted file mode 100644 index 489f8658..00000000 --- a/packages/backend/convex/model/polarOrderAmounts.test.ts +++ /dev/null @@ -1,40 +0,0 @@ -import { describe, expect, test } from "bun:test"; -import { parsePolarOrderAmounts } from "./polarOrderAmounts"; - -describe("parsePolarOrderAmounts", () => { - test("uses the gross customer payment rather than the tax-exclusive amount", () => { - expect( - parsePolarOrderAmounts({ - amount: 417, - subtotal_amount: 500, - discount_amount: 0, - tax_amount: 83, - total_amount: 500, - net_amount: 417, - refunded_amount: 0, - refunded_tax_amount: 0, - }), - ).toEqual({ - subtotalMinor: 500n, - discountMinor: 0n, - taxMinor: 83n, - grossMinor: 500n, - netMinor: 417n, - refundedGrossMinor: 0n, - }); - }); - - test("includes refunded tax in the customer-facing refund amount", () => { - expect( - parsePolarOrderAmounts({ - subtotal_amount: 1_000, - discount_amount: 0, - tax_amount: 167, - total_amount: 1_000, - net_amount: 833, - refunded_amount: 417, - refunded_tax_amount: 83, - }).refundedGrossMinor, - ).toBe(500n); - }); -}); diff --git a/packages/backend/convex/model/polarOrderAmounts.ts b/packages/backend/convex/model/polarOrderAmounts.ts deleted file mode 100644 index e181e1b2..00000000 --- a/packages/backend/convex/model/polarOrderAmounts.ts +++ /dev/null @@ -1,66 +0,0 @@ -export type PolarOrderAmounts = { - subtotalMinor: bigint; - discountMinor: bigint; - taxMinor: bigint; - grossMinor: bigint; - netMinor: bigint; - refundedGrossMinor: bigint; -}; - -function integerMinor( - value: unknown, - field: string, - allowNegative = false, -): bigint { - if ( - typeof value !== "number" || - !Number.isSafeInteger(value) || - (!allowNegative && value < 0) - ) { - throw new Error(`Polar returned an invalid ${field}`); - } - return BigInt(value); -} - -/** - * Polar's `amount` and `net_amount` exclude tax. Customer-facing credit value - * is based on the gross amount the customer paid, represented by - * `total_amount`. Refunds likewise include both the net and tax portions. - */ -export function parsePolarOrderAmounts( - data: Record, -): PolarOrderAmounts { - const subtotalMinor = integerMinor( - data.subtotal_amount, - "subtotal_amount", - true, - ); - const discountMinor = integerMinor(data.discount_amount, "discount_amount"); - const taxMinor = integerMinor(data.tax_amount, "tax_amount", true); - const grossMinor = integerMinor(data.total_amount, "total_amount", true); - const netMinor = integerMinor(data.net_amount, "net_amount", true); - const refundedNetMinor = integerMinor( - data.refunded_amount, - "refunded_amount", - ); - const refundedTaxMinor = integerMinor( - data.refunded_tax_amount, - "refunded_tax_amount", - ); - - if (grossMinor !== subtotalMinor - discountMinor) { - throw new Error("Polar order total does not match subtotal minus discount"); - } - if (netMinor + taxMinor !== grossMinor) { - throw new Error("Polar order total does not match net plus tax"); - } - - return { - subtotalMinor, - discountMinor, - taxMinor, - grossMinor, - netMinor, - refundedGrossMinor: refundedNetMinor + refundedTaxMinor, - }; -} diff --git a/packages/backend/convex/schema/billing.ts b/packages/backend/convex/schema/billing.ts index b2abf9d3..29b817ba 100644 --- a/packages/backend/convex/schema/billing.ts +++ b/packages/backend/convex/schema/billing.ts @@ -256,42 +256,18 @@ export const billingTables = { .index("by_subscription_created", ["subscriptionId", "createdAt"]), billingWebhookEvents: defineTable({ - provider: v.literal("polar"), providerEnvironment, deliveryId: v.string(), eventType: v.string(), eventOccurredAt: v.number(), - providerModifiedAt: v.optional(v.number()), - resourceType: v.optional(v.string()), - resourceId: v.optional(v.string()), - organizationId: v.optional(v.string()), - providerCustomerId: v.optional(v.string()), - providerSubscriptionId: v.optional(v.string()), - providerOrderId: v.optional(v.string()), - payloadHash: v.string(), - rawPayload: v.string(), - payload: v.any(), - status: v.union( - v.literal("pending"), - v.literal("processed"), - v.literal("ignored"), - v.literal("failed"), - v.literal("deadLettered"), - ), - attemptCount: v.number(), - nextAttemptAt: v.optional(v.number()), - failureCode: v.optional(v.string()), - failureMessage: v.optional(v.string()), + resourceId: v.string(), + organizationId: v.string(), + outcome: v.union(v.literal("applied"), v.literal("ignored")), + ignoredReason: v.optional(v.string()), receivedAt: v.number(), - processedAt: v.optional(v.number()), - updatedAt: v.number(), + processedAt: v.number(), }) .index("by_environment_delivery", ["providerEnvironment", "deliveryId"]) - .index("by_resource_occurred", [ - "resourceType", - "resourceId", - "eventOccurredAt", - ]) .index("by_organization_received", ["organizationId", "receivedAt"]) .index("by_type_occurred", ["eventType", "eventOccurredAt"]), diff --git a/packages/backend/package.json b/packages/backend/package.json index fe404f7f..e4ea7941 100644 --- a/packages/backend/package.json +++ b/packages/backend/package.json @@ -24,7 +24,8 @@ }, "devDependencies": { "@baseblocks/tsconfig": "workspace:*", - "@types/bun": "1.3.14" + "@types/bun": "1.3.14", + "convex-test": "0.0.55" }, "peerDependencies": { "typescript": ">=5 <8" @@ -34,8 +35,8 @@ "@aws-sdk/s3-presigned-post": "^3.1084.0", "@aws-sdk/s3-request-presigner": "^3.1084.0", "@baseblocks/anydoc-convex": "0.1.0-alpha.15", - "@baseblocks/domain": "workspace:*", "@baseblocks/custom-blocks": "workspace:*", + "@baseblocks/domain": "workspace:*", "@baseblocks/openeditor-contracts": "workspace:*", "@convex-dev/better-auth": "^0.12.5", "@convex-dev/workflow": "^0.4.4", @@ -47,6 +48,7 @@ "@polar-sh/sdk": "0.49.0", "ai": "7.0.49", "better-auth": "^1.6.23", + "buffer": "6.0.3", "convex": "^1.42.1", "files-sdk": "2.2.3" }