From a73d996f65fc29de29d4364526c83e84b0732030 Mon Sep 17 00:00:00 2001 From: Jesse_Chen Date: Mon, 31 Aug 2026 04:08:20 +0800 Subject: [PATCH] feat(billing): charge personal reports with durable settlement --- frontend/src/app/api/reports/route.ts | 28 ++++++ frontend/src/lib/personal-report-billing.ts | 27 ++++++ frontend/src/lib/personal-report-codes.ts | 3 + .../src/lib/personal-report-generation.ts | 5 +- .../src/lib/personal-report-route-core.ts | 86 ++++++++++++++++--- .../src/lib/personal-report-worker-core.ts | 22 +++++ frontend/src/lib/personal-report-worker.ts | 31 ++++++- frontend/src/mastra/personal-report.ts | 17 ++++ frontend/tests/personal-report-api.test.ts | 81 +++++++++++++++++ 9 files changed, 286 insertions(+), 14 deletions(-) create mode 100644 frontend/src/lib/personal-report-billing.ts diff --git a/frontend/src/app/api/reports/route.ts b/frontend/src/app/api/reports/route.ts index b6a4db04..d21b0c1d 100644 --- a/frontend/src/app/api/reports/route.ts +++ b/frontend/src/app/api/reports/route.ts @@ -24,6 +24,8 @@ import { } from "@/lib/personal-report-service"; import { createSupabasePersonalReportJobService } from "@/lib/personal-report-job-service"; import { createAdminSupabaseClient } from "@/lib/supabase/admin"; +import { authorizeUsage, completeUsage, releaseUsage } from "@/lib/consultation-billing"; +import { resolveFeaturePricing } from "@/lib/feature-pricing"; import { isSupabaseConfigurationError } from "@/lib/supabase/config"; import { createServerSupabaseClient } from "@/lib/supabase/server"; @@ -195,6 +197,32 @@ export async function POST(request: Request) { }, persistence, model: defaultModel, + billing: { + reserve: async ({ userId: billingUserId, requestId, modelId }) => { + const pricing = await resolveFeaturePricing(admin, "report.full", modelId); + return authorizeUsage(admin, { + userId: billingUserId, requestId, featureKey: "report.full", + requestedModelId: modelId, creditCost: pricing.credit_cost, + }); + }, + complete: async ({ userId: billingUserId, requestId, usage }) => { + if (!defaultModel) return false; + const costMicrousd = Math.round(( + usage.inputTokens * (defaultModel.inputCostMicrousdPerMillion ?? 0) + + usage.outputTokens * (defaultModel.outputCostMicrousdPerMillion ?? 0) + ) / 1_000_000); + const settled = await completeUsage(admin, billingUserId, requestId, { + eventKey: "report.full", actualModelId: usage.actualModelId, + modelConfigVersion: usage.modelConfigVersion, inputTokens: usage.inputTokens, + outputTokens: usage.outputTokens, costMicrousd, durationMs: usage.durationMs, + }); + return settled.success; + }, + release: async ({ userId: billingUserId, requestId, reason }) => { + const released = await releaseUsage(admin, billingUserId, requestId, reason); + return released.success; + }, + }, runWorkflow: (input) => runConsultationWorkflow(input), createAgent: (model) => createPersonalReportAgent(model as Parameters[0]), skillSnapshot: resolveSkillSnapshot(), diff --git a/frontend/src/lib/personal-report-billing.ts b/frontend/src/lib/personal-report-billing.ts new file mode 100644 index 00000000..8541ab19 --- /dev/null +++ b/frontend/src/lib/personal-report-billing.ts @@ -0,0 +1,27 @@ +import type { UsageAuthorization } from "./consultation-billing"; + +export type ReportBillingUsage = Readonly<{ + actualModelId: string; + modelConfigVersion?: number; + inputTokens: number; + outputTokens: number; + durationMs: number; +}>; + +export type ReportBillingPort = Readonly<{ + reserve(input: Readonly<{ + userId: string; + requestId: string; + modelId: string; + }>): Promise; + complete(input: Readonly<{ + userId: string; + requestId: string; + usage: ReportBillingUsage; + }>): Promise; + release(input: Readonly<{ + userId: string; + requestId: string; + reason: string; + }>): Promise; +}>; diff --git a/frontend/src/lib/personal-report-codes.ts b/frontend/src/lib/personal-report-codes.ts index 790d4712..98fe550d 100644 --- a/frontend/src/lib/personal-report-codes.ts +++ b/frontend/src/lib/personal-report-codes.ts @@ -20,4 +20,7 @@ export const REPORT_STABLE_CODES = { invalidRequest: "invalid_request", resourceForbidden: "report_resource_forbidden", generationFailed: "report_generation_failed", + billingUnavailable: "report_billing_unavailable", + insufficientCredits: "insufficient_credits", + fairUseLimited: "fair_use_limited", } as const; diff --git a/frontend/src/lib/personal-report-generation.ts b/frontend/src/lib/personal-report-generation.ts index fffcaf93..78fe2c4d 100644 --- a/frontend/src/lib/personal-report-generation.ts +++ b/frontend/src/lib/personal-report-generation.ts @@ -2132,7 +2132,7 @@ export type GeneratePersonalReportDeps = GeneratePersonalReportBaseDeps & Readon export type ReportSchemaInnerReason = string; export type GeneratePersonalReportResult = Readonly< - | { status: "ready"; document: ReportDocumentV2; evidenceHash: string } + | { status: "ready"; document: ReportDocumentV2; evidenceHash: string; usage?: Readonly<{ inputTokens: number; outputTokens: number; actualModelId?: string; modelConfigVersion?: number }> } | { status: "failed"; failureCode: "report_schema_invalid"; innerReason: ReportSchemaInnerReason } | { status: "failed"; failureCode: "report_guard_rejected" } >; @@ -2316,7 +2316,7 @@ async function generateSectionedPersonalReport( if (!guarded.ok) return { status: "failed", failureCode: "report_guard_rejected" }; const parsed = safeParseServerReportDocument(guarded.document); if (!parsed.ok || parsed.document.schemaVersion !== "report_document.v2") return failSchema("final_parse_rejected"); - return { status: "ready", document: parsed.document, evidenceHash: computeEvidenceHash(parsed.document.evidenceAppendix) }; + return { status: "ready", document: parsed.document, evidenceHash: computeEvidenceHash(parsed.document.evidenceAppendix), usage: deps.agent.getUsage?.() }; } catch (error) { rethrowIfAborted(error, deps.signal); return failSchema(error instanceof ReportEvidenceInsufficientError ? "assemble_invalid" : classifyReportSchemaInnerReason(error)); @@ -2422,5 +2422,6 @@ export async function generatePersonalReport( status: "ready", document: parsed.document, evidenceHash: computeEvidenceHash(parsed.document.evidenceAppendix), + usage: deps.agent.getUsage?.(), }; } diff --git a/frontend/src/lib/personal-report-route-core.ts b/frontend/src/lib/personal-report-route-core.ts index d3af39b8..877fa6c9 100644 --- a/frontend/src/lib/personal-report-route-core.ts +++ b/frontend/src/lib/personal-report-route-core.ts @@ -24,6 +24,7 @@ import { type SkillSnapshot, } from "./personal-report-generation"; import { REPORT_STABLE_CODES } from "./personal-report-codes"; +import type { ReportBillingPort } from "./personal-report-billing"; import { checkSameOrigin } from "./personal-report-entitlement"; import type { PersonalReportJobRecord, PersonalReportJobService } from "./personal-report-job-service-core"; import type { @@ -190,7 +191,8 @@ export type ReportCreateCoreDeps = Readonly<{ countCreatedToday: () => Promise; }>; persistence: ReportServicePort; - model: Readonly<{ id: string }> | null; + model: Readonly<{ id: string; configVersion?: number; inputCostMicrousdPerMillion?: number; outputCostMicrousdPerMillion?: number }> | null; + billing?: ReportBillingPort; runWorkflow: (input: ConsultationInput) => Promise; createAgent: (model: Readonly<{ id: string }>) => ReportAgentPort; skillSnapshot: SkillSnapshot; @@ -311,6 +313,39 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise { + if (!billingReserved || !deps.billing) return; + billingReserved = false; + try { await deps.billing.release({ userId, requestId: payload.requestId, reason }); } catch { /* release is best effort; the RPC is idempotent */ } + }; const createInput: CreateGeneratingInput = { userId, @@ -327,14 +362,22 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise => { // Real workflow evidence (main chain), never mock/example/random data. - if (!deps.model) { - await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.modelUnavailable); - return { - status: 502, - body: { error: "报告模型暂不可用", code: REPORT_STABLE_CODES.modelUnavailable }, - }; - } - const workflowInputs: ConsultationInput[] = payload.themes.map((theme) => ({ year: birthDate.year, month: birthDate.month, @@ -380,6 +422,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise) => Promise | void; }>; @@ -79,6 +81,7 @@ export type PersonalReportWorkerDeps = Readonly<{ context: PersonalReportWorkerGenerationContext, ) => Promise; sectionService?: PersonalReportSectionService; + billing?: ReportBillingPort; leaseSeconds?: number; heartbeatIntervalMs?: number; recoveryLimit?: number; @@ -212,6 +215,10 @@ export function createPersonalReportWorker(deps: PersonalReportWorkerDeps) { try { if (terminal) { + if (deps.billing) { + await deps.billing.release({ userId: job.userId, requestId: job.requestId, reason: code }); + } + if (!report || report.status !== "generating") { await deps.jobs.markFailed({ jobId: job.id, @@ -280,6 +287,7 @@ export function createPersonalReportWorker(deps: PersonalReportWorkerDeps) { timerUnref(heartbeatTimer); let report: PersonalReportRecord | null = null; + const generationStartedAt = now().getTime(); try { report = await deps.reports.getByUserAndRequestId(job.userId, job.requestId); if (!report) { @@ -359,6 +367,20 @@ export function createPersonalReportWorker(deps: PersonalReportWorkerDeps) { report, document: generated.document, }); + if (deps.billing) { + const settled = await deps.billing.complete({ + userId: job.userId, + requestId: job.requestId, + usage: { + actualModelId: generated.usage?.actualModelId ?? "unknown", + modelConfigVersion: generated.usage?.modelConfigVersion, + inputTokens: generated.usage?.inputTokens ?? 0, + outputTokens: generated.usage?.outputTokens ?? 0, + durationMs: Math.max(0, now().getTime() - generationStartedAt), + }, + }); + if (!settled) throw new PersonalReportWorkerError("report_generation_failed", false, "report billing settlement failed"); + } return "ready"; } catch (error) { if (report?.status === "ready") { diff --git a/frontend/src/lib/personal-report-worker.ts b/frontend/src/lib/personal-report-worker.ts index d23c9630..e566182c 100644 --- a/frontend/src/lib/personal-report-worker.ts +++ b/frontend/src/lib/personal-report-worker.ts @@ -19,6 +19,7 @@ import { createSupabasePersonalReportService } from "@/lib/personal-report-servi import { createPersonalReportSectionService } from "@/lib/personal-report-section-service-core"; import { resolveReportBirthClock } from "@/lib/personal-report-route-core"; import { createAdminSupabaseClient } from "@/lib/supabase/admin"; +import { completeUsage, releaseUsage } from "@/lib/consultation-billing"; import type { ConsultationInput } from "@/mastra/consultation-workflow"; const PROFILE_COLUMNS = [ @@ -159,7 +160,7 @@ async function generateProductionReport(context: PersonalReportWorkerGenerationC throw new PersonalReportWorkerError("calculation_unavailable", false); } - return generatePersonalReport({ + const result = await generatePersonalReport({ reportId: context.report.id, bundle, depth: context.report.depth, @@ -170,6 +171,16 @@ async function generateProductionReport(context: PersonalReportWorkerGenerationC sectionService: context.sectionService, onProgress: context.onProgress, }); + if (result.status !== "ready") return result; + return { + ...result, + usage: { + inputTokens: result.usage?.inputTokens ?? 0, + outputTokens: result.usage?.outputTokens ?? 0, + actualModelId: model.id, + modelConfigVersion: model.configVersion, + }, + }; } function createProductionWorker(workerId: string) { @@ -188,6 +199,24 @@ function createProductionWorker(workerId: string) { jobs: createSupabasePersonalReportJobService(admin), reports: createSupabasePersonalReportService(admin), sectionService: createPersonalReportSectionService(admin as never), + billing: { + complete: async ({ userId, requestId, usage }) => { + const model = (await loadLanguageModelCatalog()).models.find((entry) => entry.id === usage.actualModelId); + if (!model) return false; + const costMicrousd = Math.round(( + usage.inputTokens * (model.inputCostMicrousdPerMillion ?? 0) + + usage.outputTokens * (model.outputCostMicrousdPerMillion ?? 0) + ) / 1_000_000); + const settled = await completeUsage(admin, userId, requestId, { + eventKey: "report.full", actualModelId: usage.actualModelId, + modelConfigVersion: usage.modelConfigVersion, inputTokens: usage.inputTokens, + outputTokens: usage.outputTokens, costMicrousd, durationMs: usage.durationMs, + }); + return settled.success; + }, + release: async ({ userId, requestId, reason }) => (await releaseUsage(admin, userId, requestId, reason)).success, + reserve: async () => { throw new Error("report worker does not reserve usage"); }, + }, loadProfile: async (userId) => { const { data, error } = await backend .from("profiles") diff --git a/frontend/src/mastra/personal-report.ts b/frontend/src/mastra/personal-report.ts index 315694d9..65cb60a0 100644 --- a/frontend/src/mastra/personal-report.ts +++ b/frontend/src/mastra/personal-report.ts @@ -234,8 +234,11 @@ export type ReportAgentSummaryOptions = Readonly<{ maxOutputTokens?: number; }>; +export type ReportAgentUsage = Readonly<{ inputTokens: number; outputTokens: number }>; + export type ReportAgentPort = Readonly<{ modelId: string; + getUsage?: () => ReportAgentUsage; generate( bundle: ReportEvidenceBundleV2, plan: PersonalReportSectionPlan, @@ -284,6 +287,15 @@ export function createPersonalReportAgent(model: ResolvedLanguageModel): ReportA instructions: personalReportInstructions, }); + let usageTotals: ReportAgentUsage = { inputTokens: 0, outputTokens: 0 }; + const recordUsage = async (usage: unknown) => { + const tokens = await readUsage(usage); + usageTotals = { + inputTokens: usageTotals.inputTokens + Math.max(0, tokens.inputTokens ?? 0), + outputTokens: usageTotals.outputTokens + Math.max(0, tokens.outputTokens ?? 0), + }; + }; + const runStructured = async (input: Readonly<{ prompt: string; schema: z.ZodType; @@ -324,9 +336,11 @@ export function createPersonalReportAgent(model: ResolvedLanguageModel): ReportA attemptReturned = true; const accepted = accept(first); if (accepted.ok) { + await recordUsage(first.usage); await logTelemetry(model.id, startedAt, false, "resolved", first.usage, first.finishReason); return accepted.data; } + await recordUsage(first.usage); await logTelemetry(model.id, startedAt, false, "failed", first.usage, first.finishReason); repairAttempted = true; attemptReturned = false; @@ -334,9 +348,11 @@ export function createPersonalReportAgent(model: ResolvedLanguageModel): ReportA attemptReturned = true; const repairedAccepted = accept(repaired); if (repairedAccepted.ok) { + await recordUsage(repaired.usage); await logTelemetry(model.id, startedAt, true, "resolved", repaired.usage, repaired.finishReason); return repairedAccepted.data; } + await recordUsage(repaired.usage); await logTelemetry(model.id, startedAt, true, "failed", repaired.usage, repaired.finishReason); if (repairedAccepted.error) throw repairedAccepted.error; throw new PersonalReportAgentOutputError(); @@ -351,6 +367,7 @@ export function createPersonalReportAgent(model: ResolvedLanguageModel): ReportA return { modelId: model.id, + getUsage: () => usageTotals, generate: (bundle, plan, options) => { const signal = options?.signal; return runStructured({ diff --git a/frontend/tests/personal-report-api.test.ts b/frontend/tests/personal-report-api.test.ts index 5be9b5ca..cd97a3ce 100644 --- a/frontend/tests/personal-report-api.test.ts +++ b/frontend/tests/personal-report-api.test.ts @@ -22,6 +22,8 @@ import type { } from "../src/lib/personal-report-service-core.ts"; import type { ReportEvidenceBundleV2 } from "../src/lib/report-evidence-bundle-v2.ts"; import type { ReportAgentPort, PersonalReportAgentOutput } from "../src/mastra/personal-report.ts"; +import type { ReportBillingPort } from "../src/lib/personal-report-billing.ts"; +import type { UsageAuthorization } from "../src/lib/consultation-billing.ts"; const createRoute = readFileSync( new URL("../src/app/api/reports/route.ts", import.meta.url), @@ -273,6 +275,31 @@ class MemoryPersistence implements ReportServicePort { } } +function billingHarness(overrides: Partial = {}): { + billing: ReportBillingPort; + readonly reservations: number; + readonly completions: number; + releases: string[]; +} { + let reservations = 0; + let completions = 0; + const releases: string[] = []; + const authorization: UsageAuthorization = { + success: true, reservation_id: null, source: "credits", credits: 0, + subscription_id: null, reason: null, retry_after_seconds: null, ...overrides, + }; + return { + billing: { + async reserve() { reservations += 1; return authorization; }, + async complete() { completions += 1; return true; }, + async release({ reason }) { releases.push(reason); return true; }, + }, + get reservations() { return reservations; }, + get completions() { return completions; }, + releases, + }; +} + function baseDeps(overrides: Partial = {}): ReportCreateCoreDeps { const persistence = new MemoryPersistence(); return { @@ -309,6 +336,60 @@ function baseDeps(overrides: Partial = {}): ReportCreateCo // Executable route core behavior (no network, no model) // --------------------------------------------------------------------------- +test("core create: billing denial returns 402 before persistence", async () => { + const persistence = new MemoryPersistence(); + const billing = billingHarness({ success: false, reason: "insufficient_credits" }); + const response = await resolveReportCreate(baseDeps({ persistence, billing: billing.billing })); + assert.equal(response.status, 402); + assert.equal(response.body.code, "insufficient_credits"); + assert.equal(persistence.rows.size, 0); + assert.equal(billing.reservations, 1); +}); + +test("core create: non-credit billing denial returns 503", async () => { + const billing = billingHarness({ success: false, reason: "fair_use_minute" }); + const response = await resolveReportCreate(baseDeps({ billing: billing.billing })); + assert.equal(response.status, 503); + assert.equal(response.body.code, "fair_use_limited"); +}); + +test("core create: inline success completes once and replay does not complete twice", async () => { + const persistence = new MemoryPersistence(); + const billing = billingHarness(); + const deps = baseDeps({ persistence, billing: billing.billing }); + const first = await resolveReportCreate(deps); + const second = await resolveReportCreate(deps); + assert.equal(first.status, 201); + assert.equal(second.status, 200); + // Replay is resolved before billing, so the idempotency key reserves once. + assert.equal(billing.reservations, 1); + assert.equal(billing.completions, 1); + assert.deepEqual(billing.releases, []); +}); + +test("core create: inline generation failure releases the reservation", async () => { + const billing = billingHarness(); + const response = await resolveReportCreate(baseDeps({ + billing: billing.billing, + runWorkflow: async () => { throw new Error("engine down"); }, + })); + assert.equal(response.status, 502); + assert.deepEqual(billing.releases, ["calculation_unavailable"]); + assert.equal(billing.completions, 0); +}); + +test("core create: enqueue failure releases the reservation", async () => { + const billing = billingHarness(); + const response = await resolveReportCreate(baseDeps({ + billing: billing.billing, + jobs: { + async enqueue() { throw new Error("queue down"); }, + }, + })); + assert.equal(response.status, 503); + assert.deepEqual(billing.releases, ["report_generation_failed"]); +}); + test("core create: 401 when not logged in", async () => { const response = await resolveReportCreate(baseDeps({ userId: null })); assert.equal(response.status, 401);