From e639a06306e433c1ac328ccb558e7b1885298ad6 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Mon, 5 Oct 2026 00:16:51 +0300 Subject: [PATCH 1/9] refactor(billing): one billing mode and one Autumn access layer billingMode() decides selfhost, live or disabled in one place. A single strict Autumn client and autumnCall() turn SDK, network and validation failures into BillingUnavailableError, mapped to 503 for ORPC and the agent routes instead of 500. Removes hasHostedBilling, isAgentBillingConfigured, the duplicated insights client and the raw balance fetch. Investigation allowance errors are 402 everywhere. --- apps/api/src/billing/autumn.ts | 7 +- apps/api/src/index.ts | 4 +- apps/api/src/routes/agent.ts | 17 ++ apps/api/src/routes/webhooks/autumn.test.ts | 14 +- apps/api/src/routes/webhooks/autumn.ts | 4 +- apps/basket/src/lib/billing.test.ts | 17 +- apps/basket/src/lib/billing.ts | 33 ++- .../generation-billing.integration.test.ts | 8 +- .../investigation-billing.integration.test.ts | 11 +- apps/insights/src/investigation-billing.ts | 225 ++++++++---------- .../resume-clarification.integration.test.ts | 14 +- packages/ai/src/ai/agents/execution.test.ts | 5 +- packages/ai/src/ai/agents/execution.ts | 53 +++-- .../ai/src/ai/config/enrich-context.test.ts | 14 +- packages/ai/src/ai/config/enrich-context.ts | 4 +- packages/ai/src/query/index.ts | 4 +- packages/auth/src/auth.ts | 30 +-- packages/rpc/src/lib/autumn-client.ts | 80 ++++--- .../rpc/src/lib/business-context-access.ts | 33 +-- packages/rpc/src/orpc.ts | 47 ++-- packages/rpc/src/procedures/with-workspace.ts | 4 +- packages/rpc/src/routers/billing.test.ts | 6 +- packages/rpc/src/routers/billing.ts | 52 ++-- .../rpc/src/routers/business-context.test.ts | 2 +- .../rpc/src/routers/feedback.selfhost.test.ts | 6 +- packages/rpc/src/routers/feedback.ts | 4 +- .../rpc/src/routers/insight-generation.ts | 19 +- packages/rpc/src/routers/insights.ts | 14 +- packages/rpc/src/routers/organizations.ts | 61 ++--- packages/rpc/src/types/billing.ts | 6 +- packages/rpc/src/utils/autumn-balance.test.ts | 146 ++++-------- packages/rpc/src/utils/autumn-balance.ts | 65 ++--- packages/rpc/src/utils/billing.test.ts | 14 +- packages/rpc/src/utils/billing.ts | 12 +- packages/shared/src/billing.ts | 15 ++ 35 files changed, 535 insertions(+), 515 deletions(-) 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.ts b/apps/api/src/routes/agent.ts index f927f3eb20..0ee089a2d8 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) { @@ -573,6 +576,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 +1239,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..e7f0a448c0 100644 --- a/apps/api/src/routes/webhooks/autumn.test.ts +++ b/apps/api/src/routes/webhooks/autumn.test.ts @@ -213,10 +213,16 @@ 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" }, + services: 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..24f726b14b 100644 --- a/apps/basket/src/lib/billing.test.ts +++ b/apps/basket/src/lib/billing.test.ts @@ -1,4 +1,5 @@ import { afterEach, beforeEach, describe, expect, test, vi } from "vitest"; +import { BillingUnavailableError } from "@databuddy/shared/billing"; import { EvlogError } from "evlog"; const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({ @@ -13,7 +14,8 @@ const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({ mockLoggerWarn: vi.fn(() => {}), })); -vi.mock("@databuddy/rpc/autumn", () => ({ +vi.mock("@databuddy/rpc/autumn", async (importOriginal) => ({ + ...(await importOriginal()), getAutumn: () => ({ check: mockCheck }), })); @@ -35,6 +37,7 @@ 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(); @@ -133,6 +136,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..c8ff7d8011 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 } 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)) { + 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/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..999df82a85 100644 --- a/apps/insights/src/investigation-billing.ts +++ b/apps/insights/src/investigation-billing.ts @@ -1,8 +1,14 @@ -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 { @@ -12,56 +18,26 @@ export interface InvestigationBilling { 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 +49,23 @@ export async function canRunInvestigation( if (billing.mode === "unconfigured") { return true; } - if (!billing.customerId) { - throw new Error("The investigation billing customer is unavailable"); + const customerId = billing.customerId; + if (!customerId) { + throw new BillingUnavailableError( + "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 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; } @@ -127,29 +110,36 @@ export async function reserveInvestigationCharge( return reservation; } assertInvestigationReservationActive(reservation); - if (!reservation.customerId) { - throw new Error("The investigation billing customer is unavailable"); + const customerId = reservation.customerId; + if (!customerId) { + throw new BillingUnavailableError( + "The investigation billing customer is unavailable" + ); } - const autumn = client ?? createInvestigationBillingClient(); + 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 +200,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.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..0b8ca90b0a 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,7 +386,7 @@ 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(); @@ -389,35 +394,37 @@ export const organizationsRouter = { const isOrganization = billing?.isOrganization ?? false; const canUserUpgrade = billing?.canUserUpgrade ?? true; - try { - const response = await 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); - + const response = await autumnCall("check", () => + getAutumn().check({ customerId, featureId: "events" }) + ).catch((error: unknown) => { + if (isBillingUnavailable(error)) { + logger.warn({ error }, "Usage is unavailable while billing is down"); + return null; + } + throw error; + }); + if (!response) { return { - used, - limit: unlimited ? null : granted, - unlimited, - balance: b?.remaining ?? 0, - remaining, - includedUsage, - overageAllowed, + unavailable: true, isOrganizationUsage: isOrganization, canUserUpgrade, }; - } catch (error) { - logger.error({ error }, "Failed to check usage"); - throw rpcError.internal("Failed to retrieve usage data"); } + + const b = response.balance; + const unlimited = b?.unlimited ?? false; + const granted = b?.granted ?? 0; + return { + used: b?.usage ?? 0, + limit: unlimited ? null : granted, + unlimited, + balance: b?.remaining ?? 0, + remaining: unlimited ? null : Math.max(0, b?.remaining ?? 0), + includedUsage: granted, + overageAllowed: b?.overageAllowed ?? false, + isOrganizationUsage: isOrganization, + canUserUpgrade, + }; }), getBillingContext: publicProcedure @@ -465,7 +472,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..b389e049f9 100644 --- a/packages/rpc/src/utils/autumn-balance.test.ts +++ b/packages/rpc/src/utils/autumn-balance.test.ts @@ -1,41 +1,53 @@ -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 = fetchMock as unknown as typeof fetch; + + 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 [input, init] = fetchMock.mock.calls[0]; + 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,49 +60,22 @@ 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 = mock(async () => + Response.json({ message: "autumn error" }, { status }) + ) as unknown as typeof fetch; + + 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 = fetchMock as unknown as typeof fetch; + process.env.SELFHOST = "true"; + + const error = await update("redemption-4").catch((caught) => caught); expect(isDefinitiveAutumnBalanceFailure(error)).toBe(true); expect(fetchMock).not.toHaveBeenCalled(); @@ -98,47 +83,12 @@ describe("updateAutumnBalance", () => { 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; - } + throw new TypeError("socket closed after write"); + }) as unknown as typeof fetch; + + 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; +} From 42a7a3a16a8a46bcd6fcbf5f1ed0ce77195db096 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Mon, 5 Oct 2026 02:28:53 +0300 Subject: [PATCH 2/9] refactor(billing): narrow investigation customers and keep test service reads live Model fixed billing customers as a discriminated union and keep the webhook mock compatible with live service configuration. Includes the billing cleanup from 669a8c2ce. --- apps/api/src/routes/webhooks/autumn.test.ts | 4 +++- apps/insights/src/investigation-billing.ts | 25 ++++++--------------- 2 files changed, 10 insertions(+), 19 deletions(-) diff --git a/apps/api/src/routes/webhooks/autumn.test.ts b/apps/api/src/routes/webhooks/autumn.test.ts index e7f0a448c0..79dd3919b4 100644 --- a/apps/api/src/routes/webhooks/autumn.test.ts +++ b/apps/api/src/routes/webhooks/autumn.test.ts @@ -219,7 +219,9 @@ vi.mock("@databuddy/env/app", async (importOriginal) => { ...actual, config: { email: { alertsFrom: "alerts@databuddy.cc" }, - services: actual.config.services, + get services() { + return actual.config.services; + }, }, }; }); diff --git a/apps/insights/src/investigation-billing.ts b/apps/insights/src/investigation-billing.ts index 999df82a85..4e6f3efe00 100644 --- a/apps/insights/src/investigation-billing.ts +++ b/apps/insights/src/investigation-billing.ts @@ -11,10 +11,9 @@ import { INVESTIGATION_USAGE } from "@databuddy/shared/billing"; 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; @@ -49,12 +48,7 @@ export async function canRunInvestigation( if (billing.mode === "unconfigured") { return true; } - const customerId = billing.customerId; - if (!customerId) { - throw new BillingUnavailableError( - "The investigation billing customer is unavailable" - ); - } + const { customerId } = billing; const result = await autumnCall("check", () => (client ?? getAutumn()).check({ customerId, @@ -76,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") @@ -110,12 +104,7 @@ export async function reserveInvestigationCharge( return reservation; } assertInvestigationReservationActive(reservation); - const customerId = reservation.customerId; - if (!customerId) { - throw new BillingUnavailableError( - "The investigation billing customer is unavailable" - ); - } + 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: From d92aa5590fb8de6aded643d987ea4a77edcd0b57 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 6 Oct 2026 13:51:21 +0300 Subject: [PATCH 3/9] test(rpc): cover billing outages at HTTP boundary --- .../rpc/src/routers/billing.boundary.test.ts | 115 ++++++++++++++++++ 1 file changed, 115 insertions(+) create mode 100644 packages/rpc/src/routers/billing.boundary.test.ts 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..87ab1c3262 --- /dev/null +++ b/packages/rpc/src/routers/billing.boundary.test.ts @@ -0,0 +1,115 @@ +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"; + +mock.module("@databuddy/auth", () => ({ + auth: { api: { getSession: async () => null } }, +})); +mock.module("@databuddy/api-keys/resolve", () => ({ + getApiKeyFromHeader: async () => null, +})); +mock.module("@databuddy/db", () => ({ db: {} })); +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" }), +})); +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 transport = spyOn(globalThis, "fetch").mockImplementation( + async (input, init) => { + 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(url.pathname).toBe("/v1/customers.get_or_create"); + expect(await request.json()).toMatchObject({ + customer_id: "owner-example", + }); + return Response.json({}, { status }); + } +); + +const { billingRouter } = await import("./billing"); +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) { + delete 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", + ]); +}); From e76da7105b00d3b342c9a06b0b8e310c576ceed2 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 6 Oct 2026 13:56:29 +0300 Subject: [PATCH 4/9] fix(dashboard): explain unavailable billing usage --- apps/dashboard/app/(main)/billing/page.tsx | 19 ++++++++++++++++- .../analytics/event-limit-indicator.tsx | 21 +++++++++++++++++-- 2 files changed, 37 insertions(+), 3 deletions(-) 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..643c353230 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); From 46eec81ac2c8f750f8d8dea4f10cec35d3d7df2a Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:01:02 +0000 Subject: [PATCH 5/9] fix(dashboard): use shared rounding for unavailable usage --- apps/dashboard/components/analytics/event-limit-indicator.tsx | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/dashboard/components/analytics/event-limit-indicator.tsx b/apps/dashboard/components/analytics/event-limit-indicator.tsx index 643c353230..1f4c6b9b9b 100644 --- a/apps/dashboard/components/analytics/event-limit-indicator.tsx +++ b/apps/dashboard/components/analytics/event-limit-indicator.tsx @@ -27,7 +27,7 @@ export function EventLimitIndicator() { if (data.unavailable === true) { return (
Date: Wed, 7 Oct 2026 07:17:47 +0000 Subject: [PATCH 6/9] fix(api): preserve billing availability failures --- .../src/routes/agent-business-context.test.ts | 29 ++++++++++++++--- apps/api/src/routes/agent.ts | 15 ++++++--- apps/basket/src/lib/billing.test.ts | 32 ++++++++++++++++--- apps/basket/src/lib/billing.ts | 4 +-- 4 files changed, 65 insertions(+), 15 deletions(-) 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 0ee089a2d8..c2af9fe98e 100644 --- a/apps/api/src/routes/agent.ts +++ b/apps/api/src/routes/agent.ts @@ -427,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)); @@ -535,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, diff --git a/apps/basket/src/lib/billing.test.ts b/apps/basket/src/lib/billing.test.ts index 24f726b14b..670e7db7f1 100644 --- a/apps/basket/src/lib/billing.test.ts +++ b/apps/basket/src/lib/billing.test.ts @@ -1,4 +1,5 @@ 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"; @@ -14,10 +15,16 @@ const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({ mockLoggerWarn: vi.fn(() => {}), })); -vi.mock("@databuddy/rpc/autumn", async (importOriginal) => ({ - ...(await importOriginal()), - 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: () => ({ @@ -33,16 +40,31 @@ vi.mock("@lib/tracing", () => ({ })); const { checkAutumnUsage } = await import("./billing"); +const originalAutumnSecretKey = config.services.autumnSecretKey; describe("checkAutumnUsage", () => { beforeEach(() => { vi.stubEnv("SELFHOST", "false"); vi.stubEnv("AUTUMN_SECRET_KEY", "am_sk_test_synthetic"); + config.services.autumnSecretKey = "am_sk_test_synthetic"; mockCheck.mockReset(); mockLoggerSet.mockReset(); mockLoggerWarn.mockReset(); }); - afterEach(() => vi.unstubAllEnvs()); + afterEach(() => { + vi.unstubAllEnvs(); + config.services.autumnSecretKey = originalAutumnSecretKey; + }); + + test("hosted events reject missing billing configuration", async () => { + vi.stubEnv("NODE_ENV", "production"); + vi.stubEnv("AUTUMN_SECRET_KEY", undefined); + config.services.autumnSecretKey = 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"); diff --git a/apps/basket/src/lib/billing.ts b/apps/basket/src/lib/billing.ts index c8ff7d8011..288165bfb5 100644 --- a/apps/basket/src/lib/billing.ts +++ b/apps/basket/src/lib/billing.ts @@ -1,4 +1,4 @@ -import { billingMode } from "@databuddy/env/app"; +import { billingMode, config } from "@databuddy/env/app"; import { autumnCall, getAutumn, @@ -66,7 +66,7 @@ export function checkAutumnUsage( if (error instanceof EvlogError) { throw error; } - if (isBillingUnavailable(error)) { + if (isBillingUnavailable(error) && config.services.autumnSecretKey) { log.set({ billing: { allowed: true, checkFailed: true } }); captureError(error, { message: "Autumn check failed, accepting event", From 426d4b5e3547b3e28a601af0b836b747a48be7aa Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 07:44:56 +0000 Subject: [PATCH 7/9] test(basket): remove ineffective billing config writes --- apps/basket/src/lib/billing.test.ts | 4 ---- 1 file changed, 4 deletions(-) diff --git a/apps/basket/src/lib/billing.test.ts b/apps/basket/src/lib/billing.test.ts index 670e7db7f1..b66971c49e 100644 --- a/apps/basket/src/lib/billing.test.ts +++ b/apps/basket/src/lib/billing.test.ts @@ -40,26 +40,22 @@ vi.mock("@lib/tracing", () => ({ })); const { checkAutumnUsage } = await import("./billing"); -const originalAutumnSecretKey = config.services.autumnSecretKey; describe("checkAutumnUsage", () => { beforeEach(() => { vi.stubEnv("SELFHOST", "false"); vi.stubEnv("AUTUMN_SECRET_KEY", "am_sk_test_synthetic"); - config.services.autumnSecretKey = "am_sk_test_synthetic"; mockCheck.mockReset(); mockLoggerSet.mockReset(); mockLoggerWarn.mockReset(); }); afterEach(() => { vi.unstubAllEnvs(); - config.services.autumnSecretKey = originalAutumnSecretKey; }); test("hosted events reject missing billing configuration", async () => { vi.stubEnv("NODE_ENV", "production"); vi.stubEnv("AUTUMN_SECRET_KEY", undefined); - config.services.autumnSecretKey = undefined; await expect(checkAutumnUsage("cust_1", "events")).rejects.toMatchObject({ status: 503, message: "Billing check unavailable", From 7ab83051a69d1c4bb305ceab9f396f4dc39832a7 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 07:51:23 +0000 Subject: [PATCH 8/9] test(rpc): type native billing transports --- .../rpc/src/routers/billing.boundary.test.ts | 13 +++++--- packages/rpc/src/utils/autumn-balance.test.ts | 30 +++++++++++++------ 2 files changed, 30 insertions(+), 13 deletions(-) diff --git a/packages/rpc/src/routers/billing.boundary.test.ts b/packages/rpc/src/routers/billing.boundary.test.ts index 87ab1c3262..6305f72a2a 100644 --- a/packages/rpc/src/routers/billing.boundary.test.ts +++ b/packages/rpc/src/routers/billing.boundary.test.ts @@ -43,8 +43,11 @@ mock.module("@databuddy/db/clickhouse", () => ({ let status = 202; const requests: string[] = []; const originalSecret = process.env.AUTUMN_SECRET_KEY; -const transport = spyOn(globalThis, "fetch").mockImplementation( - async (input, init) => { +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); @@ -54,8 +57,10 @@ const transport = spyOn(globalThis, "fetch").mockImplementation( customer_id: "owner-example", }); return Response.json({}, { status }); - } + }, + { preconnect: globalThis.fetch.preconnect } ); +const transport = spyOn(globalThis, "fetch").mockImplementation(fetcher); const { billingRouter } = await import("./billing"); const context = { @@ -84,7 +89,7 @@ beforeEach(() => { afterAll(() => { transport.mockRestore(); if (originalSecret === undefined) { - delete process.env.AUTUMN_SECRET_KEY; + Reflect.deleteProperty(process.env, "AUTUMN_SECRET_KEY"); } else { process.env.AUTUMN_SECRET_KEY = originalSecret; } diff --git a/packages/rpc/src/utils/autumn-balance.test.ts b/packages/rpc/src/utils/autumn-balance.test.ts index b389e049f9..43e40a026b 100644 --- a/packages/rpc/src/utils/autumn-balance.test.ts +++ b/packages/rpc/src/utils/autumn-balance.test.ts @@ -36,12 +36,18 @@ describe("updateAutumnBalance", () => { async (_input: string | URL | Request, _init?: RequestInit) => Response.json({ success: true }) ); - globalThis.fetch = fetchMock as unknown as typeof fetch; + globalThis.fetch = Object.assign(fetchMock, { + preconnect: originalFetch.preconnect, + }); await update("redemption-1"); expect(fetchMock).toHaveBeenCalledTimes(1); - const [input, init] = fetchMock.mock.calls[0]; + 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( @@ -60,9 +66,10 @@ describe("updateAutumnBalance", () => { [500, false], [503, false], ])("treats an HTTP %i Autumn response as definitive=%p for rollback", async (status, definitive) => { - globalThis.fetch = mock(async () => - Response.json({ message: "autumn error" }, { status }) - ) as unknown as typeof fetch; + globalThis.fetch = Object.assign( + mock(async () => Response.json({ message: "autumn error" }, { status })), + { preconnect: originalFetch.preconnect } + ); const error = await update("redemption-2").catch((caught) => caught); @@ -72,7 +79,9 @@ describe("updateAutumnBalance", () => { it("fails definitively without calling Autumn when billing is not live", async () => { const fetchMock = mock(async () => Response.json({ success: true })); - globalThis.fetch = fetchMock as unknown as typeof fetch; + globalThis.fetch = Object.assign(fetchMock, { + preconnect: originalFetch.preconnect, + }); process.env.SELFHOST = "true"; const error = await update("redemption-4").catch((caught) => caught); @@ -82,9 +91,12 @@ describe("updateAutumnBalance", () => { }); it("marks network failures as ambiguous so callers do not roll back spent credits", async () => { - globalThis.fetch = mock(async () => { - throw new TypeError("socket closed after write"); - }) as unknown as typeof fetch; + 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); From 130294ac31bbbfe3d0d3f0a8cfc17346098a5716 Mon Sep 17 00:00:00 2001 From: iza <59828082+izadoesdev@users.noreply.github.com> Date: Wed, 7 Oct 2026 08:00:22 +0000 Subject: [PATCH 9/9] fix(rpc): show usage outages during owner lookup --- .../rpc/src/routers/billing.boundary.test.ts | 91 ++++++++++++++++++- packages/rpc/src/routers/organizations.ts | 65 ++++++------- 2 files changed, 123 insertions(+), 33 deletions(-) diff --git a/packages/rpc/src/routers/billing.boundary.test.ts b/packages/rpc/src/routers/billing.boundary.test.ts index 6305f72a2a..ecbb88e38e 100644 --- a/packages/rpc/src/routers/billing.boundary.test.ts +++ b/packages/rpc/src/routers/billing.boundary.test.ts @@ -2,6 +2,7 @@ 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 } }, @@ -9,7 +10,8 @@ mock.module("@databuddy/auth", () => ({ mock.module("@databuddy/api-keys/resolve", () => ({ getApiKeyFromHeader: async () => null, })); -mock.module("@databuddy/db", () => ({ db: {} })); +const database = { ...(await import("@databuddy/db")), db: {} }; +mock.module("@databuddy/db", () => database); mock.module("@databuddy/services/audit", () => ({ appendAuditEvent: async () => undefined, appendAuditEventInTransaction: async () => undefined, @@ -25,6 +27,10 @@ mock.module("../utils/billing", () => ({ })); mock.module("../procedures/with-workspace", () => ({ withWorkspace: async () => ({ organizationId: "org-example", role: "admin" }), + withPublicWorkspace: async () => ({ + organizationId: "org-example", + role: "admin", + }), })); mock.module("../lib/logger", () => ({ logger: { @@ -52,10 +58,29 @@ const fetcher = Object.assign( const url = new URL(request.url); requests.push(url.pathname); expect(url.origin).toBe("https://api.useautumn.com"); - expect(url.pathname).toBe("/v1/customers.get_or_create"); + 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 } @@ -63,6 +88,7 @@ const fetcher = Object.assign( 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" }, @@ -118,3 +144,64 @@ test.each([ "/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/organizations.ts b/packages/rpc/src/routers/organizations.ts index 0b8ca90b0a..b2e2a2da61 100644 --- a/packages/rpc/src/routers/organizations.ts +++ b/packages/rpc/src/routers/organizations.ts @@ -389,42 +389,45 @@ export const organizationsRouter = { 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; - - const response = await autumnCall("check", () => - getAutumn().check({ customerId, featureId: "events" }) - ).catch((error: unknown) => { - if (isBillingUnavailable(error)) { - logger.warn({ error }, "Usage is unavailable while billing is down"); - return null; - } - throw error; - }); - if (!response) { + let billing: Awaited>; + try { + 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 granted = b?.granted ?? 0; return { - unavailable: true, + used: b?.usage ?? 0, + limit: unlimited ? null : granted, + unlimited, + balance: b?.remaining ?? 0, + remaining: unlimited ? null : Math.max(0, b?.remaining ?? 0), + includedUsage: granted, + overageAllowed: b?.overageAllowed ?? false, isOrganizationUsage: isOrganization, canUserUpgrade, }; + } catch (error) { + 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, + } + : {}), + }; } - - const b = response.balance; - const unlimited = b?.unlimited ?? false; - const granted = b?.granted ?? 0; - return { - used: b?.usage ?? 0, - limit: unlimited ? null : granted, - unlimited, - balance: b?.remaining ?? 0, - remaining: unlimited ? null : Math.max(0, b?.remaining ?? 0), - includedUsage: granted, - overageAllowed: b?.overageAllowed ?? false, - isOrganizationUsage: isOrganization, - canUserUpgrade, - }; }), getBillingContext: publicProcedure