diff --git a/apps/api/src/billing/autumn.ts b/apps/api/src/billing/autumn.ts index 91ffa448e3..621b316e56 100644 --- a/apps/api/src/billing/autumn.ts +++ b/apps/api/src/billing/autumn.ts @@ -4,7 +4,7 @@ import { isInvestigationPurchaseValid } from "./investigation-purchase"; import { auth } from "@databuddy/auth"; import { getRedisCache } from "@databuddy/redis"; import { getAutumn, getBillingCustomerId, getMemberRole } from "@databuddy/rpc"; -import { AutumnError } from "autumn-js"; +import { AutumnError } from "@databuddy/rpc/autumn"; import { autumnHandler } from "autumn-js/fetch"; import { useLogger } from "evlog/elysia"; import { withAutumnApiPath } from "@/lib/autumn-mount"; @@ -105,7 +105,10 @@ async function attachWithDubCustomer( await invalidateAutumnCustomerCache(request); return Response.json(result); } catch (error) { - const statusCode = error instanceof AutumnError ? error.statusCode : 500; + const statusCode = + error instanceof AutumnError && error.statusCode >= 400 + ? error.statusCode + : 503; let parsed: { message?: string; code?: string } = {}; if (error instanceof AutumnError) { try { diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index 2f5607a9ff..408e798f64 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -1,6 +1,6 @@ import "./polyfills/compression"; import { assertAuthSecretMatchesDashboard } from "@databuddy/auth"; -import { assertConfigured } from "@databuddy/env/app"; +import { assertConfigured, billingMode } from "@databuddy/env/app"; import { readBooleanEnv } from "@databuddy/env/boolean"; import { buildHttpErrorResponse } from "@databuddy/shared/http-error-response"; import cors from "@elysiajs/cors"; @@ -152,7 +152,7 @@ const app = new Elysia({ precompile: true }) .use(discovery) .use(webhooks) .mount(AUTUMN_API_PREFIX, (request) => { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { const response = buildHttpErrorResponse({ code: "NOT_FOUND", error: null, diff --git a/apps/api/src/routes/agent-business-context.test.ts b/apps/api/src/routes/agent-business-context.test.ts index ab3d4cfd9c..c46ef97b1c 100644 --- a/apps/api/src/routes/agent-business-context.test.ts +++ b/apps/api/src/routes/agent-business-context.test.ts @@ -1,4 +1,5 @@ import type { MockLanguageModelV3 } from "ai/test"; +import { BillingUnavailableError } from "@databuddy/shared/billing"; import { type OrganizationBusinessProfile, PROFILE_ORIGIN_PROVENANCE, @@ -445,7 +446,7 @@ describe("dashboard canonical business context through the native HTTP/model str }); it("delivers organization-wide context with no selected website", async () => { expect((await chat({ websiteId: undefined })).status).toBe(200); - expect(JSON.stringify(state.prompts[0].prompt)).toContain(meaning); + expect(JSON.stringify(state.prompts[0]?.prompt)).toContain(meaning); }); it("delivers team-only settings and preserves mixed legacy meanings as assertions", async () => { for (const content of ["", `${meaning}. Public capability claims.`]) { @@ -467,9 +468,9 @@ describe("dashboard canonical business context through the native HTTP/model str 200 ); expect(state.read).not.toHaveBeenCalled(); - expect(JSON.stringify(state.prompts[0].prompt)).not.toContain(meaning); + expect(JSON.stringify(state.prompts[0]?.prompt)).not.toContain(meaning); for (const assertion of Object.values(teamContext)) { - expect(JSON.stringify(state.prompts[0].prompt)).not.toContain(assertion); + expect(JSON.stringify(state.prompts[0]?.prompt)).not.toContain(assertion); } }); it("rejects an inaccessible organization, site or existing chat before reading profiles", async () => { @@ -483,7 +484,7 @@ describe("dashboard canonical business context through the native HTTP/model str it("continues the stream with explicit uncertainty when the profile read fails", async () => { state.read.mockRejectedValueOnce(new Error("synthetic read failure")); expect((await chat()).status).toBe(200); - expect(JSON.stringify(state.prompts[0].prompt)).toContain( + expect(JSON.stringify(state.prompts[0]?.prompt)).toContain( "unavailable for this turn" ); expect(state.read).toHaveBeenCalledTimes(1); @@ -567,6 +568,26 @@ describe("dashboard memory writes", () => { }); describe("ask route permissions", () => { + it("returns HTTP 503 when billing fails before streamed text", async () => { + state.stream.mockImplementationOnce(async function* () { + yield ""; + throw new BillingUnavailableError("Synthetic provider detail"); + }); + const response = await agent.handle( + new Request("http://localhost/v1/agent/ask", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ question: "Show recent events", stream: true }), + }) + ); + expect(response.status).toBe(503); + expect(response.headers.get("content-type")).toBe("application/json"); + expect(await response.json()).toMatchObject({ + success: false, + code: "BILLING_UNAVAILABLE", + }); + }); + it("runs the shared agent read-only for answers and streams", async () => { for (const stream of [false, true]) { const response = await agent.handle( diff --git a/apps/api/src/routes/agent.ts b/apps/api/src/routes/agent.ts index f927f3eb20..c2af9fe98e 100644 --- a/apps/api/src/routes/agent.ts +++ b/apps/api/src/routes/agent.ts @@ -40,6 +40,7 @@ import { } from "@databuddy/redis/stream-buffer"; import { getRedisCache } from "@databuddy/redis"; import { ratelimit } from "@databuddy/redis/rate-limit"; +import { isBillingUnavailable } from "@databuddy/shared/billing"; import { convertToModelMessages, generateId, @@ -87,6 +88,8 @@ function jsonError(status: number, code: string, message: string): Response { const INTERNAL_AGENT_ERROR_MESSAGE = "Agent request failed. Please try again shortly."; +const BILLING_UNAVAILABLE_MESSAGE = + "Billing is temporarily unavailable. Please try again shortly."; function getErrorName(error: unknown, fallback = "UnknownError"): string { if (error instanceof Error) { @@ -424,14 +427,21 @@ function createAgentUsageInjector( }); } -function createPlainTextStreamResponse( - stream: AsyncIterable -): Response { +async function createPlainTextStreamResponse( + stream: AsyncGenerator +): Promise { + let first = await stream.next(); + while (!(first.done || first.value)) { + first = await stream.next(); + } const encoder = new TextEncoder(); return new Response( new ReadableStream({ async start(controller) { try { + if (!first.done) { + controller.enqueue(encoder.encode(first.value)); + } for await (const chunk of stream) { if (chunk) { controller.enqueue(encoder.encode(chunk)); @@ -532,7 +542,7 @@ export const agent = new Elysia({ prefix: "/v1/agent" }) } : createSessionAgentActor(user, request.headers); if (body.stream) { - return createPlainTextStreamResponse( + return await createPlainTextStreamResponse( streamDatabuddyAgent({ actor, conversationId, @@ -573,6 +583,13 @@ export const agent = new Elysia({ prefix: "/v1/agent" }) error_type: getErrorName(error), source: "slack", }); + if (isBillingUnavailable(error)) { + return jsonError( + 503, + "BILLING_UNAVAILABLE", + BILLING_UNAVAILABLE_MESSAGE + ); + } return jsonError(500, "INTERNAL_ERROR", INTERNAL_AGENT_ERROR_MESSAGE); } }, @@ -1229,6 +1246,13 @@ export const agent = new Elysia({ prefix: "/v1/agent" }) ...(user?.id ? { agent_user_id: user.id } : {}), error_type: getErrorName(error), }); + if (isBillingUnavailable(error)) { + return jsonError( + 503, + "BILLING_UNAVAILABLE", + BILLING_UNAVAILABLE_MESSAGE + ); + } return jsonError(500, "INTERNAL_ERROR", INTERNAL_AGENT_ERROR_MESSAGE); } })(); diff --git a/apps/api/src/routes/webhooks/autumn.test.ts b/apps/api/src/routes/webhooks/autumn.test.ts index c5e340320d..79dd3919b4 100644 --- a/apps/api/src/routes/webhooks/autumn.test.ts +++ b/apps/api/src/routes/webhooks/autumn.test.ts @@ -213,10 +213,18 @@ vi.mock("@databuddy/email", async (importOriginal) => ({ UsageLimitEmail: vi.fn(() => ({ type: "limit" })), })); -vi.mock("@databuddy/env/app", async (importOriginal) => ({ - ...(await importOriginal()), - config: { email: { alertsFrom: "alerts@databuddy.cc" } }, -})); +vi.mock("@databuddy/env/app", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + config: { + email: { alertsFrom: "alerts@databuddy.cc" }, + get services() { + return actual.config.services; + }, + }, + }; +}); vi.mock("@databuddy/notifications", () => ({ SlackProvider: class { diff --git a/apps/api/src/routes/webhooks/autumn.ts b/apps/api/src/routes/webhooks/autumn.ts index 9d27619f93..d6f37fbee8 100644 --- a/apps/api/src/routes/webhooks/autumn.ts +++ b/apps/api/src/routes/webhooks/autumn.ts @@ -56,7 +56,7 @@ const COOLDOWN_MS = 7 * 24 * 60 * 60 * 1000; const SVIX_SECRET = readBooleanEnv("SELFHOST") ? undefined : process.env.AUTUMN_WEBHOOK_SECRET; -const SLACK_URL = process.env.SLACK_WEBHOOK_URL ?? ""; +const SLACK_URL = config.services.slackWebhookUrl ?? ""; const svix = SVIX_SECRET ? new Webhook(SVIX_SECRET) : null; const slack = SLACK_URL ? new SlackProvider({ webhookUrl: SLACK_URL }) : null; @@ -347,7 +347,7 @@ export async function sendAlertEmail(opts: { }); return { success: true, message: "No notification recipient found" }; } - const resendApiKey = process.env.RESEND_API_KEY; + const resendApiKey = config.services.resendApiKey; if (!resendApiKey) { log.error(new Error("RESEND_API_KEY is not configured"), { autumn: { diff --git a/apps/basket/src/lib/billing.test.ts b/apps/basket/src/lib/billing.test.ts index 119b70c092..b66971c49e 100644 --- a/apps/basket/src/lib/billing.test.ts +++ b/apps/basket/src/lib/billing.test.ts @@ -1,4 +1,6 @@ import { afterEach, beforeEach, describe, expect, test, vi } from "vitest"; +import { config } from "@databuddy/env/app"; +import { BillingUnavailableError } from "@databuddy/shared/billing"; import { EvlogError } from "evlog"; const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({ @@ -13,9 +15,16 @@ const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({ mockLoggerWarn: vi.fn(() => {}), })); -vi.mock("@databuddy/rpc/autumn", () => ({ - getAutumn: () => ({ check: mockCheck }), -})); +vi.mock("@databuddy/rpc/autumn", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + getAutumn: () => + config.services.autumnSecretKey + ? { check: mockCheck } + : actual.getAutumn(), + }; +}); vi.mock("evlog/elysia", () => ({ useLogger: () => ({ @@ -35,11 +44,23 @@ const { checkAutumnUsage } = await import("./billing"); describe("checkAutumnUsage", () => { beforeEach(() => { vi.stubEnv("SELFHOST", "false"); + vi.stubEnv("AUTUMN_SECRET_KEY", "am_sk_test_synthetic"); mockCheck.mockReset(); mockLoggerSet.mockReset(); mockLoggerWarn.mockReset(); }); - afterEach(() => vi.unstubAllEnvs()); + afterEach(() => { + vi.unstubAllEnvs(); + }); + + test("hosted events reject missing billing configuration", async () => { + vi.stubEnv("NODE_ENV", "production"); + vi.stubEnv("AUTUMN_SECRET_KEY", undefined); + await expect(checkAutumnUsage("cust_1", "events")).rejects.toMatchObject({ + status: 503, + message: "Billing check unavailable", + }); + }); test("self-hosted events skip hosted billing", async () => { vi.stubEnv("SELFHOST", "true"); @@ -133,6 +154,18 @@ describe("checkAutumnUsage", () => { }); }); + test("an Autumn outage accepts the event instead of dropping it", async () => { + mockCheck.mockRejectedValue( + new BillingUnavailableError("Autumn check failed") + ); + await expect(checkAutumnUsage("cust_1", "events")).resolves.toEqual({ + allowed: true, + }); + expect(mockLoggerSet).toHaveBeenCalledWith({ + billing: { allowed: true, checkFailed: true }, + }); + }); + test("logs checkFailed on API error", async () => { mockCheck.mockRejectedValue(new Error("timeout")); await expect(checkAutumnUsage("cust_1", "events")).rejects.toThrow( diff --git a/apps/basket/src/lib/billing.ts b/apps/basket/src/lib/billing.ts index 6b5affa9f8..288165bfb5 100644 --- a/apps/basket/src/lib/billing.ts +++ b/apps/basket/src/lib/billing.ts @@ -1,5 +1,9 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; -import { getAutumn } from "@databuddy/rpc/autumn"; +import { billingMode, config } from "@databuddy/env/app"; +import { + autumnCall, + getAutumn, + isBillingUnavailable, +} from "@databuddy/rpc/autumn"; import { basketErrors } from "@lib/structured-errors"; import { captureError, record } from "@lib/tracing"; import { EvlogError } from "evlog"; @@ -15,7 +19,7 @@ export function checkAutumnUsage( properties?: Record, quantity = 1 ): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { return Promise.resolve({ allowed: true }); } return record("checkAutumnUsage", async (): Promise => { @@ -23,13 +27,15 @@ export function checkAutumnUsage( try { const response = await record("autumn.check", () => - getAutumn().check({ - customerId, - featureId, - sendEvent: true, - requiredBalance: quantity, - properties, - }) + autumnCall("check", () => + getAutumn().check({ + customerId, + featureId, + sendEvent: true, + requiredBalance: quantity, + properties, + }) + ) ); const b = response.balance; @@ -60,6 +66,13 @@ export function checkAutumnUsage( if (error instanceof EvlogError) { throw error; } + if (isBillingUnavailable(error) && config.services.autumnSecretKey) { + log.set({ billing: { allowed: true, checkFailed: true } }); + captureError(error, { + message: "Autumn check failed, accepting event", + }); + return { allowed: true }; + } log.set({ billing: { allowed: false, checkFailed: true } }); captureError(error, { diff --git a/apps/dashboard/app/(main)/billing/page.tsx b/apps/dashboard/app/(main)/billing/page.tsx index d6da0120cd..ae2e515efe 100644 --- a/apps/dashboard/app/(main)/billing/page.tsx +++ b/apps/dashboard/app/(main)/billing/page.tsx @@ -42,6 +42,7 @@ import { CrownIcon, PlusIcon, TrendUpIcon, + WarningIcon, XMarkIcon as XIcon, } from "@databuddy/ui/icons"; import { @@ -62,6 +63,7 @@ const INTELLIGENCE_PLAN_ID_SET = new Set( interface OrgUsageData { balance?: number | null; includedUsage?: number | null; + unavailable?: boolean; unlimited: boolean; } @@ -289,7 +291,7 @@ export default function BillingPage() { const orgUsage = orgUsageRaw as OrgUsageData | undefined; const overageInfo = useMemo(() => { - if (!orgUsage) { + if (!orgUsage || orgUsage.unavailable === true) { return null; } const eventsFeature = usage?.features.find( @@ -598,6 +600,21 @@ export default function BillingPage() { )} + {orgUsage?.unavailable === true && ( + + + + + )} + {usageStats.length === 0 ? ( diff --git a/apps/dashboard/components/analytics/event-limit-indicator.tsx b/apps/dashboard/components/analytics/event-limit-indicator.tsx index e522faaf69..1f4c6b9b9b 100644 --- a/apps/dashboard/components/analytics/event-limit-indicator.tsx +++ b/apps/dashboard/components/analytics/event-limit-indicator.tsx @@ -9,7 +9,7 @@ import { formatLocaleNumber } from "@/lib/format-locale-number"; import { orpc } from "@/lib/orpc"; import { cn } from "@/lib/utils"; import { WarningIcon } from "@databuddy/ui/icons"; -import { buttonVariants } from "@databuddy/ui"; +import { buttonVariants, Text } from "@databuddy/ui"; export function EventLimitIndicator() { const pathname = usePathname(); @@ -20,10 +20,27 @@ export function EventLimitIndicator() { enabled: !(isDemoRoute || isSelfHosted), }); - if (isSelfHosted || !data || data.unlimited) { + if (isDemoRoute || isSelfHosted || !data || data.unlimited) { return null; } + if (data.unavailable === true) { + return ( +
+
+ ); + } + const planLimit = Number(data.includedUsage ?? 0); const overageAllowed = Boolean(data.overageAllowed); const used = Number(data.used ?? 0); diff --git a/apps/insights/src/generation-billing.integration.test.ts b/apps/insights/src/generation-billing.integration.test.ts index d5f140b7bc..a73d51cdea 100644 --- a/apps/insights/src/generation-billing.integration.test.ts +++ b/apps/insights/src/generation-billing.integration.test.ts @@ -514,7 +514,7 @@ integration("native generation fixed-unit persistence", () => { const input = await fixture(); const before = calls; await expect(generateWebsiteInsights(input)).rejects.toThrow( - "unconfirmed response" + "Autumn balances.finalize failed" ); const reservation = reservationsSince(0)[0]!; expect(holds.get(reservation.lock.lock_id)?.state).toBe("held"); @@ -551,7 +551,7 @@ integration("native generation fixed-unit persistence", () => { attemptsStarted: finalAttempt ? 2 : 1, }; await expect(processInsightsJob(job)).rejects.toThrow( - "unconfirmed response" + "Autumn balances.finalize failed" ); const reservation = reservationsSince(0)[0]!; const [pending] = await db @@ -559,7 +559,7 @@ integration("native generation fixed-unit persistence", () => { .from(insightRunItems) .where(eq(insightRunItems.id, input.itemId)); expect(pending?.status).toBe(finalAttempt ? "failed" : "queued"); - expect(pending?.errorMessage).toContain("unconfirmed response"); + expect(pending?.errorMessage).toContain("Autumn balances.finalize failed"); expect(pending?.preparedStatus).toBe("succeeded"); expect(holds.get(reservation.lock.lock_id)?.state).toBe("held"); expect(calls).toBe(1); @@ -932,7 +932,7 @@ integration("native generation fixed-unit persistence", () => { const input = await fixture(); loseFinalizeReceipt = true; await expect(generateWebsiteInsights(input)).rejects.toThrow( - "unconfirmed response" + "Autumn balances.finalize failed" ); const reservation = reservationsSince(0)[0]!; expect(holds.get(reservation.lock.lock_id)?.state).toBe("confirmed"); diff --git a/apps/insights/src/investigation-billing.integration.test.ts b/apps/insights/src/investigation-billing.integration.test.ts index b16a3ad044..1417a15780 100644 --- a/apps/insights/src/investigation-billing.integration.test.ts +++ b/apps/insights/src/investigation-billing.integration.test.ts @@ -1,4 +1,5 @@ -import { afterEach, describe, expect, it, mock } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, mock } from "bun:test"; +import { createAutumnClient } from "@databuddy/rpc/autumn"; import { INVESTIGATION_USAGE } from "@databuddy/shared/billing"; // The provider contract is exercised through the installed SDK. Customer ownership @@ -10,7 +11,6 @@ mock.module("@databuddy/ai/agents/execution", () => ({ const { assertInvestigationReservationActive, canRunInvestigation, - createInvestigationBillingClient, releaseInvestigationCharge, reserveInvestigationCharge, resolveInvestigationBilling, @@ -23,6 +23,10 @@ const customerId = "synthetic-investigation-customer"; const originalSecret = process.env.AUTUMN_SECRET_KEY; const originalNodeEnv = process.env.NODE_ENV; +beforeEach(() => { + process.env.AUTUMN_SECRET_KEY = "synthetic-local-only"; +}); + afterEach(() => { if (originalSecret === undefined) { delete process.env.AUTUMN_SECRET_KEY; @@ -128,7 +132,7 @@ function provider( { status: fault.status ?? 500 } ); }; - const client = createInvestigationBillingClient({ + const client = createAutumnClient({ secretKey: "synthetic-local-only", fetcher: async (request) => { if (!(request instanceof Request)) { @@ -294,7 +298,6 @@ integration("investigation billing through the native Autumn SDK", () => { it("rejects missing production configuration and mismatched native customer identities", async () => { delete process.env.AUTUMN_SECRET_KEY; process.env.NODE_ENV = "production"; - expect(() => createInvestigationBillingClient()).toThrow("not configured"); await expect( resolveInvestigationBilling({ organizationId: "synthetic-org" }) ).rejects.toThrow("not configured"); diff --git a/apps/insights/src/investigation-billing.ts b/apps/insights/src/investigation-billing.ts index 4833ab2a98..4e6f3efe00 100644 --- a/apps/insights/src/investigation-billing.ts +++ b/apps/insights/src/investigation-billing.ts @@ -1,67 +1,42 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { resolveAgentBillingCustomerId } from "@databuddy/ai/agents/execution"; import { createHash } from "node:crypto"; +import { + AutumnError, + autumnCall, + BillingUnavailableError, + getAutumn, +} from "@databuddy/rpc/autumn"; import { INVESTIGATION_USAGE } from "@databuddy/shared/billing"; -import { Autumn, HTTPClient } from "autumn-js"; +import type { Autumn } from "autumn-js"; import { captureInsightsError } from "./lib/evlog-insights"; -export interface InvestigationBilling { - customerId: string | null; - mode: "fixed" | "unconfigured"; -} +export type InvestigationBilling = + | { customerId: string; mode: "fixed" } + | { customerId: null; mode: "unconfigured" }; const LOCK_MS = 23 * 60 * 60 * 1000; -export function createInvestigationBillingClient( - options: { - secretKey?: string; - fetcher?: NonNullable< - ConstructorParameters[0] - >["fetcher"]; - } = {} -): Autumn { - const secretKey = options.secretKey ?? process.env.AUTUMN_SECRET_KEY; - if (!secretKey?.trim()) { - throw new Error("Investigation billing is not configured"); - } - const httpClient = new HTTPClient({ fetcher: options.fetcher }); - httpClient.addHook("response", (response) => { - // SDK 1.2.23 also accepts a degraded 202 body. It is not a receipt. - if (response.status === 202) { - throw new Error("Investigation billing returned an unconfirmed response"); - } - }); - return new Autumn({ - secretKey, - httpClient, - failOpen: false, - timeoutMs: 5000, - retryConfig: { strategy: "none" }, - }); -} - export async function resolveInvestigationBilling( principal: { organizationId: string; userId?: string | null }, client?: Autumn ): Promise { - if (readBooleanEnv("SELFHOST")) { - return { mode: "unconfigured", customerId: null }; - } - if (!(process.env.AUTUMN_SECRET_KEY?.trim() || client)) { - if (process.env.NODE_ENV === "production") { - throw new Error("Investigation billing is not configured"); - } + if (billingMode() !== "live") { return { mode: "unconfigured", customerId: null }; } const customerId = await resolveAgentBillingCustomerId(principal); if (!customerId) { - throw new Error("The investigation billing customer is unavailable"); + throw new BillingUnavailableError( + "The investigation billing customer is unavailable" + ); } - const customer = await ( - client ?? createInvestigationBillingClient() - ).customers.get({ customerId }); + const customer = await autumnCall("customers.get", () => + (client ?? getAutumn()).customers.get({ customerId }) + ); if (customer.id !== customerId) { - throw new Error("The investigation billing customer could not be verified"); + throw new BillingUnavailableError( + "The investigation billing customer could not be verified" + ); } return { customerId, mode: "fixed" }; } @@ -73,16 +48,18 @@ export async function canRunInvestigation( if (billing.mode === "unconfigured") { return true; } - if (!billing.customerId) { - throw new Error("The investigation billing customer is unavailable"); - } - const result = await (client ?? createInvestigationBillingClient()).check({ - customerId: billing.customerId, - featureId: INVESTIGATION_USAGE.featureId, - requiredBalance: 1, - }); - if (result.customerId !== billing.customerId) { - throw new Error("Investigation access could not be verified"); + const { customerId } = billing; + const result = await autumnCall("check", () => + (client ?? getAutumn()).check({ + customerId, + featureId: INVESTIGATION_USAGE.featureId, + requiredBalance: 1, + }) + ); + if (result.customerId !== customerId) { + throw new BillingUnavailableError( + "Investigation access could not be verified" + ); } return result.allowed === true; } @@ -93,10 +70,10 @@ interface InvestigationOperation { websiteId: string; } -interface InvestigationReservation extends InvestigationBilling { +type InvestigationReservation = InvestigationBilling & { expiresAt: Date; id: string; -} +}; function reservationId(input: InvestigationOperation): string { return `investigation:${createHash("sha256") @@ -127,29 +104,31 @@ export async function reserveInvestigationCharge( return reservation; } assertInvestigationReservationActive(reservation); - if (!reservation.customerId) { - throw new Error("The investigation billing customer is unavailable"); - } - const autumn = client ?? createInvestigationBillingClient(); + const { customerId } = reservation; + const autumn = client ?? getAutumn(); // Autumn owns the hold. A duplicate/ambiguous response never authorizes work. // The immutable expiry is shorter than the provider's idempotency window: // a released/confirmed operation cannot be reserved again after that window. - const result = await autumn.check( - { - customerId: reservation.customerId, - featureId: INVESTIGATION_USAGE.featureId, - requiredBalance: 1, - sendEvent: true, - lock: { - enabled: true, - lockId: reservation.id, - expiresAt: reservation.expiresAt.getTime(), + const result = await autumnCall("check", () => + autumn.check( + { + customerId, + featureId: INVESTIGATION_USAGE.featureId, + requiredBalance: 1, + sendEvent: true, + lock: { + enabled: true, + lockId: reservation.id, + expiresAt: reservation.expiresAt.getTime(), + }, }, - }, - { headers: { "Idempotency-Key": `${reservation.id}:reserve` } } + { headers: { "Idempotency-Key": `${reservation.id}:reserve` } } + ) ); - if (result.customerId !== reservation.customerId) { - throw new Error("Investigation reservation could not be verified"); + if (result.customerId !== customerId) { + throw new BillingUnavailableError( + "Investigation reservation could not be verified" + ); } if (!result.allowed) { throw new Error( @@ -210,67 +189,62 @@ export function assertInvestigationReservationActive( } } +function isMissingLock(error: unknown, id: string): boolean { + if (!(error instanceof AutumnError && error.statusCode === 400)) { + return false; + } + let body: unknown; + try { + body = JSON.parse(error.body); + } catch { + return false; + } + return ( + typeof body === "object" && + body !== null && + "code" in body && + "message" in body && + body.code === "invalid_request" && + body.message === `Lock not found for ID: ${id}` + ); +} + async function finalizeReservation( id: string, complete: boolean, client?: Autumn ): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { return; } - if (!(client || process.env.AUTUMN_SECRET_KEY?.trim())) { - if (process.env.NODE_ENV !== "production") { - return; - } - throw new Error("Investigation billing is not configured"); - } - try { - // Full confirmation does not debit again. Replays only finalize this lock; - // they must never reserve or track a replacement unit for a saved result. - const result = await ( - client ?? createInvestigationBillingClient() - ).balances.finalize({ - lockId: id, - action: complete ? "confirm" : "release", - }); - if (!result.success) { - throw new Error("Investigation settlement was not confirmed"); - } - } catch (error) { - if ( - error instanceof Error && - "statusCode" in error && - error.statusCode === 400 && - "body" in error && - typeof error.body === "string" - ) { - let body: unknown; - try { - body = JSON.parse(error.body); - } catch { - /* Non-JSON provider failures remain failures. */ - } - if ( - body && - typeof body === "object" && - "code" in body && - "message" in body && - body.code === "invalid_request" && - body.message === `Lock not found for ID: ${id}` - ) { - // Missing can mean confirmed, released, or expired. Nothing remains to - // finalize; it is not proof of payment and never warrants a new debit. - if (complete) { - captureInsightsError( - new Error("Investigation charge lock was gone before confirmation"), - "investigation_billing.settlement_unconfirmed", - { lock_id: id } - ); + // Full confirmation does not debit again. Replays only finalize this lock; + // they must never reserve or track a replacement unit for a saved result. + const result = await autumnCall("balances.finalize", () => + (client ?? getAutumn()).balances + .finalize({ lockId: id, action: complete ? "confirm" : "release" }) + .catch((error: unknown) => { + if (isMissingLock(error, id)) { + return null; } - return; - } + throw error; + }) + ); + if (!result) { + // Missing can mean confirmed, released, or expired. Nothing remains to + // finalize; it is not proof of payment and never warrants a new debit. + if (complete) { + captureInsightsError( + new Error("Investigation charge lock was gone before confirmation"), + "investigation_billing.settlement_unconfirmed", + { lock_id: id } + ); } - throw error; + return; + } + if (!result.success) { + throw new BillingUnavailableError( + "Investigation settlement was not confirmed" + ); } } diff --git a/apps/insights/src/resume-clarification.integration.test.ts b/apps/insights/src/resume-clarification.integration.test.ts index ab880cc7b5..8f669e78b9 100644 --- a/apps/insights/src/resume-clarification.integration.test.ts +++ b/apps/insights/src/resume-clarification.integration.test.ts @@ -27,9 +27,11 @@ import { rankInvestigationBusinessContext } from "./business-context-ranking"; import { createEvidenceSnapshot } from "./evidence-snapshot"; import { resumeInsightReply, recordInsightReplyFailure } from "./resume"; import * as billing from "./investigation-billing"; +import { createAutumnClient } from "@databuddy/rpc/autumn"; const integration = process.env.INSIGHTS_INTEGRATION_TESTS === "true" ? describe : describe.skip; +const originalSecret = process.env.AUTUMN_SECRET_KEY; const ids: string[] = []; const signal: InvestigationSignal = { signalKey: "goal:workspace", @@ -185,7 +187,7 @@ function nativeProvider( released = 0; let loseConfirmation = options.loseConfirmation === true; let failConfirmation = options.failConfirmation === true; - const client = billing.createInvestigationBillingClient({ + const client = createAutumnClient({ secretKey: "synthetic-local-only", fetcher: async (request) => { if (!(request instanceof Request)) { @@ -285,6 +287,7 @@ function wireProvider(remote: ReturnType) { const nativeReserve = billing.reserveInvestigationCharge; const nativeSettle = billing.settleInvestigationCharge; const nativeRelease = billing.releaseInvestigationCharge; + process.env.AUTUMN_SECRET_KEY = "synthetic-local-only"; spyOn(billing, "resolveInvestigationBilling").mockResolvedValue({ mode: "fixed", customerId: "synthetic-customer", @@ -303,7 +306,14 @@ function wireProvider(remote: ReturnType) { } integration("included saved-evidence replies", () => { - afterEach(() => mock.restore()); + afterEach(() => { + mock.restore(); + if (originalSecret === undefined) { + Reflect.deleteProperty(process.env, "AUTUMN_SECRET_KEY"); + } else { + process.env.AUTUMN_SECRET_KEY = originalSecret; + } + }); afterAll(async () => { if (ids.length) { await db.delete(organization).where(inArray(organization.id, ids)); diff --git a/packages/ai/src/ai/agents/execution.test.ts b/packages/ai/src/ai/agents/execution.test.ts index 554d03c9e0..54a93fc0fc 100644 --- a/packages/ai/src/ai/agents/execution.test.ts +++ b/packages/ai/src/ai/agents/execution.test.ts @@ -2,6 +2,7 @@ import { afterAll, beforeEach, describe, expect, it, mock } from "bun:test"; import { createLogger } from "evlog"; const originalAutumnSecretKey = process.env.AUTUMN_SECRET_KEY; +const actualAutumn = { ...(await import("@databuddy/rpc/autumn")) }; const mockAutumnCheck = mock(async (input: { customerId: string }) => ({ allowed: true, @@ -25,6 +26,7 @@ const mockGetOrganizationOwnerId = mock(async (organizationId: string) => const mockMergeWideEvent = mock((_: Record) => {}); mock.module("@databuddy/rpc/autumn", () => ({ + ...actualAutumn, getAutumn: () => ({ customers: { get: async (input: { customerId: string }) => ({ @@ -65,7 +67,6 @@ mock.module("../../lib/tracing", () => ({ const { getAgentBillingAccess, - isAgentBillingConfigured, resolveAgentBillingCustomerId, trackAgentUsage, trackAgentUsageAndBill, @@ -163,7 +164,6 @@ describe("resolveAgentBillingCustomerId", () => { userId: "self-hosted-user", }); - expect(isAgentBillingConfigured()).toBe(false); expect(customerId).toBeNull(); expect(mockGetBillingCustomerId).not.toHaveBeenCalled(); expect(mockGetOrganizationOwnerId).not.toHaveBeenCalled(); @@ -297,7 +297,6 @@ it("self-hosted AI keeps provider setup and skips all hosted billing", async () AI_GATEWAY_API_KEY: "synthetic-ai-key", }; try { - expect(isAgentBillingConfigured()).toBe(false); expect( await resolveAgentBillingCustomerId({ organizationId: "synthetic-org", diff --git a/packages/ai/src/ai/agents/execution.ts b/packages/ai/src/ai/agents/execution.ts index 0d94c8fa18..1b82bd0c25 100644 --- a/packages/ai/src/ai/agents/execution.ts +++ b/packages/ai/src/ai/agents/execution.ts @@ -1,7 +1,11 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode, readBooleanEnv } from "@databuddy/env/app"; import type { ApiKeyRow } from "@databuddy/api-keys/resolve"; import { MIN_AGENT_CREDIT_CHECK_BALANCE } from "@databuddy/shared/agent-credits"; -import { getAutumn } from "@databuddy/rpc/autumn"; +import { + autumnCall, + BillingUnavailableError, + getAutumn, +} from "@databuddy/rpc/autumn"; import { getBillingCustomerId } from "@databuddy/rpc/billing"; import { getOrganizationOwnerId } from "@databuddy/rpc/organization"; import type { RequestLogger } from "evlog"; @@ -33,19 +37,12 @@ export interface AgentBillingAccess { customerId: string | null; } -export function isAgentBillingConfigured(): boolean { - return ( - !readBooleanEnv("SELFHOST") && - Boolean(process.env.AUTUMN_SECRET_KEY?.trim()) - ); -} - export async function resolveAgentBillingCustomerId(principal: { apiKey?: ApiKeyRow | null; organizationId?: string | null; userId?: string | null; }): Promise { - if (!isAgentBillingConfigured()) { + if (billingMode() !== "live") { mergeAgentBillingFields({ billingCustomerId: null, organizationId: @@ -102,7 +99,7 @@ export async function getAgentBillingAccess( "Ask your administrator to configure AI before using Databunny." ); } - if (!isAgentBillingConfigured()) { + if (billingMode() !== "live") { mergeWideEvent({ agent_credits_allowed: true, agent_credits_check_skipped: true, @@ -110,28 +107,36 @@ export async function getAgentBillingAccess( return { allowed: true, customerId: billingCustomerId }; } if (!billingCustomerId) { - throw new Error("The agent billing customer is unavailable"); + throw new BillingUnavailableError( + "The agent billing customer is unavailable" + ); } const startedAt = performance.now(); try { - const autumn = getAutumn({ strict: true }); - const customer = await autumn.customers.get({ - customerId: billingCustomerId, - }); + const autumn = getAutumn(); + const customer = await autumnCall("customers.get", () => + autumn.customers.get({ customerId: billingCustomerId }) + ); if (customer.id !== billingCustomerId) { - throw new Error("The agent billing customer could not be verified"); + throw new BillingUnavailableError( + "The agent billing customer could not be verified" + ); } - const result = await autumn.check({ - customerId: billingCustomerId, - featureId: "agent_credits", - requiredBalance: MIN_AGENT_CREDIT_CHECK_BALANCE, - }); + const result = await autumnCall("check", () => + autumn.check({ + customerId: billingCustomerId, + featureId: "agent_credits", + requiredBalance: MIN_AGENT_CREDIT_CHECK_BALANCE, + }) + ); if ( result.customerId !== billingCustomerId || (result.allowed && result.balance?.featureId !== "agent_credits") ) { - throw new Error("The agent credit balance could not be verified"); + throw new BillingUnavailableError( + "The agent credit balance could not be verified" + ); } const allowed = result.allowed === true; const balance = result.balance; @@ -209,7 +214,7 @@ export async function trackAgentUsageAndBill( ): Promise { const summary = trackAgentUsage(input); - if (!(isAgentBillingConfigured() && input.billingCustomerId)) { + if (!(billingMode() === "live" && input.billingCustomerId)) { return summary; } if (input.source !== "insights") { diff --git a/packages/ai/src/ai/config/enrich-context.test.ts b/packages/ai/src/ai/config/enrich-context.test.ts index 8a2af6c323..91ffe9000c 100644 --- a/packages/ai/src/ai/config/enrich-context.test.ts +++ b/packages/ai/src/ai/config/enrich-context.test.ts @@ -18,8 +18,12 @@ mock.module("../../lib/tracing", () => ({ captureError })); const { enrichAgentContext } = await import("./enrich-context"); test("self-hosted agent context keeps local entity counts without hosted plan lookup", async () => { - const original = process.env.SELFHOST; - process.env.SELFHOST = "true"; + const original = process.env; + process.env = { + ...original, + SELFHOST: "true", + AUTUMN_SECRET_KEY: "synthetic-stale-key", + }; try { const input = { userId: "synthetic-user", @@ -36,10 +40,6 @@ test("self-hosted agent context keeps local entity counts without hosted plan lo expect(await enrichAgentContext(input)).toContain("free"); expect(billingOwner).toHaveBeenCalledTimes(1); } finally { - if (original === undefined) { - Reflect.deleteProperty(process.env, "SELFHOST"); - } else { - process.env.SELFHOST = original; - } + process.env = original; } }); diff --git a/packages/ai/src/ai/config/enrich-context.ts b/packages/ai/src/ai/config/enrich-context.ts index 2a9f7b52f2..6dec72ee1f 100644 --- a/packages/ai/src/ai/config/enrich-context.ts +++ b/packages/ai/src/ai/config/enrich-context.ts @@ -1,5 +1,5 @@ import { and, count, db, eq, isNull } from "@databuddy/db"; -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { annotations, funnelDefinitions, @@ -17,7 +17,7 @@ async function fetchPlanContext( userId: string, organizationId: string | null ): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { return "Analytics features have no plan limits."; } try { diff --git a/packages/ai/src/query/index.ts b/packages/ai/src/query/index.ts index 08e97ebffd..64f14330ba 100644 --- a/packages/ai/src/query/index.ts +++ b/packages/ai/src/query/index.ts @@ -1,5 +1,5 @@ /** biome-ignore-all lint/performance/noBarrelFile: this is a barrel file */ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { getBillingOwner } from "@databuddy/rpc/billing"; import { getOrganizationOwnerId } from "@databuddy/rpc/organization"; import { @@ -133,7 +133,7 @@ export async function queryPlanGateError( queryTypes: string[], scope: { organizationId: string | null } | { websiteId: string } ): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { return null; } const required = new Set(); diff --git a/packages/auth/src/auth.ts b/packages/auth/src/auth.ts index a1c5ee656b..05136ca895 100644 --- a/packages/auth/src/auth.ts +++ b/packages/auth/src/auth.ts @@ -30,8 +30,7 @@ import { ResetPasswordEmail, VerificationEmail, } from "@databuddy/email"; -import { config } from "@databuddy/env/app"; -import { readBooleanEnv } from "@databuddy/env/app"; +import { billingMode, config, readBooleanEnv } from "@databuddy/env/app"; import { SlackProvider } from "@databuddy/notifications"; import { getRedisCache, @@ -179,7 +178,7 @@ function shouldRequireEmailVerification() { if ( required && isSelfHosted() && - !(process.env.RESEND_API_KEY?.trim() && process.env.EMAIL_FROM?.trim()) + !(config.services.resendApiKey && process.env.EMAIL_FROM?.trim()) ) { throw new Error( "Self-hosted email verification requires RESEND_API_KEY and EMAIL_FROM on a verified domain." @@ -235,7 +234,7 @@ async function sendAuthEmail(input: { log.info({ service: "auth", auth_email_skipped: true }); return; } - const apiKey = process.env.RESEND_API_KEY; + const apiKey = config.services.resendApiKey; if (!apiKey) { log.error({ service: "auth", @@ -278,7 +277,7 @@ function formatInvitationRole(role: string | string[]): string { .join(", "); } -const SLACK_WEBHOOK_URL = process.env.SLACK_WEBHOOK_URL ?? ""; +const SLACK_WEBHOOK_URL = config.services.slackWebhookUrl ?? ""; function notifySlack( title: string, @@ -305,7 +304,7 @@ function notifySlack( }); } -const DUB_API_KEY = process.env.DUB_API_KEY ?? ""; +const DUB_API_KEY = config.services.dubApiKey ?? ""; function trackDubSignUp(user: { id: string; @@ -528,22 +527,13 @@ async function assertUserRowDeletable( } async function assertNoRenewingSubscription(userId: string): Promise { - const secretKey = process.env.AUTUMN_SECRET_KEY?.trim(); - if (isSelfHosted()) { - return; - } - if (!secretKey) { - if (isProduction()) { - log.error({ - service: "auth", - component: "account_deletion", - message: - "AUTUMN_SECRET_KEY is not set, so the subscription check was skipped", - }); - } + if (billingMode() !== "live") { return; } - const customer = await new Autumn({ secretKey, timeoutMs: 5000 }).customers + const customer = await new Autumn({ + secretKey: config.services.autumnSecretKey, + timeoutMs: 5000, + }).customers .get({ customerId: userId }) .catch((error: unknown) => { if (error instanceof AutumnError && error.statusCode === 404) { diff --git a/packages/rpc/src/lib/autumn-client.ts b/packages/rpc/src/lib/autumn-client.ts index e058807070..42e2e59b50 100644 --- a/packages/rpc/src/lib/autumn-client.ts +++ b/packages/rpc/src/lib/autumn-client.ts @@ -1,57 +1,59 @@ -/** - * Server-side Autumn SDK — see: - * https://docs.useautumn.com/documentation/getting-started/setup - * https://docs.useautumn.com/documentation/modelling-pricing/spend-limits - */ -import { Autumn, HTTPClient } from "autumn-js"; -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode, config } from "@databuddy/env/app"; +import { BillingUnavailableError } from "@databuddy/shared/billing"; +import { Autumn, AutumnError, HTTPClient, HTTPClientError } from "autumn-js"; -function createClient(strict = false): Autumn { - const secretKey = process.env.AUTUMN_SECRET_KEY; - if (!secretKey) { - throw new Error("AUTUMN_SECRET_KEY is not set"); - } - if (!strict) { - return new Autumn({ secretKey, timeoutMs: 3000 }); - } - const httpClient = new HTTPClient(); +export { + BillingUnavailableError, + isBillingUnavailable, +} from "@databuddy/shared/billing"; +export { AutumnError } from "autumn-js"; + +type Fetcher = NonNullable< + ConstructorParameters[0] +>["fetcher"]; + +export function createAutumnClient(options: { + secretKey: string; + fetcher?: Fetcher; +}): Autumn { + const httpClient = new HTTPClient({ fetcher: options.fetcher }); httpClient.addHook("response", (response) => { if (response.status === 202) { throw new Error("Autumn returned an unconfirmed billing response"); } }); return new Autumn({ - secretKey, - timeoutMs: 3000, + secretKey: options.secretKey, httpClient, failOpen: false, + timeoutMs: 5000, retryConfig: { strategy: "none" }, }); } -export function hasHostedBilling(): boolean { - if (readBooleanEnv("SELFHOST")) { - return false; - } - return Boolean( - process.env.AUTUMN_SECRET_KEY?.trim() || - process.env.NODE_ENV === "production" - ); -} - let instance: Autumn | null = null; -let strictInstance: Autumn | null = null; -export function getAutumn(options?: { strict?: boolean }): Autumn { - if (readBooleanEnv("SELFHOST")) { - throw new Error("Hosted billing is disabled for self-hosted instances"); - } - if (options?.strict) { - strictInstance ??= createClient(true); - return strictInstance; - } - if (!instance) { - instance = createClient(); +export function getAutumn(): Autumn { + const secretKey = config.services.autumnSecretKey; + if (billingMode() !== "live" || !secretKey) { + throw new BillingUnavailableError("Autumn billing is not configured"); } + instance ??= createAutumnClient({ secretKey }); return instance; } + +export async function autumnCall( + operation: string, + call: () => Promise +): Promise { + try { + return await call(); + } catch (error) { + if (error instanceof AutumnError || error instanceof HTTPClientError) { + throw new BillingUnavailableError(`Autumn ${operation} failed`, { + cause: error, + }); + } + throw error; + } +} diff --git a/packages/rpc/src/lib/business-context-access.ts b/packages/rpc/src/lib/business-context-access.ts index 7f2241e8d4..c5ffa2aa65 100644 --- a/packages/rpc/src/lib/business-context-access.ts +++ b/packages/rpc/src/lib/business-context-access.ts @@ -1,9 +1,13 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { roleHasPermission } from "@databuddy/auth/permissions"; import { MIN_AGENT_CREDIT_CHECK_BALANCE } from "@databuddy/shared/agent-credits"; import { z } from "zod"; import { getOrganizationOwnerId } from "../utils/organization"; -import { getAutumn } from "./autumn-client"; +import { + autumnCall, + BillingUnavailableError, + getAutumn, +} from "./autumn-client"; import { logger } from "./logger"; export const businessContextGenerationAccessSchema = z.object({ @@ -36,9 +40,7 @@ export async function businessContextGenerationAccess( !( process.env.AI_GATEWAY_API_KEY?.trim() && process.env.CONTEXT_DEV_API_KEY?.trim() - ) || - (!(readBooleanEnv("SELFHOST") || process.env.AUTUMN_SECRET_KEY?.trim()) && - process.env.NODE_ENV === "production") + ) ) { return { status: "not-configured", @@ -47,8 +49,7 @@ export async function businessContextGenerationAccess( action: "contact-admin", }; } - // Match generation's local/self-hosted policy when billing is not configured. - if (readBooleanEnv("SELFHOST") || !process.env.AUTUMN_SECRET_KEY?.trim()) { + if (billingMode() !== "live") { return { status: "allowed", message: "You can generate a draft from your website.", @@ -65,15 +66,19 @@ export async function businessContextGenerationAccess( try { const customerId = await getOrganizationOwnerId(organizationId); if (!customerId) { - throw new Error("The organization billing owner is unavailable"); + throw new BillingUnavailableError( + "The organization billing owner is unavailable" + ); } - const access = await getAutumn({ strict: true }).check({ - customerId, - featureId: "agent_credits", - requiredBalance: MIN_AGENT_CREDIT_CHECK_BALANCE, - }); + const access = await autumnCall("check", () => + getAutumn().check({ + customerId, + featureId: "agent_credits", + requiredBalance: MIN_AGENT_CREDIT_CHECK_BALANCE, + }) + ); if (access.customerId !== customerId) { - throw new Error( + throw new BillingUnavailableError( "The organization generation access could not be verified" ); } diff --git a/packages/rpc/src/orpc.ts b/packages/rpc/src/orpc.ts index 815697de61..4921b7919a 100644 --- a/packages/rpc/src/orpc.ts +++ b/packages/rpc/src/orpc.ts @@ -6,7 +6,8 @@ import { } from "@databuddy/api-keys/resolve"; import { auth, type User } from "@databuddy/auth"; import { db } from "@databuddy/db"; -import { os as createOS } from "@orpc/server"; +import { billingMode } from "@databuddy/env/app"; +import { ORPCError, os as createOS } from "@orpc/server"; import { baseErrors } from "./errors"; import { enrichRpcWideEventContext, @@ -15,7 +16,7 @@ import { setRpcProcedureType, setRpcAuthTiming, } from "./lib/rpc-log-context"; -import { hasHostedBilling } from "./lib/autumn-client"; +import { isBillingUnavailable } from "./lib/autumn-client"; import { runTracked } from "./middleware/track-mutation"; import { runAuditedMutation } from "./middleware/audit-mutation"; import { type BillingOwner, getBillingOwner } from "./utils/billing"; @@ -123,7 +124,7 @@ export const createRPCContext = async ( const getBilling = async ( billingOrganizationId: string | null = organizationId ): Promise => { - if (!hasHostedBilling()) { + if (billingMode() !== "live") { return; } if (user && billingOrganizationId !== organizationId) { @@ -164,25 +165,43 @@ export type Context = Awaited>; const os = createOS.$context().errors(baseErrors); -export const publicProcedure = os.use(({ context, next, path }) => { +const procedure = os.use(async ({ next }) => { + try { + return await next(); + } catch (error) { + if (isBillingUnavailable(error)) { + throw new ORPCError("SERVICE_UNAVAILABLE", { + status: 503, + message: "Billing is temporarily unavailable", + data: { retryAfter: 30 }, + cause: error, + }); + } + throw error; + } +}); + +export const publicProcedure = procedure.use(({ context, next, path }) => { setRpcProcedureType("public"); setRpcProcedurePath(path); enrichRpcWideEventContext(context); return next(); }); -export const protectedProcedure = os.use(({ context, next, errors, path }) => { - setRpcProcedureType("protected"); - setRpcProcedurePath(path); - enrichRpcWideEventContext(context); +export const protectedProcedure = procedure.use( + ({ context, next, errors, path }) => { + setRpcProcedureType("protected"); + setRpcProcedurePath(path); + enrichRpcWideEventContext(context); - if (!(context.user || context.apiKey)) { - recordORPCError({ code: "UNAUTHORIZED" }); - throw errors.UNAUTHORIZED(); - } + if (!(context.user || context.apiKey)) { + recordORPCError({ code: "UNAUTHORIZED" }); + throw errors.UNAUTHORIZED(); + } - return next({ context }); -}); + return next({ context }); + } +); export const sessionProcedure = protectedProcedure.use( ({ context, next, errors }) => { diff --git a/packages/rpc/src/procedures/with-workspace.ts b/packages/rpc/src/procedures/with-workspace.ts index 87fb9e51ae..20efc567db 100644 --- a/packages/rpc/src/procedures/with-workspace.ts +++ b/packages/rpc/src/procedures/with-workspace.ts @@ -7,12 +7,12 @@ import { roleHasPermission, } from "@databuddy/auth/permissions"; import { db } from "@databuddy/db"; +import { billingMode } from "@databuddy/env/app"; import { cacheNamespaces, cacheable } from "@databuddy/redis"; import { normalizePlanId, type PlanId } from "@databuddy/shared/types/features"; import { ORPCError } from "@orpc/server"; import { z } from "zod"; import { rpcError } from "../errors"; -import { hasHostedBilling } from "../lib/autumn-client"; import { type Context, os } from "../orpc"; import { getMemberRole, getOrganizationOwnerId } from "../utils/organization"; @@ -156,7 +156,7 @@ async function getPlanId( } function requirePlan(plan: PlanId, requiredPlans: PlanId[] | undefined): void { - if (!(hasHostedBilling() && requiredPlans?.length)) { + if (!(billingMode() === "live" && requiredPlans?.length)) { return; } if (!requiredPlans.includes(plan)) { diff --git a/packages/rpc/src/routers/billing.boundary.test.ts b/packages/rpc/src/routers/billing.boundary.test.ts new file mode 100644 index 0000000000..ecbb88e38e --- /dev/null +++ b/packages/rpc/src/routers/billing.boundary.test.ts @@ -0,0 +1,207 @@ +import { afterAll, beforeEach, expect, mock, spyOn, test } from "bun:test"; +import { createProcedureClient } from "@orpc/server"; +import { RPCHandler } from "@orpc/server/fetch"; +import type { Context } from "../orpc"; +import { BillingUnavailableError } from "@databuddy/shared/billing"; + +mock.module("@databuddy/auth", () => ({ + auth: { api: { getSession: async () => null } }, +})); +mock.module("@databuddy/api-keys/resolve", () => ({ + getApiKeyFromHeader: async () => null, +})); +const database = { ...(await import("@databuddy/db")), db: {} }; +mock.module("@databuddy/db", () => database); +mock.module("@databuddy/services/audit", () => ({ + appendAuditEvent: async () => undefined, + appendAuditEventInTransaction: async () => undefined, +})); +mock.module("../utils/organization", () => ({ + getOrganizationOwnerId: async () => "owner-example", +})); +mock.module("../utils/billing", () => ({ + getBillingOwner: async () => ({ + customerId: "owner-example", + canUserUpgrade: true, + }), +})); +mock.module("../procedures/with-workspace", () => ({ + withWorkspace: async () => ({ organizationId: "org-example", role: "admin" }), + withPublicWorkspace: async () => ({ + organizationId: "org-example", + role: "admin", + }), +})); +mock.module("../lib/logger", () => ({ + logger: { + error: () => undefined, + info: () => undefined, + warn: () => undefined, + }, +})); +mock.module("@databuddy/db/clickhouse", () => ({ + EXCLUDE_IMPORTED_ROWS: "NOT startsWith(anonymous_id, 'imp_')", + chQuery: () => { + throw new Error("Unexpected analytics query"); + }, +})); + +let status = 202; +const requests: string[] = []; +const originalSecret = process.env.AUTUMN_SECRET_KEY; +const fetcher = Object.assign( + async ( + input: Parameters[0], + init?: Parameters[1] + ) => { + const request = input instanceof Request ? input : new Request(input, init); + const url = new URL(request.url); + requests.push(url.pathname); + expect(url.origin).toBe("https://api.useautumn.com"); + expect(["/v1/customers.get_or_create", "/v1/balances.check"]).toContain( + url.pathname + ); + expect(await request.json()).toMatchObject({ + customer_id: "owner-example", + }); + if (url.pathname === "/v1/balances.check" && status === 200) { + return Response.json({ + allowed: true, + customer_id: "owner-example", + flag: null, + balance: { + feature_id: "events", + usage: 25, + granted: 100, + remaining: 75, + unlimited: false, + overage_allowed: false, + max_purchase: null, + next_reset_at: null, + }, + }); + } + return Response.json({}, { status }); + }, + { preconnect: globalThis.fetch.preconnect } +); +const transport = spyOn(globalThis, "fetch").mockImplementation(fetcher); + +const { billingRouter } = await import("./billing"); +const { organizationsRouter } = await import("./organizations"); +const context = { + user: { id: "admin", email: "admin@example.com", name: "Admin" }, + session: { activeOrganizationId: "org-example" }, + organizationId: "org-example", + headers: new Headers(), + db: {}, +} as Context; +const input = { + featureId: "investigation_runs" as const, + enabled: true, + overageLimit: 50, +}; +const setLimit = createProcedureClient(billingRouter.setSpendLimit, { + path: ["billing", "setSpendLimit"], + context, +}); +const handler = new RPCHandler({ billing: billingRouter }); + +beforeEach(() => { + process.env.AUTUMN_SECRET_KEY = "synthetic-native-transport-only"; + requests.length = 0; +}); + +afterAll(() => { + transport.mockRestore(); + if (originalSecret === undefined) { + Reflect.deleteProperty(process.env, "AUTUMN_SECRET_KEY"); + } else { + process.env.AUTUMN_SECRET_KEY = originalSecret; + } +}); + +test.each([ + 202, 500, +])("billing outages (%s) cross the real session middleware as HTTP 503", async (responseStatus) => { + status = responseStatus; + await expect(setLimit(input)).rejects.toMatchObject({ + code: "SERVICE_UNAVAILABLE", + status: 503, + }); + const result = await handler.handle( + new Request("http://localhost/rpc/billing/setSpendLimit", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ json: input }), + }), + { prefix: "/rpc", context } + ); + expect(result.matched).toBe(true); + expect(result.response?.status).toBe(503); + expect(requests).toEqual([ + "/v1/customers.get_or_create", + "/v1/customers.get_or_create", + ]); +}); + +test("usage shows owner outages and recovers without inventing permissions", async () => { + const owner: NonNullable>> = { + customerId: "owner-example", + canUserUpgrade: false, + isOrganization: true, + planId: "free", + }; + const lookup = mock(async () => owner); + lookup.mockRejectedValueOnce( + new BillingUnavailableError("SYNTHETIC_OWNER_DOWN") + ); + const usageContext = { ...context, getBilling: lookup }; + const usage = createProcedureClient(organizationsRouter.getUsage, { + context: usageContext, + }); + const usageHandler = new RPCHandler({ organizations: organizationsRouter }); + const result = await usageHandler.handle( + new Request("http://localhost/rpc/organizations/getUsage", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ json: {} }), + }), + { prefix: "/rpc", context: usageContext } + ); + expect(result.response?.status).toBe(200); + expect(await result.response?.json()).toEqual({ + json: { unavailable: true }, + }); + expect(requests).toEqual([]); + + status = 500; + await expect(usage()).resolves.toEqual({ + unavailable: true, + isOrganizationUsage: true, + canUserUpgrade: false, + }); + status = 200; + const recovered = await usage(); + expect(recovered).toMatchObject({ + used: 25, + limit: 100, + remaining: 75, + canUserUpgrade: false, + isOrganizationUsage: true, + }); + expect(recovered).not.toHaveProperty("unavailable"); +}); + +test("usage preserves unexpected owner errors", async () => { + const error = new Error("SYNTHETIC_UNEXPECTED_OWNER_FAILURE"); + const usage = createProcedureClient(organizationsRouter.getUsage, { + context: { + ...context, + getBilling: async () => { + throw error; + }, + }, + }); + await expect(usage()).rejects.toBe(error); +}); diff --git a/packages/rpc/src/routers/billing.test.ts b/packages/rpc/src/routers/billing.test.ts index 22f7747446..250a465fb3 100644 --- a/packages/rpc/src/routers/billing.test.ts +++ b/packages/rpc/src/routers/billing.test.ts @@ -212,7 +212,7 @@ describe("native investigation spending limits", () => { enabled: true, overageLimit: 50, }) - ).rejects.toMatchObject({ code: "INTERNAL_SERVER_ERROR" }); + ).rejects.toMatchObject({ code: "billing_unavailable" }); expect(reads).toBe(1); expect(writes).toEqual([]); expect(limits).toEqual(originalLimits); @@ -226,7 +226,7 @@ describe("native investigation spending limits", () => { enabled: true, overageLimit: 50, }) - ).rejects.toMatchObject({ code: "INTERNAL_SERVER_ERROR" }); + ).rejects.toMatchObject({ code: "billing_unavailable" }); expect(reads).toBe(1); expect(writes).toEqual([]); }); @@ -242,7 +242,7 @@ describe("native investigation spending limits", () => { enabled: true, overageLimit: 50, }) - ).rejects.toMatchObject({ code: "INTERNAL_SERVER_ERROR" }); + ).rejects.toMatchObject({ code: "billing_unavailable" }); expect(reads).toBe(1); expect(writes).toHaveLength(1); expect(limits).toEqual(originalLimits); diff --git a/packages/rpc/src/routers/billing.ts b/packages/rpc/src/routers/billing.ts index d7cea708c7..4b13e2a05f 100644 --- a/packages/rpc/src/routers/billing.ts +++ b/packages/rpc/src/routers/billing.ts @@ -1,8 +1,12 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { chQuery, EXCLUDE_IMPORTED_ROWS } from "@databuddy/db/clickhouse"; import { z } from "zod"; import { rpcError } from "../errors"; -import { getAutumn } from "../lib/autumn-client"; +import { + autumnCall, + BillingUnavailableError, + getAutumn, +} from "../lib/autumn-client"; import { logger } from "../lib/logger"; import { setTrackProperties } from "../middleware/track-mutation"; import { protectedProcedure, trackedSessionProcedure } from "../orpc"; @@ -263,7 +267,7 @@ async function upsertBillingControl< entry: BillingControlEntries[K]; operation: string; }): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { throw rpcError.badRequest( "Billing is turned off on this Databuddy instance." ); @@ -278,30 +282,28 @@ async function upsertBillingControl< ); } - try { - const autumn = getAutumn({ strict: true }); - const customer = await autumn.customers.getOrCreate({ customerId }); - if (customer.id !== customerId) { - throw new Error("The billing customer could not be verified"); - } - const existing = (customer.billingControls?.[args.key] ?? []) as Array<{ - featureId: string; - }>; - const merged = [ - ...existing.filter((e) => e.featureId !== args.entry.featureId), - args.entry, - ]; - await autumn.customers.update({ - customerId, - billingControls: { [args.key]: merged } as Record, - }); - } catch (error) { - logger.error( - { error, customerId, userId: args.context.user.id }, - `Failed to update ${args.operation} configuration` + const autumn = getAutumn(); + const customer = await autumnCall("customers.getOrCreate", () => + autumn.customers.getOrCreate({ customerId }) + ); + if (customer.id !== customerId) { + throw new BillingUnavailableError( + "The billing customer could not be verified" ); - throw rpcError.internal(`Failed to update ${args.operation} settings`); } + const existing = (customer.billingControls?.[args.key] ?? []) as Array<{ + featureId: string; + }>; + const merged = [ + ...existing.filter((e) => e.featureId !== args.entry.featureId), + args.entry, + ]; + await autumnCall("customers.update", () => + autumn.customers.update({ + customerId, + billingControls: { [args.key]: merged } as Record, + }) + ); } export const billingRouter = { diff --git a/packages/rpc/src/routers/business-context.test.ts b/packages/rpc/src/routers/business-context.test.ts index 4642dbe58f..b4b8fda586 100644 --- a/packages/rpc/src/routers/business-context.test.ts +++ b/packages/rpc/src/routers/business-context.test.ts @@ -703,7 +703,7 @@ test("unconfigured billing preserves the local billing policy and fails closed i delete process.env.AUTUMN_SECRET_KEY; expect(await access()).toMatchObject({ status: "allowed" }); process.env.NODE_ENV = "production"; - expect(await access()).toMatchObject({ status: "not-configured" }); + expect(await access()).toMatchObject({ status: "unavailable" }); expect(billingRequests).toEqual([]); }); diff --git a/packages/rpc/src/routers/feedback.selfhost.test.ts b/packages/rpc/src/routers/feedback.selfhost.test.ts index 6f3b39e0d5..f54f668ddc 100644 --- a/packages/rpc/src/routers/feedback.selfhost.test.ts +++ b/packages/rpc/src/routers/feedback.selfhost.test.ts @@ -1,7 +1,7 @@ import { expect, mock, test } from "bun:test"; import { createProcedureClient, os } from "@orpc/server"; import type { Context } from "../orpc"; -import { getAutumn } from "../lib/autumn-client"; +import { BillingUnavailableError, getAutumn } from "../lib/autumn-client"; const transaction = mock(() => { throw new Error("Unexpected credit mutation"); @@ -33,10 +33,8 @@ test("self-hosted rewards stop before credit mutation, and copied billing keys s }; try { getAutumn(); - getAutumn({ strict: true }); process.env.SELFHOST = "true"; - expect(() => getAutumn()).toThrow("disabled"); - expect(() => getAutumn({ strict: true })).toThrow("disabled"); + expect(() => getAutumn()).toThrow(BillingUnavailableError); const redeem = createProcedureClient(feedbackRouter.redeemCredits, { context: { user: { id: "synthetic-user" }, diff --git a/packages/rpc/src/routers/feedback.ts b/packages/rpc/src/routers/feedback.ts index c14f933d89..533f096833 100644 --- a/packages/rpc/src/routers/feedback.ts +++ b/packages/rpc/src/routers/feedback.ts @@ -1,4 +1,4 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { and, desc, eq, sql, withTransaction } from "@databuddy/db"; import type { db as DbType } from "@databuddy/db"; import { feedback, feedbackRedemptions } from "@databuddy/db/schema"; @@ -248,7 +248,7 @@ export const feedbackRouter = { }) ) .handler(async ({ context, input }) => { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { throw rpcError.badRequest( "Credit rewards are only available on Databuddy Cloud." ); diff --git a/packages/rpc/src/routers/insight-generation.ts b/packages/rpc/src/routers/insight-generation.ts index cb3e473d1d..00469235e2 100644 --- a/packages/rpc/src/routers/insight-generation.ts +++ b/packages/rpc/src/routers/insight-generation.ts @@ -1,4 +1,4 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; +import { billingMode } from "@databuddy/env/app"; import { and, db, @@ -37,7 +37,11 @@ import { rpcError } from "../errors"; import { setAuditOrganization } from "../lib/audit"; import { logger } from "../lib/logger"; import { getOrganizationOwnerId } from "../utils/organization"; -import { getAutumn } from "../lib/autumn-client"; +import { + autumnCall, + getAutumn, + isBillingUnavailable, +} from "../lib/autumn-client"; import { hasInvestigationAllowance, INVESTIGATION_USAGE, @@ -972,7 +976,7 @@ async function insertInsightRunOrFindActive( async function requireInvestigationsAccess( organizationId: string ): Promise { - if (readBooleanEnv("SELFHOST")) { + if (billingMode() !== "live") { if (!process.env.AI_GATEWAY_API_KEY?.trim()) { throw rpcError.badRequest( "AI is not set up on this Databuddy instance. Ask your administrator to configure it before running investigations." @@ -982,7 +986,9 @@ async function requireInvestigationsAccess( } const customerId = await getOrganizationOwnerId(organizationId); const customer = customerId - ? await getAutumn().customers.get({ customerId }) + ? await autumnCall("customers.get", () => + getAutumn().customers.get({ customerId }) + ) : null; if ( !hasInvestigationAllowance( @@ -1003,7 +1009,10 @@ async function hasInvestigationsAccess( try { await requireInvestigationsAccess(organizationId); return true; - } catch { + } catch (error) { + if (isBillingUnavailable(error)) { + throw error; + } return false; } } diff --git a/packages/rpc/src/routers/insights.ts b/packages/rpc/src/routers/insights.ts index 606d253c2b..0c818d20c4 100644 --- a/packages/rpc/src/routers/insights.ts +++ b/packages/rpc/src/routers/insights.ts @@ -1,5 +1,5 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; -import { getAutumn } from "../lib/autumn-client"; +import { billingMode, readBooleanEnv } from "@databuddy/env/app"; +import { autumnCall, getAutumn } from "../lib/autumn-client"; import { getBillingCustomerId } from "../utils/billing"; import { getClientIp } from "@databuddy/shared/utils/client-ip"; import { @@ -619,20 +619,24 @@ export async function appendInvestigationReply( if ( parsed.intent === "analysis" && author.authorId && - !readBooleanEnv("SELFHOST") + billingMode() === "live" ) { const customerId = await getBillingCustomerId( author.authorId, insight.organizationId ); - const customer = await getAutumn().customers.get({ customerId }); + const customer = await autumnCall("customers.get", () => + getAutumn().customers.get({ customerId }) + ); if ( customer.id !== customerId || !hasInvestigationAllowance( customer.balances[INVESTIGATION_USAGE.featureId] ) ) { - throw rpcError.badRequest( + throw rpcError.featureUnavailable( + INVESTIGATION_USAGE.featureId, + undefined, "Activate investigation billing to start a new analysis. Clarifications remain included." ); } diff --git a/packages/rpc/src/routers/organizations.ts b/packages/rpc/src/routers/organizations.ts index 943385eb49..b2e2a2da61 100644 --- a/packages/rpc/src/routers/organizations.ts +++ b/packages/rpc/src/routers/organizations.ts @@ -10,13 +10,18 @@ import { or, } from "@databuddy/db"; import { invitation, organization } from "@databuddy/db/schema"; +import { billingMode } from "@databuddy/env/app"; import { clearExpiredInvitationsSchema, getPendingInvitationsSchema, } from "@databuddy/validation"; import { z } from "zod"; import { rpcError } from "../errors"; -import { getAutumn, hasHostedBilling } from "../lib/autumn-client"; +import { + autumnCall, + getAutumn, + isBillingUnavailable, +} from "../lib/autumn-client"; import { logger } from "../lib/logger"; import { setTrackProperties } from "../middleware/track-mutation"; import { @@ -381,42 +386,47 @@ export const organizationsRouter = { }) .output(z.record(z.string(), z.unknown())) .handler(async ({ context }) => { - if (!hasHostedBilling()) { + if (billingMode() !== "live") { return { unlimited: true, canUserUpgrade: false }; } - const billing = await context.getBilling(); - const customerId = billing?.customerId ?? context.user.id; - const isOrganization = billing?.isOrganization ?? false; - const canUserUpgrade = billing?.canUserUpgrade ?? true; - + let billing: Awaited>; try { - const response = await getAutumn().check({ - customerId, - featureId: "events", - }); + billing = await context.getBilling(); + const customerId = billing?.customerId ?? context.user.id; + const isOrganization = billing?.isOrganization ?? false; + const canUserUpgrade = billing?.canUserUpgrade ?? true; + const response = await autumnCall("check", () => + getAutumn().check({ customerId, featureId: "events" }) + ); const b = response.balance; const unlimited = b?.unlimited ?? false; - const used = b?.usage ?? 0; const granted = b?.granted ?? 0; - const includedUsage = granted; - const overageAllowed = b?.overageAllowed ?? false; - const remaining = unlimited ? null : Math.max(0, b?.remaining ?? 0); - return { - used, + used: b?.usage ?? 0, limit: unlimited ? null : granted, unlimited, balance: b?.remaining ?? 0, - remaining, - includedUsage, - overageAllowed, + remaining: unlimited ? null : Math.max(0, b?.remaining ?? 0), + includedUsage: granted, + overageAllowed: b?.overageAllowed ?? false, isOrganizationUsage: isOrganization, canUserUpgrade, }; } catch (error) { - logger.error({ error }, "Failed to check usage"); - throw rpcError.internal("Failed to retrieve usage data"); + if (!isBillingUnavailable(error)) { + throw error; + } + logger.warn({ error }, "Usage is unavailable while billing is down"); + return { + unavailable: true, + ...(billing + ? { + isOrganizationUsage: billing.isOrganization, + canUserUpgrade: billing.canUserUpgrade, + } + : {}), + }; } }), @@ -465,7 +475,7 @@ export const organizationsRouter = { } } - if (!hasHostedBilling()) { + if (billingMode() !== "live") { return { planId: null, isOrganization: Boolean(context.organizationId), diff --git a/packages/rpc/src/types/billing.ts b/packages/rpc/src/types/billing.ts index 851ac8d19c..7ea41d94df 100644 --- a/packages/rpc/src/types/billing.ts +++ b/packages/rpc/src/types/billing.ts @@ -7,14 +7,14 @@ import { isFeatureAvailable, isWithinLimit, } from "@databuddy/shared/types/features"; +import { billingMode } from "@databuddy/env/app"; import { rpcError } from "../errors"; -import { hasHostedBilling } from "../lib/autumn-client"; function requireFeature( planId: string | undefined, feature: GatedFeatureId ): void { - if (!hasHostedBilling()) { + if (billingMode() !== "live") { return; } if (!isFeatureAvailable(planId ?? null, feature)) { @@ -41,7 +41,7 @@ export function requireUsageWithinLimit( feature: GatedFeatureId, currentUsage: number ): void { - if (!hasHostedBilling()) { + if (billingMode() !== "live") { return; } if (!isWithinLimit(planId ?? null, feature, currentUsage)) { diff --git a/packages/rpc/src/utils/autumn-balance.test.ts b/packages/rpc/src/utils/autumn-balance.test.ts index 507f0af9d5..43e40a026b 100644 --- a/packages/rpc/src/utils/autumn-balance.test.ts +++ b/packages/rpc/src/utils/autumn-balance.test.ts @@ -1,41 +1,59 @@ -import { afterEach, describe, expect, it, mock } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, mock } from "bun:test"; import { isDefinitiveAutumnBalanceFailure, updateAutumnBalance, } from "./autumn-balance"; const originalFetch = globalThis.fetch; +const originalEnv = process.env; + +beforeEach(() => { + process.env = { + ...originalEnv, + AUTUMN_SECRET_KEY: "synthetic-balance-key", + NODE_ENV: "test", + SELFHOST: "false", + }; +}); afterEach(() => { globalThis.fetch = originalFetch; + process.env = originalEnv; }); +function update(redemptionId: string) { + return updateAutumnBalance({ + amount: 2500, + customerId: "cus_1", + featureId: "events", + redemptionId, + }); +} + describe("updateAutumnBalance", () => { it("posts the balance update with a redemption-scoped idempotency key", async () => { const fetchMock = mock( - async (_url: string | URL | Request, _init?: RequestInit) => - new Response("{}", { status: 200 }) + async (_input: string | URL | Request, _init?: RequestInit) => + Response.json({ success: true }) ); - globalThis.fetch = fetchMock as typeof fetch; - - await updateAutumnBalance({ - amount: 2500, - customerId: "cus_1", - featureId: "events", - redemptionId: "redemption-1", - secretKey: "secret", + globalThis.fetch = Object.assign(fetchMock, { + preconnect: originalFetch.preconnect, }); + await update("redemption-1"); + expect(fetchMock).toHaveBeenCalledTimes(1); - const [url, init] = fetchMock.mock.calls[0]; - expect(url).toBe("https://api.useautumn.com/v1/balances.update"); - expect(init?.method).toBe("POST"); - expect(init?.headers).toMatchObject({ - Authorization: "Bearer secret", - "Content-Type": "application/json", - "Idempotency-Key": "feedback-redemption:redemption-1", - }); - expect(JSON.parse(String(init?.body))).toEqual({ + const call = fetchMock.mock.calls[0]; + if (!call) { + throw new Error("Expected an Autumn balance request"); + } + const [input, init] = call; + const request = new Request(input, init); + expect(request.url).toBe("https://api.useautumn.com/v1/balances.update"); + expect(request.headers.get("Idempotency-Key")).toBe( + "feedback-redemption:redemption-1" + ); + expect(await request.json()).toEqual({ customer_id: "cus_1", feature_id: "events", add_to_balance: 2500, @@ -48,97 +66,41 @@ describe("updateAutumnBalance", () => { [500, false], [503, false], ])("treats an HTTP %i Autumn response as definitive=%p for rollback", async (status, definitive) => { - globalThis.fetch = mock( - async () => new Response("autumn error", { status }) - ) as typeof fetch; - - let error: unknown; - try { - await updateAutumnBalance({ - amount: 10, - customerId: "cus_1", - featureId: "agent-credits", - redemptionId: "redemption-2", - secretKey: "secret", - }); - } catch (caught) { - error = caught; - } + globalThis.fetch = Object.assign( + mock(async () => Response.json({ message: "autumn error" }, { status })), + { preconnect: originalFetch.preconnect } + ); + + const error = await update("redemption-2").catch((caught) => caught); expect(error).toBeInstanceOf(Error); expect(isDefinitiveAutumnBalanceFailure(error)).toBe(definitive); }); - it("fails definitively without calling Autumn when no secret key is configured", async () => { - const fetchMock = mock(async () => new Response("{}", { status: 200 })); - globalThis.fetch = fetchMock as typeof fetch; - const originalSecret = process.env.AUTUMN_SECRET_KEY; - delete process.env.AUTUMN_SECRET_KEY; - - let error: unknown; - try { - await updateAutumnBalance({ - amount: 10, - customerId: "cus_1", - featureId: "agent-credits", - redemptionId: "redemption-4", - secretKey: null, - }); - } catch (caught) { - error = caught; - } finally { - if (originalSecret !== undefined) { - process.env.AUTUMN_SECRET_KEY = originalSecret; - } - } + it("fails definitively without calling Autumn when billing is not live", async () => { + const fetchMock = mock(async () => Response.json({ success: true })); + globalThis.fetch = Object.assign(fetchMock, { + preconnect: originalFetch.preconnect, + }); + process.env.SELFHOST = "true"; + + const error = await update("redemption-4").catch((caught) => caught); expect(isDefinitiveAutumnBalanceFailure(error)).toBe(true); expect(fetchMock).not.toHaveBeenCalled(); }); it("marks network failures as ambiguous so callers do not roll back spent credits", async () => { - globalThis.fetch = mock(async () => { - throw new Error("socket closed after write"); - }) as typeof fetch; - - let error: unknown; - try { - await updateAutumnBalance({ - amount: 10, - customerId: "cus_1", - featureId: "agent-credits", - redemptionId: "redemption-3", - secretKey: "secret", - }); - } catch (caught) { - error = caught; - } + globalThis.fetch = Object.assign( + mock(async () => { + throw new TypeError("socket closed after write"); + }), + { preconnect: originalFetch.preconnect } + ); + + const error = await update("redemption-3").catch((caught) => caught); expect(error).toBeInstanceOf(Error); expect(isDefinitiveAutumnBalanceFailure(error)).toBe(false); }); }); - -it("self-hosted balance writes fail definitively before any provider request", async () => { - const original = process.env.SELFHOST; - process.env.SELFHOST = "true"; - const request = mock(async () => new Response("{}")); - globalThis.fetch = request as typeof fetch; - try { - const definitiveFailure = await updateAutumnBalance({ - amount: 1, - customerId: "synthetic-customer", - featureId: "events", - redemptionId: "synthetic-redemption", - secretKey: "synthetic-stale-key", - }).catch(isDefinitiveAutumnBalanceFailure); - expect(definitiveFailure).toBe(true); - expect(request).not.toHaveBeenCalled(); - } finally { - if (original === undefined) { - Reflect.deleteProperty(process.env, "SELFHOST"); - } else { - process.env.SELFHOST = original; - } - } -}); diff --git a/packages/rpc/src/utils/autumn-balance.ts b/packages/rpc/src/utils/autumn-balance.ts index c03c5c4437..d9a03ed7e7 100644 --- a/packages/rpc/src/utils/autumn-balance.ts +++ b/packages/rpc/src/utils/autumn-balance.ts @@ -1,14 +1,21 @@ -import { readBooleanEnv } from "@databuddy/env/boolean"; - -const AUTUMN_BALANCE_TIMEOUT_MS = 10_000; +import { AutumnError, ResponseValidationError } from "autumn-js"; +import { getAutumn, isBillingUnavailable } from "../lib/autumn-client"; class AutumnBalanceUpdateError extends Error { readonly definitiveFailure: boolean; - constructor(message: string, definitiveFailure: boolean) { - super(message); - this.definitiveFailure = definitiveFailure; + constructor(cause: unknown) { + super( + `Autumn balance update failed: ${cause instanceof Error ? cause.message : String(cause)}`, + { cause } + ); this.name = "AutumnBalanceUpdateError"; + this.definitiveFailure = + isBillingUnavailable(cause) || + (cause instanceof AutumnError && + !(cause instanceof ResponseValidationError) && + cause.statusCode >= 400 && + cause.statusCode < 500); } } @@ -21,51 +28,21 @@ export async function updateAutumnBalance(input: { customerId: string; featureId: string; redemptionId: string; - secretKey?: string | null; }): Promise { - if (readBooleanEnv("SELFHOST")) { - throw new AutumnBalanceUpdateError( - "Hosted billing is disabled for self-hosted instances", - true - ); - } - const secretKey = input.secretKey ?? process.env.AUTUMN_SECRET_KEY; - if (!secretKey) { - throw new AutumnBalanceUpdateError("AUTUMN_SECRET_KEY is not set", true); - } - - const controller = new AbortController(); - const timeoutId = setTimeout( - () => controller.abort(), - AUTUMN_BALANCE_TIMEOUT_MS - ); try { - const response = await fetch( - "https://api.useautumn.com/v1/balances.update", + await getAutumn().balances.update( + { + customerId: input.customerId, + featureId: input.featureId, + addToBalance: input.amount, + }, { - method: "POST", headers: { - "Content-Type": "application/json", - Authorization: `Bearer ${secretKey}`, "Idempotency-Key": `feedback-redemption:${input.redemptionId}`, }, - body: JSON.stringify({ - customer_id: input.customerId, - feature_id: input.featureId, - add_to_balance: input.amount, - }), - signal: controller.signal, } ); - - if (!response.ok) { - const body = await response.text(); - throw new AutumnBalanceUpdateError( - `Autumn API ${response.status}: ${body}`, - response.status < 500 - ); - } - } finally { - clearTimeout(timeoutId); + } catch (error) { + throw new AutumnBalanceUpdateError(error); } } diff --git a/packages/rpc/src/utils/billing.test.ts b/packages/rpc/src/utils/billing.test.ts index b965d58a53..5597a41e22 100644 --- a/packages/rpc/src/utils/billing.test.ts +++ b/packages/rpc/src/utils/billing.test.ts @@ -7,11 +7,12 @@ let ownerId: string | null = OWNER_ID; let memberRole: string | null = null; const mockGetOrCreate = mock(async () => ({ subscriptions: [] })); -const mockLoggerError = mock(() => undefined); const mockGetOrganizationOwnerId = mock(async () => ownerId); const mockGetMemberRole = mock(async () => memberRole); +const actualAutumnClient = { ...(await import("../lib/autumn-client")) }; mock.module("../lib/autumn-client", () => ({ + ...actualAutumnClient, getAutumn: () => ({ customers: { getOrCreate: mockGetOrCreate, @@ -19,15 +20,6 @@ mock.module("../lib/autumn-client", () => ({ }), })); -mock.module("../lib/logger", () => ({ - logger: { - error: mockLoggerError, - info: mock(() => undefined), - warn: mock(() => undefined), - }, - record: (_name: string, fn: () => Promise | T) => fn(), -})); - mock.module("./organization", () => ({ getMemberRole: mockGetMemberRole, getOrganizationOwnerId: mockGetOrganizationOwnerId, @@ -43,7 +35,6 @@ beforeEach(() => { ownerId = OWNER_ID; memberRole = null; mockGetOrCreate.mockClear(); - mockLoggerError.mockClear(); mockGetOrganizationOwnerId.mockClear(); mockGetMemberRole.mockClear(); mockGetOrCreate.mockImplementation(async () => ({ subscriptions: [] })); @@ -58,7 +49,6 @@ describe("resolveBillingOwner", () => { await expect(resolveBillingOwner("user-1", null)).rejects.toThrow( "autumn unavailable" ); - expect(mockLoggerError).toHaveBeenCalledTimes(1); }); it("resolves and normalizes the active billing plan when Autumn succeeds", async () => { diff --git a/packages/rpc/src/utils/billing.ts b/packages/rpc/src/utils/billing.ts index 92bfc816a6..0f7fea6184 100644 --- a/packages/rpc/src/utils/billing.ts +++ b/packages/rpc/src/utils/billing.ts @@ -1,6 +1,5 @@ import { cacheNamespaces, cacheTags, cacheable } from "@databuddy/redis"; -import { getAutumn } from "../lib/autumn-client"; -import { logger } from "../lib/logger"; +import { autumnCall, getAutumn } from "../lib/autumn-client"; import { getMemberRole, getOrganizationOwnerId } from "./organization"; export interface BillingOwner { @@ -43,12 +42,9 @@ export async function resolveBillingOwner( } } - const customer = await getAutumn() - .customers.getOrCreate({ customerId }) - .catch((error: unknown) => { - logger.error({ error, customerId }, "Error resolving billing owner plan"); - throw error; - }); + const customer = await autumnCall("customers.getOrCreate", () => + getAutumn().customers.getOrCreate({ customerId }) + ); const subs = customer.subscriptions; const activeSub = diff --git a/packages/shared/src/billing.ts b/packages/shared/src/billing.ts index a916ada5aa..363ccf6723 100644 --- a/packages/shared/src/billing.ts +++ b/packages/shared/src/billing.ts @@ -95,3 +95,18 @@ export const investigationQuantitySchema = number() .int() .min(1) .max(INVESTIGATION_USAGE.maxPurchase); + +export class BillingUnavailableError extends Error { + readonly code = "billing_unavailable"; + + constructor(message: string, options?: ErrorOptions) { + super(message, options); + this.name = "BillingUnavailableError"; + } +} + +export function isBillingUnavailable( + error: unknown +): error is BillingUnavailableError { + return error instanceof BillingUnavailableError; +}