Files
EvoScientist-WebUI/src/lib/server/usageStore.ts
T
m4 c8e8305107 fix(webui): read usage deployment id from EVOSCIENTIST_USAGE_DEPLOYMENT_ID
Matches the backend split: usage attribution no longer shares
EVOSCIENTIST_DEPLOYMENT_ID with workspace scope partitioning.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-07-30 22:11:31 +08:00

761 lines
27 KiB
TypeScript

import "server-only";
import { createHash, randomUUID } from "node:crypto";
import fs from "node:fs";
import path from "node:path";
import Database from "better-sqlite3";
import {
canonicalJson,
type UsageEventV1,
type UsageHeartbeatV1,
} from "@/lib/usageTypes";
import {
USAGE_HEARTBEAT_TTL_SECONDS,
USAGE_INBOX_RETENTION_DAYS,
usageDataDir,
} from "./usageConfig";
type IngestResult = {
status: "accepted" | "duplicate" | "conflict";
event_id: string;
reason_code: string;
};
type SqlValue = string | bigint | null;
export interface UsageSummaryResult {
deployment_id: string | null;
workspace_id: string | null;
input_tokens: string;
output_tokens: string;
total_tokens: string;
confirmed_call_count: string;
unknown_call_count: string;
by_scope: Array<Record<string, string>>;
by_model: Array<Record<string, string>>;
by_provider: Array<Record<string, string>>;
}
export interface UsageStatusResult {
availability: string;
reason_code: string;
deployment_id?: string | null;
workspace_id?: string | null;
[key: string]: string | null | undefined;
}
export interface ProviderAggregateObservation {
aggregate_id: string;
deployment_id: string;
provider_profile_id: string;
upstream_model_id: string | null;
window_start: string;
window_end: string;
input_tokens: string;
output_tokens: string;
provider_total_tokens: string | null;
request_count: string;
source_name: string;
source_payload: unknown;
observed_at: string;
}
interface UsageStoreGlobal {
db?: Database.Database;
collectorInstanceId?: string;
lastMaintenanceAt?: number;
}
const globalUsage = globalThis as typeof globalThis & {
__evoscientistUsage?: UsageStoreGlobal;
};
globalUsage.__evoscientistUsage ??= {};
function durableIdentifier(filePath: string): string {
fs.mkdirSync(path.dirname(filePath), { recursive: true, mode: 0o700 });
try {
const fd = fs.openSync(filePath, "wx", 0o600);
const value = randomUUID();
try {
fs.writeFileSync(fd, `${value}\n`, "utf8");
fs.fsyncSync(fd);
} finally {
fs.closeSync(fd);
}
try {
const dir = fs.openSync(path.dirname(filePath), "r");
try {
fs.fsyncSync(dir);
} finally {
fs.closeSync(dir);
}
} catch {
/* Directory fsync is unavailable on some platforms. */
}
return value;
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error;
return fs.readFileSync(filePath, "utf8").trim();
}
}
export function collectorInstanceId(): string {
if (!globalUsage.__evoscientistUsage!.collectorInstanceId) {
globalUsage.__evoscientistUsage!.collectorInstanceId = durableIdentifier(
path.join(usageDataDir(), "webui", "collector-id")
);
}
return globalUsage.__evoscientistUsage!.collectorInstanceId!;
}
function migrate(db: Database.Database): void {
const version = Number(db.pragma("user_version", { simple: true }));
if (version > 3)
throw new Error(
`usage database version ${version} is newer than this WebUI supports`
);
if (version < 1)
db.exec(`
CREATE TABLE usage_sources (
deployment_id TEXT NOT NULL, workspace_id TEXT NOT NULL,
emitter_version TEXT NOT NULL, schema_version INTEGER NOT NULL,
sender_status TEXT NOT NULL CHECK(sender_status IN ('healthy','degraded')),
spool_pending INTEGER NOT NULL DEFAULT 0, spool_inflight INTEGER NOT NULL DEFAULT 0,
spool_quarantined INTEGER NOT NULL DEFAULT 0, spool_bytes INTEGER NOT NULL DEFAULT 0,
first_loss_at TEXT, tracking_degraded_reason TEXT, last_error_code TEXT,
last_seen_at TEXT NOT NULL, PRIMARY KEY(deployment_id, workspace_id)
);
CREATE TABLE usage_event_inbox (
event_id TEXT PRIMARY KEY, deployment_id TEXT NOT NULL, model_call_id TEXT NOT NULL,
source TEXT NOT NULL CHECK(source='callback_final'),
authority_class TEXT NOT NULL CHECK(authority_class='observed_final'),
revision INTEGER NOT NULL CHECK(revision=1), payload_hash TEXT NOT NULL,
payload_json TEXT NOT NULL, received_at TEXT NOT NULL
);
CREATE TABLE usage_event_conflicts (
conflict_id TEXT PRIMARY KEY, event_id TEXT NOT NULL, stored_payload_hash TEXT NOT NULL,
received_payload_hash TEXT NOT NULL, received_payload_json TEXT NOT NULL,
received_at TEXT NOT NULL, UNIQUE(event_id, received_payload_hash)
);
CREATE TABLE model_usage (
deployment_id TEXT NOT NULL, workspace_id TEXT NOT NULL, model_call_id TEXT NOT NULL,
source TEXT NOT NULL CHECK(source='callback_final'),
authority_class TEXT NOT NULL CHECK(authority_class='observed_final'),
revision INTEGER NOT NULL CHECK(revision=1), parent_run_id TEXT, provider_request_id TEXT,
thread_id TEXT, source_session_id TEXT, turn_id TEXT, workspace_dir TEXT,
scope TEXT NOT NULL DEFAULT 'unattributed', source_agent TEXT,
provider_profile_id TEXT NOT NULL, provider_revision TEXT, provider_adapter TEXT NOT NULL,
model_alias TEXT NOT NULL, upstream_model_id TEXT NOT NULL,
usage_status TEXT NOT NULL CHECK(usage_status IN ('confirmed','unknown')),
input_tokens INTEGER CHECK(input_tokens IS NULL OR input_tokens BETWEEN 0 AND 9007199254740991),
output_tokens INTEGER CHECK(output_tokens IS NULL OR output_tokens BETWEEN 0 AND 9007199254740991),
provider_total_tokens INTEGER CHECK(provider_total_tokens IS NULL OR provider_total_tokens BETWEEN 0 AND 9007199254740991),
input_details_json TEXT NOT NULL, output_details_json TEXT NOT NULL,
started_at TEXT, observed_at TEXT NOT NULL, completed_at TEXT NOT NULL, updated_at TEXT NOT NULL,
PRIMARY KEY(deployment_id, model_call_id),
CHECK((usage_status='confirmed' AND input_tokens IS NOT NULL AND output_tokens IS NOT NULL)
OR (usage_status='unknown' AND input_tokens IS NULL AND output_tokens IS NULL AND provider_total_tokens IS NULL)),
CHECK(usage_status='unknown' OR input_tokens + output_tokens <= 9007199254740991)
);
CREATE INDEX model_usage_thread_idx ON model_usage(deployment_id, thread_id, completed_at);
CREATE INDEX model_usage_source_session_idx ON model_usage(deployment_id, source_session_id, completed_at);
CREATE INDEX model_usage_turn_idx ON model_usage(deployment_id, turn_id, completed_at);
CREATE INDEX model_usage_provider_idx ON model_usage(deployment_id, provider_profile_id, upstream_model_id, completed_at);
PRAGMA user_version=1;
`);
if (version < 2)
db.exec(`
CREATE TABLE usage_event_idempotency (
event_id TEXT PRIMARY KEY,
payload_hash TEXT NOT NULL,
retained_at TEXT NOT NULL
);
INSERT OR IGNORE INTO usage_event_idempotency(event_id,payload_hash,retained_at)
SELECT event_id,payload_hash,received_at FROM usage_event_inbox;
PRAGMA user_version=2;
`);
if (version < 3)
db.exec(`
CREATE TABLE provider_usage_aggregates (
aggregate_id TEXT PRIMARY KEY,
deployment_id TEXT NOT NULL,
provider_profile_id TEXT NOT NULL,
upstream_model_id TEXT,
window_start TEXT NOT NULL,
window_end TEXT NOT NULL,
input_tokens TEXT NOT NULL,
output_tokens TEXT NOT NULL,
provider_total_tokens TEXT,
request_count TEXT NOT NULL,
source_name TEXT NOT NULL,
source_payload_json TEXT NOT NULL,
payload_hash TEXT NOT NULL,
observed_at TEXT NOT NULL,
projection_policy_version TEXT NOT NULL
CHECK(projection_policy_version='aggregate-report-v1')
);
CREATE INDEX provider_usage_aggregate_window_idx
ON provider_usage_aggregates(
deployment_id, provider_profile_id, window_start, window_end
);
PRAGMA user_version=3;
`);
}
export function usageDb(): Database.Database {
const state = globalUsage.__evoscientistUsage!;
if (state.db?.open) {
maybeRunUsageMaintenance(state.db);
return state.db;
}
const dbPath = path.join(usageDataDir(), "webui", "usage.db");
fs.mkdirSync(path.dirname(dbPath), { recursive: true, mode: 0o700 });
const db = new Database(dbPath);
db.pragma("journal_mode = WAL");
db.pragma("synchronous = FULL");
db.pragma("busy_timeout = 5000");
db.pragma("foreign_keys = ON");
db.defaultSafeIntegers(true);
db.aggregate("decimal_sum", {
start: BigInt(0),
step: (total: bigint, value: bigint | null) => total + (value ?? BigInt(0)),
result: (total: bigint) => total.toString(),
});
migrate(db);
state.db = db;
runUsageMaintenance(Date.now(), db);
return db;
}
export function closeUsageDb(): void {
const state = globalUsage.__evoscientistUsage!;
if (state.db?.open) state.db.close();
state.db = undefined;
state.lastMaintenanceAt = undefined;
}
function maybeRunUsageMaintenance(db: Database.Database): void {
const state = globalUsage.__evoscientistUsage!;
const now = Date.now();
if (!state.lastMaintenanceAt || now - state.lastMaintenanceAt >= 3_600_000) {
runUsageMaintenance(now, db);
}
}
export function runUsageMaintenance(
now = Date.now(),
database: Database.Database = usageDb()
): { conflicts_deleted: string; inbox_deleted: string } {
const inboxCutoff = new Date(
now - USAGE_INBOX_RETENTION_DAYS * 86_400_000
).toISOString();
const result = database.transaction(() => {
const conflicts = database
.prepare("DELETE FROM usage_event_conflicts WHERE received_at < ?")
.run(inboxCutoff);
const inbox = database
.prepare("DELETE FROM usage_event_inbox WHERE received_at < ?")
.run(inboxCutoff);
return {
conflicts_deleted: conflicts.changes.toString(),
inbox_deleted: inbox.changes.toString(),
};
})();
globalUsage.__evoscientistUsage!.lastMaintenanceAt = now;
return result;
}
export async function backupUsageDatabase(destination: string): Promise<void> {
if (!path.isAbsolute(destination)) {
throw new Error("usage backup destination must be absolute");
}
fs.mkdirSync(path.dirname(destination), { recursive: true, mode: 0o700 });
await usageDb().backup(destination);
}
function eventParams(event: UsageEventV1, now: string): SqlValue[] {
return [
event.deployment_id,
event.workspace_id,
event.model_call_id,
event.source,
event.authority_class,
BigInt(event.revision),
event.parent_run_id,
event.provider_request_id,
event.thread_id,
event.source_session_id,
event.turn_id,
event.workspace_dir,
event.scope,
event.source_agent,
event.provider_profile_id,
event.provider_revision,
event.provider_adapter,
event.model_alias,
event.upstream_model_id,
event.usage_status,
event.input_tokens === null ? null : BigInt(event.input_tokens),
event.output_tokens === null ? null : BigInt(event.output_tokens),
event.provider_total_tokens === null
? null
: BigInt(event.provider_total_tokens),
canonicalJson(event.input_token_details),
canonicalJson(event.output_token_details),
event.started_at,
event.observed_at,
event.completed_at,
now,
];
}
export function ingestUsageEvent(event: UsageEventV1): IngestResult {
const db = usageDb();
const payload = canonicalJson(event);
const hash = createHash("sha256").update(payload).digest("hex");
const now = new Date().toISOString();
return db.transaction((): IngestResult => {
const stored = db
.prepare(
"SELECT payload_hash FROM usage_event_idempotency WHERE event_id=?"
)
.get(event.event_id) as { payload_hash: string } | undefined;
if (stored) {
if (stored.payload_hash === hash)
return {
status: "duplicate",
event_id: event.event_id,
reason_code: "event_already_stored",
};
const conflictId = createHash("sha256")
.update(`${event.event_id}\0${hash}`)
.digest("hex");
db.prepare(
`INSERT OR IGNORE INTO usage_event_conflicts
(conflict_id,event_id,stored_payload_hash,received_payload_hash,received_payload_json,received_at)
VALUES(?,?,?,?,?,?)`
).run(
conflictId,
event.event_id,
stored.payload_hash,
hash,
payload,
now
);
return {
status: "conflict",
event_id: event.event_id,
reason_code: "event_payload_conflict",
};
}
db.prepare(
`INSERT INTO usage_event_inbox
(event_id,deployment_id,model_call_id,source,authority_class,revision,payload_hash,payload_json,received_at)
VALUES(?,?,?,?,?,?,?,?,?)`
).run(
event.event_id,
event.deployment_id,
event.model_call_id,
event.source,
event.authority_class,
BigInt(event.revision),
hash,
payload,
now
);
db.prepare(
`INSERT INTO usage_event_idempotency(event_id,payload_hash,retained_at)
VALUES(?,?,?)`
).run(event.event_id, hash, now);
db.prepare(
`INSERT INTO model_usage (
deployment_id,workspace_id,model_call_id,source,authority_class,revision,parent_run_id,
provider_request_id,thread_id,source_session_id,turn_id,workspace_dir,scope,source_agent,
provider_profile_id,provider_revision,provider_adapter,model_alias,upstream_model_id,
usage_status,input_tokens,output_tokens,provider_total_tokens,input_details_json,
output_details_json,started_at,observed_at,completed_at,updated_at)
VALUES(${Array.from({ length: 29 }, () => "?").join(",")})`
).run(...eventParams(event, now));
return {
status: "accepted",
event_id: event.event_id,
reason_code: "event_stored",
};
})();
}
export function recordHeartbeat(heartbeat: UsageHeartbeatV1): void {
const db = usageDb();
const receivedAt = new Date().toISOString();
db.prepare(
`INSERT INTO usage_sources (
deployment_id,workspace_id,emitter_version,schema_version,sender_status,spool_pending,
spool_inflight,spool_quarantined,spool_bytes,first_loss_at,tracking_degraded_reason,
last_error_code,last_seen_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)
ON CONFLICT(deployment_id,workspace_id) DO UPDATE SET
emitter_version=excluded.emitter_version, schema_version=excluded.schema_version,
sender_status=excluded.sender_status, spool_pending=excluded.spool_pending,
spool_inflight=excluded.spool_inflight, spool_quarantined=excluded.spool_quarantined,
spool_bytes=excluded.spool_bytes,
first_loss_at=CASE WHEN usage_sources.first_loss_at IS NULL THEN excluded.first_loss_at
WHEN excluded.first_loss_at IS NULL THEN usage_sources.first_loss_at
WHEN usage_sources.first_loss_at <= excluded.first_loss_at THEN usage_sources.first_loss_at ELSE excluded.first_loss_at END,
tracking_degraded_reason=COALESCE(usage_sources.tracking_degraded_reason,excluded.tracking_degraded_reason),
last_error_code=excluded.last_error_code,last_seen_at=excluded.last_seen_at`
).run(
heartbeat.deployment_id,
heartbeat.workspace_id,
heartbeat.emitter_version,
BigInt(heartbeat.schema_version),
heartbeat.sender_status,
BigInt(heartbeat.spool_pending),
BigInt(heartbeat.spool_inflight),
BigInt(heartbeat.spool_quarantined),
BigInt(heartbeat.spool_bytes),
heartbeat.first_loss_at,
heartbeat.tracking_degraded_reason,
heartbeat.last_error_code,
receivedAt
);
}
function sourceWhere(
filters: URLSearchParams,
includeTime = true
): { sql: string; values: SqlValue[] } {
const clauses: string[] = [];
const values: SqlValue[] = [];
const allowed = [
"deployment_id",
"workspace_id",
"thread_id",
"turn_id",
"workspace_dir",
"provider_profile_id",
"model_alias",
"scope",
];
for (const key of allowed) {
const value = filters.get(key === "model_alias" ? "model" : key);
if (!value) continue;
if (key === "thread_id") {
clauses.push("(thread_id=? OR source_session_id=?)");
values.push(value, value);
} else {
clauses.push(`${key}=?`);
values.push(value);
}
}
if (includeTime) {
if (filters.get("from")) {
clauses.push("completed_at>=?");
values.push(filters.get("from")!);
}
if (filters.get("to")) {
clauses.push("completed_at<=?");
values.push(filters.get("to")!);
}
}
return {
sql: clauses.length ? ` WHERE ${clauses.join(" AND ")}` : "",
values,
};
}
function stringValue(value: unknown): string {
return value === null || value === undefined ? "0" : String(value);
}
export function usageSummary(filters: URLSearchParams): UsageSummaryResult {
const db = usageDb();
const where = sourceWhere(filters);
const total = db
.prepare(
`SELECT
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens ELSE 0 END) input_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN output_tokens ELSE 0 END) output_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens+output_tokens ELSE 0 END) total_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN 1 ELSE 0 END) confirmed_call_count,
decimal_sum(CASE WHEN usage_status='unknown' THEN 1 ELSE 0 END) unknown_call_count
FROM model_usage${where.sql}`
)
.get(...where.values) as Record<string, unknown>;
const grouped = (
column: "scope" | "upstream_model_id" | "provider_profile_id"
) =>
(
db
.prepare(
`SELECT ${column} key,
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens ELSE 0 END) input_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN output_tokens ELSE 0 END) output_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens+output_tokens ELSE 0 END) total_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN 1 ELSE 0 END) confirmed_call_count,
decimal_sum(CASE WHEN usage_status='unknown' THEN 1 ELSE 0 END) unknown_call_count
FROM model_usage${where.sql} GROUP BY ${column}`
)
.all(...where.values) as Record<string, unknown>[]
)
.map((row) =>
Object.fromEntries(
Object.entries(row).map(([key, value]) => [key, stringValue(value)])
)
)
.sort((left, right) => {
const leftTotal = BigInt(left.total_tokens);
const rightTotal = BigInt(right.total_tokens);
return leftTotal === rightTotal ? 0 : leftTotal > rightTotal ? -1 : 1;
});
const serialized = Object.fromEntries(
Object.entries(total).map(([key, value]) => [key, stringValue(value)])
);
return {
deployment_id: filters.get("deployment_id"),
workspace_id: filters.get("workspace_id"),
...serialized,
by_scope: grouped("scope"),
by_model: grouped("upstream_model_id"),
by_provider: grouped("provider_profile_id"),
} as UsageSummaryResult;
}
export function usageCalls(
filters: URLSearchParams,
limit: number,
cursor: { at: string; deployment: string; id: string } | null
) {
const db = usageDb();
const where = sourceWhere(filters);
const clauses = where.sql ? [where.sql.slice(7)] : [];
const values = [...where.values];
if (cursor) {
clauses.push(`(completed_at < ? OR (completed_at = ? AND
(deployment_id < ? OR (deployment_id = ? AND model_call_id < ?))))`);
values.push(
cursor.at,
cursor.at,
cursor.deployment,
cursor.deployment,
cursor.id
);
}
const sqlWhere = clauses.length ? ` WHERE ${clauses.join(" AND ")}` : "";
const rows = db
.prepare(
`SELECT * FROM model_usage${sqlWhere}
ORDER BY completed_at DESC, deployment_id DESC, model_call_id DESC LIMIT ?`
)
.all(...values, BigInt(limit + 1)) as Record<string, unknown>[];
const hasMore = rows.length > limit;
const page = rows
.slice(0, limit)
.map((row) =>
Object.fromEntries(
Object.entries(row).map(([key, value]) => [
key,
typeof value === "bigint" ? value.toString() : value,
])
)
);
const last = page.at(-1);
const nextCursor =
hasMore && last
? Buffer.from(
JSON.stringify({
at: last.completed_at,
deployment: last.deployment_id,
id: last.model_call_id,
})
).toString("base64url")
: null;
return { items: page, next_cursor: nextCursor };
}
export function usageStatus(filters: URLSearchParams): UsageStatusResult {
const db = usageDb();
const statusFilters = new URLSearchParams();
for (const [key, environment] of [
["deployment_id", "EVOSCIENTIST_USAGE_DEPLOYMENT_ID"],
["workspace_id", "EVOSCIENTIST_WORKSPACE_ID"],
] as const) {
const value = filters.get(key) ?? process.env[environment]?.trim();
if (value) statusFilters.set(key, value);
}
const where = sourceWhere(statusFilters, false);
const row = db
.prepare(
`SELECT * FROM usage_sources${where.sql} ORDER BY last_seen_at DESC LIMIT 1`
)
.get(...where.values) as Record<string, unknown> | undefined;
if (!row)
return {
availability: "unavailable",
reason_code: "no_compatible_heartbeat",
};
const ageSeconds = (Date.now() - Date.parse(String(row.last_seen_at))) / 1000;
let availability = "available";
let reason_code = "healthy";
if (ageSeconds > USAGE_HEARTBEAT_TTL_SECONDS) {
availability = "offline";
reason_code = "heartbeat_expired";
} else if (
row.sender_status === "degraded" ||
row.first_loss_at ||
row.tracking_degraded_reason ||
row.spool_quarantined
) {
availability = "degraded";
reason_code = "tracking_incomplete";
}
return {
availability,
reason_code,
...Object.fromEntries(
Object.entries(row).map(([key, value]) => [
key,
typeof value === "bigint" ? value.toString() : value,
])
),
heartbeat_ttl_seconds: String(USAGE_HEARTBEAT_TTL_SECONDS),
};
}
function requireDecimal(value: string, field: string): void {
if (!/^(0|[1-9]\d*)$/.test(value)) {
throw new Error(`${field} must be a non-negative decimal string`);
}
}
export function recordProviderAggregate(
observation: ProviderAggregateObservation
): { status: "accepted" | "duplicate" | "conflict"; aggregate_id: string } {
for (const field of [
"input_tokens",
"output_tokens",
"request_count",
] as const) {
requireDecimal(observation[field], field);
}
if (observation.provider_total_tokens !== null) {
requireDecimal(observation.provider_total_tokens, "provider_total_tokens");
}
const start = Date.parse(observation.window_start);
const end = Date.parse(observation.window_end);
if (!Number.isFinite(start) || !Number.isFinite(end) || start >= end) {
throw new Error("provider aggregate has an invalid window");
}
if (!Number.isFinite(Date.parse(observation.observed_at))) {
throw new Error("provider aggregate has an invalid observed_at");
}
const payload = canonicalJson(observation);
const hash = createHash("sha256").update(payload).digest("hex");
const db = usageDb();
const stored = db
.prepare(
"SELECT payload_hash FROM provider_usage_aggregates WHERE aggregate_id=?"
)
.get(observation.aggregate_id) as { payload_hash: string } | undefined;
if (stored) {
return {
status: stored.payload_hash === hash ? "duplicate" : "conflict",
aggregate_id: observation.aggregate_id,
};
}
db.prepare(
`INSERT INTO provider_usage_aggregates (
aggregate_id,deployment_id,provider_profile_id,upstream_model_id,
window_start,window_end,input_tokens,output_tokens,provider_total_tokens,
request_count,source_name,source_payload_json,payload_hash,observed_at,
projection_policy_version
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`
).run(
observation.aggregate_id,
observation.deployment_id,
observation.provider_profile_id,
observation.upstream_model_id,
observation.window_start,
observation.window_end,
observation.input_tokens,
observation.output_tokens,
observation.provider_total_tokens,
observation.request_count,
observation.source_name,
canonicalJson(observation.source_payload),
hash,
observation.observed_at,
"aggregate-report-v1"
);
return { status: "accepted", aggregate_id: observation.aggregate_id };
}
export function providerAggregateDifference(aggregateId: string) {
const db = usageDb();
const aggregate = db
.prepare("SELECT * FROM provider_usage_aggregates WHERE aggregate_id=?")
.get(aggregateId) as Record<string, unknown> | undefined;
if (!aggregate) return null;
const clauses = [
"deployment_id=?",
"provider_profile_id=?",
"completed_at>=?",
"completed_at<?",
];
const values: SqlValue[] = [
String(aggregate.deployment_id),
String(aggregate.provider_profile_id),
String(aggregate.window_start),
String(aggregate.window_end),
];
if (aggregate.upstream_model_id !== null) {
clauses.push("upstream_model_id=?");
values.push(String(aggregate.upstream_model_id));
}
const local = db
.prepare(
`SELECT
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens ELSE 0 END) input_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN output_tokens ELSE 0 END) output_tokens,
decimal_sum(CASE WHEN usage_status='confirmed' THEN input_tokens+output_tokens ELSE 0 END) total_tokens,
decimal_sum(1) request_count,
decimal_sum(CASE WHEN usage_status='unknown' THEN 1 ELSE 0 END) unknown_count
FROM model_usage WHERE ${clauses.join(" AND ")}`
)
.get(...values) as Record<string, unknown>;
const providerInput = BigInt(String(aggregate.input_tokens));
const providerOutput = BigInt(String(aggregate.output_tokens));
const providerTotal =
aggregate.provider_total_tokens === null
? providerInput + providerOutput
: BigInt(String(aggregate.provider_total_tokens));
const localInput = BigInt(stringValue(local.input_tokens));
const localOutput = BigInt(stringValue(local.output_tokens));
const localTotal = BigInt(stringValue(local.total_tokens));
const localRequests = BigInt(stringValue(local.request_count));
const providerRequests = BigInt(String(aggregate.request_count));
return {
aggregate_id: aggregateId,
reconciliation_mode: "aggregate_only",
projection_policy_version: "aggregate-report-v1",
projection_updated: false,
provider: {
input_tokens: providerInput.toString(),
output_tokens: providerOutput.toString(),
total_tokens: providerTotal.toString(),
request_count: providerRequests.toString(),
},
local: {
input_tokens: localInput.toString(),
output_tokens: localOutput.toString(),
total_tokens: localTotal.toString(),
request_count: localRequests.toString(),
unknown_count: stringValue(local.unknown_count),
},
difference: {
input_tokens: (providerInput - localInput).toString(),
output_tokens: (providerOutput - localOutput).toString(),
total_tokens: (providerTotal - localTotal).toString(),
request_count: (providerRequests - localRequests).toString(),
},
};
}