Skip to content
7 changes: 5 additions & 2 deletions apps/api/src/billing/autumn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/index.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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,
Expand Down
29 changes: 25 additions & 4 deletions apps/api/src/routes/agent-business-context.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { MockLanguageModelV3 } from "ai/test";
import { BillingUnavailableError } from "@databuddy/shared/billing";
import {
type OrganizationBusinessProfile,
PROFILE_ORIGIN_PROVENANCE,
Expand Down Expand Up @@ -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.`]) {
Expand All @@ -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 () => {
Expand All @@ -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);
Expand Down Expand Up @@ -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(
Expand Down
32 changes: 28 additions & 4 deletions apps/api/src/routes/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -424,14 +427,21 @@ function createAgentUsageInjector(
});
}

function createPlainTextStreamResponse(
stream: AsyncIterable<string>
): Response {
async function createPlainTextStreamResponse(
stream: AsyncGenerator<string>
): Promise<Response> {
let first = await stream.next();
while (!(first.done || first.value)) {
first = await stream.next();
Comment thread
izadoesdev marked this conversation as resolved.
Comment thread
izadoesdev marked this conversation as resolved.
}
const encoder = new TextEncoder();
return new Response(
new ReadableStream<Uint8Array>({
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));
Expand Down Expand Up @@ -532,7 +542,7 @@ export const agent = new Elysia({ prefix: "/v1/agent" })
}
: createSessionAgentActor(user, request.headers);
if (body.stream) {
return createPlainTextStreamResponse(
return await createPlainTextStreamResponse(
streamDatabuddyAgent({
actor,
conversationId,
Expand Down Expand Up @@ -573,6 +583,13 @@ export const agent = new Elysia({ prefix: "/v1/agent" })
error_type: getErrorName(error),
source: "slack",
});
if (isBillingUnavailable(error)) {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
return jsonError(
503,
"BILLING_UNAVAILABLE",
BILLING_UNAVAILABLE_MESSAGE
);
}
return jsonError(500, "INTERNAL_ERROR", INTERNAL_AGENT_ERROR_MESSAGE);
}
},
Expand Down Expand Up @@ -1229,6 +1246,13 @@ export const agent = new Elysia({ prefix: "/v1/agent" })
...(user?.id ? { agent_user_id: user.id } : {}),
error_type: getErrorName(error),
});
if (isBillingUnavailable(error)) {
return jsonError(
503,
"BILLING_UNAVAILABLE",
BILLING_UNAVAILABLE_MESSAGE
);
}
return jsonError(500, "INTERNAL_ERROR", INTERNAL_AGENT_ERROR_MESSAGE);
}
})();
Expand Down
16 changes: 12 additions & 4 deletions apps/api/src/routes/webhooks/autumn.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -213,10 +213,18 @@ vi.mock("@databuddy/email", async (importOriginal) => ({
UsageLimitEmail: vi.fn(() => ({ type: "limit" })),
}));

vi.mock("@databuddy/env/app", async (importOriginal) => ({
...(await importOriginal<typeof import("@databuddy/env/app")>()),
config: { email: { alertsFrom: "alerts@databuddy.cc" } },
}));
vi.mock("@databuddy/env/app", async (importOriginal) => {
const actual = await importOriginal<typeof import("@databuddy/env/app")>();
return {
...actual,
config: {
email: { alertsFrom: "alerts@databuddy.cc" },
get services() {
return actual.config.services;
},
},
};
});

vi.mock("@databuddy/notifications", () => ({
SlackProvider: class {
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/routes/webhooks/autumn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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: {
Expand Down
41 changes: 37 additions & 4 deletions apps/basket/src/lib/billing.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
import { afterEach, beforeEach, describe, expect, test, vi } from "vitest";
import { config } from "@databuddy/env/app";
import { BillingUnavailableError } from "@databuddy/shared/billing";
import { EvlogError } from "evlog";

const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({
Expand All @@ -13,9 +15,16 @@ const { mockCheck, mockLoggerSet, mockLoggerWarn } = vi.hoisted(() => ({
mockLoggerWarn: vi.fn(() => {}),
}));

vi.mock("@databuddy/rpc/autumn", () => ({
getAutumn: () => ({ check: mockCheck }),
}));
vi.mock("@databuddy/rpc/autumn", async (importOriginal) => {
const actual = await importOriginal<typeof import("@databuddy/rpc/autumn")>();
return {
...actual,
getAutumn: () =>
config.services.autumnSecretKey
? { check: mockCheck }
: actual.getAutumn(),
};
});

vi.mock("evlog/elysia", () => ({
useLogger: () => ({
Expand All @@ -35,11 +44,23 @@ const { checkAutumnUsage } = await import("./billing");
describe("checkAutumnUsage", () => {
beforeEach(() => {
vi.stubEnv("SELFHOST", "false");
vi.stubEnv("AUTUMN_SECRET_KEY", "am_sk_test_synthetic");
mockCheck.mockReset();
mockLoggerSet.mockReset();
mockLoggerWarn.mockReset();
});
afterEach(() => vi.unstubAllEnvs());
afterEach(() => {
vi.unstubAllEnvs();
});

test("hosted events reject missing billing configuration", async () => {
vi.stubEnv("NODE_ENV", "production");
vi.stubEnv("AUTUMN_SECRET_KEY", undefined);
await expect(checkAutumnUsage("cust_1", "events")).rejects.toMatchObject({
status: 503,
message: "Billing check unavailable",
});
});

test("self-hosted events skip hosted billing", async () => {
vi.stubEnv("SELFHOST", "true");
Expand Down Expand Up @@ -133,6 +154,18 @@ describe("checkAutumnUsage", () => {
});
});

test("an Autumn outage accepts the event instead of dropping it", async () => {
mockCheck.mockRejectedValue(
new BillingUnavailableError("Autumn check failed")
);
await expect(checkAutumnUsage("cust_1", "events")).resolves.toEqual({
allowed: true,
});
expect(mockLoggerSet).toHaveBeenCalledWith({
billing: { allowed: true, checkFailed: true },
});
});

test("logs checkFailed on API error", async () => {
mockCheck.mockRejectedValue(new Error("timeout"));
await expect(checkAutumnUsage("cust_1", "events")).rejects.toThrow(
Expand Down
33 changes: 23 additions & 10 deletions apps/basket/src/lib/billing.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
import { readBooleanEnv } from "@databuddy/env/boolean";
import { getAutumn } from "@databuddy/rpc/autumn";
import { billingMode, config } from "@databuddy/env/app";
import {
autumnCall,
getAutumn,
isBillingUnavailable,
} from "@databuddy/rpc/autumn";
import { basketErrors } from "@lib/structured-errors";
import { captureError, record } from "@lib/tracing";
import { EvlogError } from "evlog";
Expand All @@ -15,21 +19,23 @@ export function checkAutumnUsage(
properties?: Record<string, unknown>,
quantity = 1
): Promise<BillingResult> {
if (readBooleanEnv("SELFHOST")) {
if (billingMode() !== "live") {
return Promise.resolve({ allowed: true });
}
return record("checkAutumnUsage", async (): Promise<BillingResult> => {
const log = useLogger();

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;
Expand Down Expand Up @@ -60,6 +66,13 @@ export function checkAutumnUsage(
if (error instanceof EvlogError) {
throw error;
}
if (isBillingUnavailable(error) && config.services.autumnSecretKey) {
log.set({ billing: { allowed: true, checkFailed: true } });
captureError(error, {
message: "Autumn check failed, accepting event",
});
return { allowed: true };
}

log.set({ billing: { allowed: false, checkFailed: true } });
captureError(error, {
Expand Down
19 changes: 18 additions & 1 deletion apps/dashboard/app/(main)/billing/page.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ import {
CrownIcon,
PlusIcon,
TrendUpIcon,
WarningIcon,
XMarkIcon as XIcon,
} from "@databuddy/ui/icons";
import {
Expand All @@ -62,6 +63,7 @@ const INTELLIGENCE_PLAN_ID_SET = new Set<string>(
interface OrgUsageData {
balance?: number | null;
includedUsage?: number | null;
unavailable?: boolean;
unlimited: boolean;
}

Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -598,6 +600,21 @@ export default function BillingPage() {
</Card>
)}

{orgUsage?.unavailable === true && (
<Card className="border-warning/30 bg-warning/5" role="status">
<Card.Content className="flex items-start gap-2 py-3">
<WarningIcon
aria-hidden="true"
className="size-4 shrink-0 text-warning"
/>
<Text className="text-pretty" tone="muted" variant="caption">
Billing usage and event overage estimates are temporarily
unavailable.
</Text>
</Card.Content>
</Card>
)}

{usageStats.length === 0 ? (
<Card>
<Card.Content className="py-8">
Expand Down
Loading
Loading