feat(hub): usage fact ledger for run-scoped cost attribution (ADR-0026)

Replace the single-scalar cost model on AgentRun with an append-only
UsageFact ledger. One AgentRun owns zero or more UsageFact rows; each
records one billable consumption event (model completion, external
capability, or tool proxy) with its own provider/model/tokens/quantity/
cost. AgentRun.costUsd/inputTokens/outputTokens become a derived rollup
cache.

This unblocks external capabilities (PDF->MD bundle, audio/video->text)
that bill in non-token units (pages, seconds) through a different
provider than the main agent loop, without per-capability schema changes
or nested AgentRuns (which would pollute lock/admission/session
semantics).

Contract:
- spec/Spec/System/Agent/Usage.lean pins UsageFact, UsageFactKind,
  CostSource and three invariants: append-only; belongs to one run,
  never holds a lock; missing cost != zero (ADR-0022).
- ADR-0026 records the decision, the rejected nested-Run alternative,
  the rollup cache strategy, and the deferred capability registry /
  pricebook / post-hoc correction flows.

Schema:
- UsageFact model with indexes on (runId, occurredAt), (runId, kind),
  (provider, model, occurredAt), (capabilityId, occurredAt).
- Migration backfills one synthetic model_completion fact per existing
  run with recorded cost/tokens (correlationId = runId marks backfill);
  truly unrecorded runs stay runsWithoutCost per ADR-0022.

Write path (trigger finish):
- Write the UsageFact first, then mirror it onto AgentRun as two
  separate statements (not one transaction). The fact is the truth so it
  goes first; the cache is derived so it goes second. A crash between
  them leaves the cache stale but the usage service re-reads facts
  directly, so this is recoverable; the reverse order would lose the
  truth. Separate statements also avoid an AgentRun row lock held across
  the insert's FK ShareLock, which deadlocked concurrent workspace
  teardown under the Organization->Project->AgentRun->UsageFact cascade.

Read paths:
- org/usage.ts aggregates from UsageFact, ignoring the AgentRun cache.
- slash /usage buckets by (fact.provider, fact.model); a run with a
  main loop + an external call lands in two buckets.
- session detail exposes usageFacts[] + costSource for future per-run
  cost-breakdown UI.

Tests:
- usage.test.ts: 6 integration tests pin fact aggregation, missing-cost-
  !=-zero, multi-fact-per-run, empty-run, project-level, empty-org.
- trigger.test.ts: existing /usage assertion ($0.0023,
  openrouter / mock-model) passes on the fact path.
- feishu-reactions mock prisma gains usageFact.create.
This commit is contained in:
2026-07-18 14:44:23 +08:00
parent 97f7972cc5
commit f3b087371a
11 changed files with 760 additions and 46 deletions
+54 -24
View File
@@ -92,7 +92,18 @@ export function createSlashCommandRegistry(
finishedAt: { not: null },
},
orderBy: { finishedAt: "asc" },
select: { model: true, provider: true, inputTokens: true, outputTokens: true, costUsd: true },
select: {
usageFacts: {
select: {
provider: true,
model: true,
inputTokens: true,
outputTokens: true,
costUsd: true,
},
orderBy: { occurredAt: "asc" },
},
},
});
await sendText(rt, chatId, formatUsageReport(runs, scope), sendOptions);
},
@@ -118,18 +129,24 @@ async function currentRoleSessionIds(
return sessions.map((session) => session.id);
}
interface UsageRun {
readonly model: string;
readonly provider: string;
readonly inputTokens: number | null;
readonly outputTokens: number | null;
readonly costUsd: unknown;
interface UsageRunWithFacts {
readonly usageFacts: readonly {
readonly provider: string;
readonly model: string | null;
readonly inputTokens: number | null;
readonly outputTokens: number | null;
readonly costUsd: unknown;
}[];
}
function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "project"): string {
function formatUsageReport(runs: readonly UsageRunWithFacts[], scope: "current" | "project"): string {
if (runs.length === 0) return `${scope === "current" ? "当前角色会话" : "当前项目"}还没有已结束的 Agent run。`;
// Bucket by (fact.provider, fact.model) — ADR-0026: external capabilities
// carry their own provider/model and contribute their own meter, so a run
// with a main loop + an external call lands in two buckets. A run with no
// cost-bearing fact is "unrecorded" (ADR-0022: missing cost ≠ zero).
const buckets = new Map<string, {
provider: string; model: string; runs: number; inputTokens: number; outputTokens: number; costUsd: number;
provider: string; model: string; facts: number; inputTokens: number; outputTokens: number; costUsd: number;
}>();
let recordedRuns = 0;
let unrecordedRuns = 0;
@@ -137,21 +154,34 @@ function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "projec
let totalOutputTokens = 0;
let totalCostUsd = 0;
for (const run of runs) {
const costUsd = decimalToNumberOrNull(run.costUsd);
if (costUsd === null) { unrecordedRuns++; continue; }
const facts = run.usageFacts;
if (facts.length === 0 || !facts.some((f) => f.costUsd !== null && f.costUsd !== undefined)) {
unrecordedRuns++;
// Tokens from unrecorded runs still count toward totals (matches pre-0026).
for (const f of facts) {
totalInputTokens += f.inputTokens ?? 0;
totalOutputTokens += f.outputTokens ?? 0;
}
continue;
}
recordedRuns++;
const key = `${run.provider}\u0000${run.model}`;
const bucket = buckets.get(key) ?? {
provider: run.provider, model: run.model, runs: 0, inputTokens: 0, outputTokens: 0, costUsd: 0,
};
bucket.runs++;
bucket.inputTokens += run.inputTokens ?? 0;
bucket.outputTokens += run.outputTokens ?? 0;
bucket.costUsd += costUsd;
buckets.set(key, bucket);
totalInputTokens += run.inputTokens ?? 0;
totalOutputTokens += run.outputTokens ?? 0;
totalCostUsd += costUsd;
for (const f of facts) {
const costUsd = decimalToNumberOrNull(f.costUsd);
if (costUsd === null) continue;
const model = f.model ?? "(unknown model)";
const key = `${f.provider}\u0000${model}`;
const bucket = buckets.get(key) ?? {
provider: f.provider, model, facts: 0, inputTokens: 0, outputTokens: 0, costUsd: 0,
};
bucket.facts += 1;
bucket.inputTokens += f.inputTokens ?? 0;
bucket.outputTokens += f.outputTokens ?? 0;
bucket.costUsd += costUsd;
buckets.set(key, bucket);
totalInputTokens += f.inputTokens ?? 0;
totalOutputTokens += f.outputTokens ?? 0;
totalCostUsd += costUsd;
}
}
const lines = [
`${scope === "current" ? "当前角色会话" : "当前项目"}用量`,
@@ -160,7 +190,7 @@ function formatUsageReport(runs: readonly UsageRun[], scope: "current" | "projec
];
if (unrecordedRuns > 0) lines.push(`另有 ${formatInteger(unrecordedRuns)} 个 run 未记录成本。`);
for (const bucket of buckets.values()) {
lines.push(`- ${bucket.provider} / ${bucket.model}: ${formatInteger(bucket.runs)} runs, ${formatUsd(bucket.costUsd)}`);
lines.push(`- ${bucket.provider} / ${bucket.model}: ${formatInteger(bucket.facts)} , ${formatUsd(bucket.costUsd)}`);
}
return lines.join("\n");
}
+29 -2
View File
@@ -604,6 +604,32 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler {
data: { metadata: mergeSessionMetadata(session.metadata, metadataPatch) },
});
}
// ADR-0026: the UsageFact ledger is the truth; AgentRun.costUsd /
// inputTokens / outputTokens are a derived rollup cache for existing
// readers. We write the fact first, then mirror it onto the run.
// They are separate statements (not one transaction) so the AgentRun
// row lock is held for the shortest possible window and concurrent
// workspace teardown cannot deadlock on the FK ShareLock. A crash
// between the two leaves the cache stale, but the cache is derived
// (the usage service re-reads facts), so staleness is recoverable;
// the reverse order would lose the truth and is unrecoverable.
const finishedAt = new Date();
const reportedCost = result.costUsd;
const costSource = reportedCost !== undefined ? "provider_reported" : "unknown";
await deps.prisma.usageFact.create({
data: {
runId: run.id,
occurredAt: finishedAt,
kind: "model_completion",
provider: providerId,
model,
inputTokens: result.usage.inputTokens,
outputTokens: result.usage.outputTokens,
costUsd: reportedCost ?? null,
costSource,
metadata: {},
},
});
await deps.prisma.agentRun.update({
where: { id: run.id },
data: {
@@ -618,9 +644,10 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler {
: "FAILED",
inputTokens: result.usage.inputTokens,
outputTokens: result.usage.outputTokens,
...(result.costUsd !== undefined ? { costUsd: result.costUsd, costSource: "provider_reported" } : {}),
costUsd: reportedCost ?? null,
costSource,
error: result.error ?? null,
finishedAt: new Date(),
finishedAt,
},
});
await writeAudit(deps.prisma, {