feat(billing): charge personal reports with durable settlement
This commit is contained in:
@@ -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<typeof createPersonalReportAgent>[0]),
|
||||
skillSnapshot: resolveSkillSnapshot(),
|
||||
|
||||
@@ -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<UsageAuthorization>;
|
||||
complete(input: Readonly<{
|
||||
userId: string;
|
||||
requestId: string;
|
||||
usage: ReportBillingUsage;
|
||||
}>): Promise<boolean>;
|
||||
release(input: Readonly<{
|
||||
userId: string;
|
||||
requestId: string;
|
||||
reason: string;
|
||||
}>): Promise<boolean>;
|
||||
}>;
|
||||
@@ -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;
|
||||
|
||||
@@ -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?.(),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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<number>;
|
||||
}>;
|
||||
persistence: ReportServicePort;
|
||||
model: Readonly<{ id: string }> | null;
|
||||
model: Readonly<{ id: string; configVersion?: number; inputCostMicrousdPerMillion?: number; outputCostMicrousdPerMillion?: number }> | null;
|
||||
billing?: ReportBillingPort;
|
||||
runWorkflow: (input: ConsultationInput) => Promise<unknown>;
|
||||
createAgent: (model: Readonly<{ id: string }>) => ReportAgentPort;
|
||||
skillSnapshot: SkillSnapshot;
|
||||
@@ -311,6 +313,39 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
body: { error: "今日报告生成次数已达上限", code: REPORT_STABLE_CODES.rateLimited },
|
||||
};
|
||||
}
|
||||
const selectedModel = deps.model;
|
||||
|
||||
let billingReserved = false;
|
||||
if (deps.billing && selectedModel) {
|
||||
let authorization;
|
||||
try {
|
||||
authorization = await deps.billing.reserve({ userId, requestId: payload.requestId, modelId: selectedModel.id });
|
||||
} catch {
|
||||
return {
|
||||
status: 503,
|
||||
body: { error: "暂时无法确认报告点数", code: REPORT_STABLE_CODES.billingUnavailable },
|
||||
};
|
||||
}
|
||||
if (!authorization.success) {
|
||||
const insufficient = authorization.reason === REPORT_STABLE_CODES.insufficientCredits;
|
||||
return {
|
||||
status: insufficient ? 402 : 503,
|
||||
body: {
|
||||
error: insufficient ? "报告点数不足" : "暂时无法扣除报告点数",
|
||||
message: authorization.reason ?? "请稍后重试。",
|
||||
code: insufficient ? REPORT_STABLE_CODES.insufficientCredits : REPORT_STABLE_CODES.fairUseLimited,
|
||||
...(authorization.retry_after_seconds === null ? {} : { retryAfterSeconds: authorization.retry_after_seconds }),
|
||||
},
|
||||
};
|
||||
}
|
||||
billingReserved = true;
|
||||
}
|
||||
|
||||
const releaseBilling = async (reason: string) => {
|
||||
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<R
|
||||
skillSourceCommit: deps.skillSnapshot.sourceCommit,
|
||||
skillSnapshotSha256: deps.skillSnapshot.sha256,
|
||||
};
|
||||
const begun: CreateGeneratingResult = await deps.persistence.createGenerating(createInput);
|
||||
let begun: CreateGeneratingResult;
|
||||
try {
|
||||
begun = await deps.persistence.createGenerating(createInput);
|
||||
} catch (error) {
|
||||
await releaseBilling(REPORT_STABLE_CODES.generationFailed);
|
||||
throw error;
|
||||
}
|
||||
if (begun.kind === "generation_in_progress") {
|
||||
await releaseBilling(REPORT_STABLE_CODES.generationInProgress);
|
||||
return {
|
||||
status: 409,
|
||||
body: { error: "已有报告正在生成中", code: REPORT_STABLE_CODES.generationInProgress },
|
||||
};
|
||||
}
|
||||
if (begun.kind === "request_conflict") {
|
||||
await releaseBilling(REPORT_STABLE_CODES.requestConflict);
|
||||
return {
|
||||
status: 409,
|
||||
body: { error: "请求内容与已有记录不一致", code: REPORT_STABLE_CODES.requestConflict },
|
||||
@@ -344,17 +387,16 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
return replayOrConflict(begun.record, fingerprint);
|
||||
}
|
||||
const row = begun.record;
|
||||
if (!selectedModel) {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.modelUnavailable);
|
||||
return {
|
||||
status: 502,
|
||||
body: { error: "报告模型暂不可用", code: REPORT_STABLE_CODES.modelUnavailable },
|
||||
};
|
||||
}
|
||||
|
||||
const finishGeneration = async (): Promise<ReportRouteResponse> => {
|
||||
// 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<R
|
||||
}
|
||||
} catch {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.calculationUnavailable);
|
||||
await releaseBilling(REPORT_STABLE_CODES.calculationUnavailable);
|
||||
return {
|
||||
status: 502,
|
||||
body: { error: "排盘引擎暂不可用", code: REPORT_STABLE_CODES.calculationUnavailable },
|
||||
@@ -392,6 +435,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
});
|
||||
if (!hasUsableBaseChart) {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.calculationUnavailable);
|
||||
await releaseBilling(REPORT_STABLE_CODES.calculationUnavailable);
|
||||
return {
|
||||
status: 502,
|
||||
body: { error: "排盘引擎未返回可用星盘", code: REPORT_STABLE_CODES.calculationUnavailable },
|
||||
@@ -417,6 +461,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
// gap is represented inside the Bundle as a blocked section instead.
|
||||
if (error instanceof Error && error.name === "ReportEvidenceInsufficientError") {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.calculationUnavailable);
|
||||
await releaseBilling(REPORT_STABLE_CODES.calculationUnavailable);
|
||||
return {
|
||||
status: 422,
|
||||
body: { error: "排盘证据不足以生成诚实报告", code: REPORT_STABLE_CODES.calculationUnavailable },
|
||||
@@ -429,12 +474,13 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
reportId: row.id,
|
||||
bundle,
|
||||
depth: payload.depth,
|
||||
agent: deps.createAgent(deps.model),
|
||||
agent: deps.createAgent(selectedModel),
|
||||
now: deps.now,
|
||||
});
|
||||
|
||||
if (result.status === "failed") {
|
||||
await deps.persistence.markFailed(userId, row.id, result.failureCode);
|
||||
await releaseBilling(result.failureCode);
|
||||
return {
|
||||
status: 422,
|
||||
body: {
|
||||
@@ -448,12 +494,27 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
|
||||
try {
|
||||
const readyRow = await deps.persistence.completeReady(userId, row.id, result.document);
|
||||
if (deps.billing) {
|
||||
const settled = await deps.billing.complete({
|
||||
userId, requestId: payload.requestId,
|
||||
usage: {
|
||||
actualModelId: selectedModel.id,
|
||||
modelConfigVersion: selectedModel.configVersion,
|
||||
inputTokens: result.usage?.inputTokens ?? 0,
|
||||
outputTokens: result.usage?.outputTokens ?? 0,
|
||||
durationMs: 0,
|
||||
},
|
||||
});
|
||||
if (!settled) throw new Error("report billing settlement failed");
|
||||
billingReserved = false;
|
||||
}
|
||||
return {
|
||||
status: 201,
|
||||
body: { report: reportView(readyRow), reportDocument: readyRow.reportDocument },
|
||||
};
|
||||
} catch {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.schemaInvalid);
|
||||
await releaseBilling(REPORT_STABLE_CODES.schemaInvalid);
|
||||
return {
|
||||
status: 422,
|
||||
body: { error: "报告未通过合同校验", code: REPORT_STABLE_CODES.schemaInvalid },
|
||||
@@ -470,6 +531,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
});
|
||||
if (enqueued.kind === "request_conflict") {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.schemaInvalid);
|
||||
await releaseBilling(REPORT_STABLE_CODES.requestConflict);
|
||||
return {
|
||||
status: 409,
|
||||
body: { error: "请求内容与已有任务不一致", code: REPORT_STABLE_CODES.requestConflict },
|
||||
@@ -477,6 +539,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
}
|
||||
if (enqueued.kind === "active_job_exists") {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.generationInProgress);
|
||||
await releaseBilling(REPORT_STABLE_CODES.generationInProgress);
|
||||
return {
|
||||
status: 409,
|
||||
body: { error: "已有报告任务正在处理中", code: REPORT_STABLE_CODES.generationInProgress },
|
||||
@@ -488,6 +551,7 @@ export async function resolveReportCreate(deps: ReportCreateCoreDeps): Promise<R
|
||||
};
|
||||
} catch {
|
||||
await deps.persistence.markFailed(userId, row.id, REPORT_STABLE_CODES.schemaInvalid);
|
||||
await releaseBilling(REPORT_STABLE_CODES.generationFailed);
|
||||
return {
|
||||
status: 503,
|
||||
body: { error: "报告任务暂时无法入队", code: REPORT_STABLE_CODES.generationFailed },
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
} from "./personal-report-service-core.ts";
|
||||
import type { GeneratePersonalReportResult } from "./personal-report-generation.ts";
|
||||
import type { PersonalReportSectionService } from "./personal-report-section-service-core.ts";
|
||||
import type { ReportBillingPort } from "./personal-report-billing.ts";
|
||||
|
||||
/**
|
||||
* Durable personal-report worker orchestration.
|
||||
@@ -49,6 +50,7 @@ export type PersonalReportWorkerGenerationContext = Readonly<{
|
||||
profile: unknown;
|
||||
signal: AbortSignal;
|
||||
sectionService?: PersonalReportSectionService;
|
||||
billing?: ReportBillingPort;
|
||||
onProgress?: (progress: Readonly<{ phase: string; completed: number; total: number }>) => Promise<void> | void;
|
||||
}>;
|
||||
|
||||
@@ -79,6 +81,7 @@ export type PersonalReportWorkerDeps = Readonly<{
|
||||
context: PersonalReportWorkerGenerationContext,
|
||||
) => Promise<GeneratePersonalReportResult>;
|
||||
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") {
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 <T>(input: Readonly<{
|
||||
prompt: string;
|
||||
schema: z.ZodType<T>;
|
||||
@@ -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({
|
||||
|
||||
@@ -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<UsageAuthorization> = {}): {
|
||||
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> = {}): ReportCreateCoreDeps {
|
||||
const persistence = new MemoryPersistence();
|
||||
return {
|
||||
@@ -309,6 +336,60 @@ function baseDeps(overrides: Partial<ReportCreateCoreDeps> = {}): 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);
|
||||
|
||||
Reference in New Issue
Block a user