feat(stats): added token aggregation and updated window grouping logic
- Added `fetchUsageData` and fleet token summation to aggregate token burn across clients. - Changed window grouping logic from window labels to limit IDs for distinct separation. - Updated `ProviderWindowInsight` properties and added tests for label suffixing and token summation.
This commit is contained in:
Generated
+120
-120
File diff suppressed because one or more lines are too long
@@ -2,6 +2,14 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Changed
|
||||
|
||||
- Window token estimates now use broker-held fleet token burn (per-client observed-usage reports) when an auth broker is configured, matching the fleet-wide window fractions instead of undercounting with local-only message stats.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed subscription-window insights merging distinct limits that share a duration label (Anthropic's overall vs model-scoped 7-day windows, Codex base vs Spark weeklies); interleaved fractions inflated window-equivalents consumed by orders of magnitude and broke tokens-per-window estimates. Windows now group by provider limit id.
|
||||
|
||||
## [17.3.6] - 2026-08-17
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -48,7 +48,7 @@ import type {
|
||||
RequestDetails,
|
||||
ToolDashboardStats,
|
||||
} from "./types";
|
||||
import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows";
|
||||
import { computeUsageWindowStats, fetchUsageData } from "./usage-windows";
|
||||
|
||||
const STATS_SYNC_LOCK_RETRY_MS = 25;
|
||||
const STATS_SYNC_LOCK_WAIT_MS = 60 * 60 * 1000;
|
||||
@@ -544,14 +544,18 @@ export async function getToolDashboardStats(range?: string | null): Promise<Tool
|
||||
* Get the providers dashboard payload: per-provider totals, peak-burn-hours
|
||||
* histogram, provider token time series, and subscription-window analytics
|
||||
* (utilization series + insights) derived from recorded usage-limit snapshots.
|
||||
*
|
||||
* Window token estimates use broker-held fleet token burn when a broker is
|
||||
* configured — the window fractions cover every install sharing the broker's
|
||||
* credentials, so dividing them into local-only tokens would undercount.
|
||||
*/
|
||||
export async function getProviderDashboardStats(range?: string | null): Promise<ProviderDashboardStats> {
|
||||
await initDb();
|
||||
const { modelSeriesDays, modelSeriesBucketMs, cutoff } = getTimeRangeConfig(range);
|
||||
const providers = getStatsByProvider(cutoff ?? undefined);
|
||||
const tokensByProvider = new Map(providers.map(p => [p.provider, p.totalTokens]));
|
||||
const snapshots = await fetchUsageSnapshots(cutoff ?? 0);
|
||||
const { usageSeries, windowInsights } = computeUsageWindowStats(snapshots, tokensByProvider);
|
||||
const usage = await fetchUsageData(cutoff ?? 0);
|
||||
const tokensByProvider = usage.fleetTokensByProvider ?? new Map(providers.map(p => [p.provider, p.totalTokens]));
|
||||
const { usageSeries, windowInsights } = computeUsageWindowStats(usage.rows, tokensByProvider);
|
||||
return {
|
||||
providers,
|
||||
hourly: getProviderHourlyBurn(cutoff ?? undefined),
|
||||
|
||||
@@ -386,8 +386,9 @@ export interface UsageWindowSeries {
|
||||
accountKey: string;
|
||||
/** Email/account id when known, else the stable account key. */
|
||||
accountLabel: string;
|
||||
/** Groups the same limit window across accounts (window label or limit id). */
|
||||
/** Groups the same limit window across accounts (the provider limit id). */
|
||||
windowKey: string;
|
||||
/** Human label of the limit (distinguishes same-duration windows). */
|
||||
windowLabel: string;
|
||||
points: UsageWindowPoint[];
|
||||
}
|
||||
@@ -399,7 +400,9 @@ export interface UsageWindowSeries {
|
||||
*/
|
||||
export interface ProviderWindowInsight {
|
||||
provider: string;
|
||||
/** Groups the same limit window across accounts (the provider limit id). */
|
||||
windowKey: string;
|
||||
/** Human label of the limit (distinguishes same-duration windows). */
|
||||
windowLabel: string;
|
||||
/** Accounts with at least one snapshot for this window in range. */
|
||||
accounts: number;
|
||||
|
||||
@@ -14,6 +14,7 @@
|
||||
*/
|
||||
import { Database } from "bun:sqlite";
|
||||
import { AuthBrokerClient, resolveAuthBrokerConfig } from "@oh-my-pi/pi-ai/auth-broker";
|
||||
import type { ClientUsageClientSummary } from "@oh-my-pi/pi-ai/usage";
|
||||
import { getAgentDbPath, logger } from "@oh-my-pi/pi-utils";
|
||||
import type { ProviderWindowInsight, UsageWindowPoint, UsageWindowSeries } from "./shared-types";
|
||||
|
||||
@@ -40,6 +41,19 @@ export interface UsageWindowStats {
|
||||
windowInsights: ProviderWindowInsight[];
|
||||
}
|
||||
|
||||
/** Usage snapshots plus, in broker mode, fleet-wide token burn per provider. */
|
||||
export interface UsageDataSnapshot {
|
||||
rows: UsageSnapshotRow[];
|
||||
/**
|
||||
* Total tokens (input + output + cache read/write) per provider summed
|
||||
* across every install reporting to the auth broker, or `null` when no
|
||||
* broker is configured or no client reports exist for the range. Matches
|
||||
* the fleet-wide window fractions in `rows`, unlike local message stats
|
||||
* which only see this install's burn.
|
||||
*/
|
||||
fleetTokensByProvider: Map<string, number> | null;
|
||||
}
|
||||
|
||||
/** A used-fraction drop smaller than this is jitter, not a window reset. */
|
||||
const RESET_DROP_THRESHOLD = 0.05;
|
||||
/** Minimum window-equivalents consumed before extrapolating tokens/window. */
|
||||
@@ -101,35 +115,74 @@ export function readUsageSnapshots(sinceMs: number, dbPath = getAgentDbPath()):
|
||||
}
|
||||
|
||||
/**
|
||||
* Fetch usage snapshots from wherever they actually accumulate: the auth
|
||||
* broker's durable history when a broker is configured (the broker performs
|
||||
* every upstream usage fetch in that mode, so the local `usage_history` stays
|
||||
* frozen), else the local agent DB. Broker errors fall back to the local read
|
||||
* so the dashboard degrades to stale-but-present data instead of failing.
|
||||
* Fetch usage data from wherever it actually accumulates: the auth broker's
|
||||
* durable history plus per-client observed-usage reports when a broker is
|
||||
* configured (the broker performs every upstream usage fetch in that mode, so
|
||||
* the local `usage_history` stays frozen), else the local agent DB. Broker
|
||||
* errors fall back to the local read so the dashboard degrades to
|
||||
* stale-but-present data instead of failing.
|
||||
*/
|
||||
export async function fetchUsageSnapshots(sinceMs: number): Promise<UsageSnapshotRow[]> {
|
||||
export async function fetchUsageData(sinceMs: number): Promise<UsageDataSnapshot> {
|
||||
try {
|
||||
const brokerConfig = await resolveAuthBrokerConfig();
|
||||
if (brokerConfig) {
|
||||
const client = new AuthBrokerClient({ url: brokerConfig.url, token: brokerConfig.token });
|
||||
const response = await client.fetchUsageHistory({ sinceMs });
|
||||
return response.entries.map(entry => ({
|
||||
recordedAt: entry.recordedAt,
|
||||
provider: entry.provider,
|
||||
accountKey: entry.accountKey,
|
||||
email: entry.email ?? null,
|
||||
accountId: entry.accountId ?? null,
|
||||
limitId: entry.limitId,
|
||||
label: entry.label,
|
||||
windowLabel: entry.windowLabel ?? null,
|
||||
usedFraction: entry.usedFraction ?? null,
|
||||
status: entry.status ?? null,
|
||||
}));
|
||||
const [response, fleetTokensByProvider] = await Promise.all([
|
||||
client.fetchUsageHistory({ sinceMs }),
|
||||
fetchFleetTokens(client, sinceMs),
|
||||
]);
|
||||
return {
|
||||
rows: response.entries.map(entry => ({
|
||||
recordedAt: entry.recordedAt,
|
||||
provider: entry.provider,
|
||||
accountKey: entry.accountKey,
|
||||
email: entry.email ?? null,
|
||||
accountId: entry.accountId ?? null,
|
||||
limitId: entry.limitId,
|
||||
label: entry.label,
|
||||
windowLabel: entry.windowLabel ?? null,
|
||||
usedFraction: entry.usedFraction ?? null,
|
||||
status: entry.status ?? null,
|
||||
})),
|
||||
fleetTokensByProvider,
|
||||
};
|
||||
}
|
||||
} catch (err) {
|
||||
logger.debug("broker usage history unavailable, falling back to local", { error: String(err) });
|
||||
}
|
||||
return readUsageSnapshots(sinceMs);
|
||||
return { rows: readUsageSnapshots(sinceMs), fleetTokensByProvider: null };
|
||||
}
|
||||
|
||||
/**
|
||||
* Sum broker-recorded client token burn per provider since `sinceMs`.
|
||||
* Returns `null` on fetch failure or when no client has reported usage, so
|
||||
* callers fall back to local message stats instead of zeroing estimates.
|
||||
*/
|
||||
async function fetchFleetTokens(client: AuthBrokerClient, sinceMs: number): Promise<Map<string, number> | null> {
|
||||
try {
|
||||
const summary = await client.fetchClientUsageSummary({ sinceMs });
|
||||
return sumFleetTokens(summary.clients);
|
||||
} catch (err) {
|
||||
logger.debug("broker client usage summary unavailable", { error: String(err) });
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fold per-client provider aggregates into total tokens per provider
|
||||
* (input + output + cache read/write, matching message-stat `totalTokens`).
|
||||
* Returns `null` when no client reported anything, signalling "no data"
|
||||
* rather than "zero burn".
|
||||
*/
|
||||
export function sumFleetTokens(clients: ClientUsageClientSummary[]): Map<string, number> | null {
|
||||
const tokens = new Map<string, number>();
|
||||
for (const client of clients) {
|
||||
for (const p of client.providers) {
|
||||
const total = p.inputTokens + p.outputTokens + p.cacheReadTokens + p.cacheWriteTokens;
|
||||
tokens.set(p.provider, (tokens.get(p.provider) ?? 0) + total);
|
||||
}
|
||||
}
|
||||
return tokens.size > 0 ? tokens : null;
|
||||
}
|
||||
|
||||
/** True when a snapshot reports an exhausted window, by status or by fraction. */
|
||||
@@ -175,14 +228,33 @@ interface WindowGroup {
|
||||
accounts: Map<string, AccountSeries>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Display label for one limit window: the limit label, plus the window label
|
||||
* when it adds information ("Usage (Google) · Daily"). The limit label alone
|
||||
* distinguishes same-duration windows ("Claude 7 Day" vs "Claude 7 Day
|
||||
* (Fable)"); the window-label suffix distinguishes same-named limits with
|
||||
* different durations (Antigravity's daily vs weekly "Usage (Google)").
|
||||
*/
|
||||
function windowDisplayLabel(row: UsageSnapshotRow): string {
|
||||
const { label, windowLabel } = row;
|
||||
if (!windowLabel || label.toLowerCase().includes(windowLabel.toLowerCase())) return label;
|
||||
return `${label} · ${windowLabel}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive utilization series and per-window insights from raw snapshots.
|
||||
*
|
||||
* Windows are grouped by `(provider, limitId)` — never by display label.
|
||||
* Distinct limits can share a duration label (Anthropic's overall and
|
||||
* model-scoped 7-day windows, Codex's base and Spark weeklies); merging them
|
||||
* interleaves unrelated fractions per account and wildly inflates consumption.
|
||||
*
|
||||
* `tokensByProvider` supplies each provider's token burn over the same time
|
||||
* range (from the local message stats); it converts consumed window fraction
|
||||
* into an estimated token capacity per window. Attribution note: tokens are
|
||||
* per provider, not per account, so the estimate treats the account fleet as
|
||||
* one pooled subscription — which is exactly how round-robin auth uses it.
|
||||
* range (fleet-wide broker client reports when available, else local message
|
||||
* stats); it converts consumed window fraction into an estimated token
|
||||
* capacity per window. Attribution note: tokens are per provider, not per
|
||||
* account, so the estimate treats the account fleet as one pooled
|
||||
* subscription — which is exactly how round-robin auth uses it.
|
||||
*/
|
||||
export function computeUsageWindowStats(
|
||||
rows: UsageSnapshotRow[],
|
||||
@@ -190,15 +262,19 @@ export function computeUsageWindowStats(
|
||||
): UsageWindowStats {
|
||||
const groups = new Map<string, WindowGroup>();
|
||||
for (const row of rows) {
|
||||
const windowKey = row.windowLabel ?? row.limitId;
|
||||
const groupKey = `${row.provider}\u0000${windowKey}`;
|
||||
const groupKey = `${row.provider}\u0000${row.limitId}`;
|
||||
let group = groups.get(groupKey);
|
||||
if (!group) {
|
||||
group = { provider: row.provider, windowKey, windowLabel: row.windowLabel ?? row.label, accounts: new Map() };
|
||||
group = {
|
||||
provider: row.provider,
|
||||
windowKey: row.limitId,
|
||||
windowLabel: windowDisplayLabel(row),
|
||||
accounts: new Map(),
|
||||
};
|
||||
groups.set(groupKey, group);
|
||||
}
|
||||
// Labels can change across snapshots (provider renames); latest wins.
|
||||
group.windowLabel = row.windowLabel ?? row.label;
|
||||
group.windowLabel = windowDisplayLabel(row);
|
||||
let account = group.accounts.get(row.accountKey);
|
||||
if (!account) {
|
||||
account = { accountKey: row.accountKey, accountLabel: row.email ?? row.accountId ?? row.accountKey, rows: [] };
|
||||
|
||||
@@ -5,7 +5,12 @@ import * as path from "node:path";
|
||||
import { getProviderDashboardStats } from "@oh-my-pi/omp-stats/aggregator";
|
||||
import { initDb, insertMessageStats } from "@oh-my-pi/omp-stats/db";
|
||||
import type { MessageStats } from "@oh-my-pi/omp-stats/types";
|
||||
import { computeUsageWindowStats, readUsageSnapshots, type UsageSnapshotRow } from "@oh-my-pi/omp-stats/usage-windows";
|
||||
import {
|
||||
computeUsageWindowStats,
|
||||
readUsageSnapshots,
|
||||
sumFleetTokens,
|
||||
type UsageSnapshotRow,
|
||||
} from "@oh-my-pi/omp-stats/usage-windows";
|
||||
import { getAgentDbPath } from "@oh-my-pi/pi-utils";
|
||||
import { installStatsTestIsolation } from "./helpers/temp-agent";
|
||||
|
||||
@@ -147,7 +152,7 @@ describe("computeUsageWindowStats", () => {
|
||||
expect(windowInsights[0].idealAccounts).toBe(1);
|
||||
});
|
||||
|
||||
it("keeps windows with distinct labels separate", () => {
|
||||
it("keeps windows with distinct limit ids separate", () => {
|
||||
const rows: UsageSnapshotRow[] = [
|
||||
snapshot({ recordedAt: T0, usedFraction: 0.2, limitId: "5h", windowLabel: "5h" }),
|
||||
snapshot({ recordedAt: T0, usedFraction: 0.1, limitId: "weekly", windowLabel: "Weekly", label: "Weekly" }),
|
||||
@@ -161,7 +166,100 @@ describe("computeUsageWindowStats", () => {
|
||||
}),
|
||||
];
|
||||
const { windowInsights } = computeUsageWindowStats(rows, new Map());
|
||||
expect(windowInsights.map(i => i.windowKey).sort()).toEqual(["5h", "Weekly"]);
|
||||
expect(windowInsights.map(i => i.windowKey).sort()).toEqual(["5h", "weekly"]);
|
||||
});
|
||||
|
||||
it("never merges distinct limits sharing a window label", () => {
|
||||
// Anthropic reports an overall 7-day window and a model-scoped one with
|
||||
// the same "7 Day" window label. Grouping by label interleaved the two
|
||||
// fraction series per account, inflating consumption from oscillation:
|
||||
// here 0.8 - 0.1 = 0.7 fake burn per interleaved pair.
|
||||
const base = { windowLabel: "7 Day", accountKey: "acct-1" };
|
||||
const rows: UsageSnapshotRow[] = [
|
||||
snapshot({ ...base, recordedAt: T0 + 0 * MINUTE, limitId: "7d", label: "Claude 7 Day", usedFraction: 0.8 }),
|
||||
snapshot({
|
||||
...base,
|
||||
recordedAt: T0 + 1 * MINUTE,
|
||||
limitId: "7d:opus",
|
||||
label: "Claude 7 Day (Opus)",
|
||||
usedFraction: 0.1,
|
||||
}),
|
||||
snapshot({ ...base, recordedAt: T0 + 2 * MINUTE, limitId: "7d", label: "Claude 7 Day", usedFraction: 0.9 }),
|
||||
snapshot({
|
||||
...base,
|
||||
recordedAt: T0 + 3 * MINUTE,
|
||||
limitId: "7d:opus",
|
||||
label: "Claude 7 Day (Opus)",
|
||||
usedFraction: 0.15,
|
||||
}),
|
||||
];
|
||||
const { windowInsights } = computeUsageWindowStats(rows, new Map([["prov-a", 1_000_000]]));
|
||||
|
||||
expect(windowInsights.map(i => i.windowKey).sort()).toEqual(["7d", "7d:opus"]);
|
||||
const overall = windowInsights.find(i => i.windowKey === "7d");
|
||||
const opus = windowInsights.find(i => i.windowKey === "7d:opus");
|
||||
expect(overall?.fractionConsumed).toBeCloseTo(0.1, 10);
|
||||
expect(opus?.fractionConsumed).toBeCloseTo(0.05, 10);
|
||||
expect(overall?.cycles).toBe(0);
|
||||
// The limit label (not the shared window label) tells the two apart.
|
||||
expect(overall?.windowLabel).toBe("Claude 7 Day");
|
||||
expect(opus?.windowLabel).toBe("Claude 7 Day (Opus)");
|
||||
});
|
||||
|
||||
it("suffixes the window label when the limit label lacks the duration", () => {
|
||||
// Antigravity exposes daily and weekly limits under the same limit
|
||||
// label; the display label must carry the duration to tell them apart.
|
||||
const rows: UsageSnapshotRow[] = [
|
||||
snapshot({
|
||||
recordedAt: T0,
|
||||
limitId: "g:daily",
|
||||
label: "Usage (Google)",
|
||||
windowLabel: "Daily",
|
||||
usedFraction: 0.2,
|
||||
}),
|
||||
snapshot({
|
||||
recordedAt: T0,
|
||||
limitId: "g:weekly",
|
||||
label: "Usage (Google)",
|
||||
windowLabel: "Weekly",
|
||||
usedFraction: 0.1,
|
||||
}),
|
||||
];
|
||||
const { windowInsights } = computeUsageWindowStats(rows, new Map());
|
||||
expect(windowInsights.map(i => i.windowLabel).sort()).toEqual([
|
||||
"Usage (Google) · Daily",
|
||||
"Usage (Google) · Weekly",
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
describe("sumFleetTokens", () => {
|
||||
it("sums all four token components per provider across clients, null when empty", () => {
|
||||
const client = (installId: string, provider: string, tokens: [number, number, number, number]) => ({
|
||||
installId,
|
||||
firstSeen: T0,
|
||||
lastSeen: T0,
|
||||
providers: [
|
||||
{
|
||||
provider,
|
||||
requests: 1,
|
||||
inputTokens: tokens[0],
|
||||
outputTokens: tokens[1],
|
||||
cacheReadTokens: tokens[2],
|
||||
cacheWriteTokens: tokens[3],
|
||||
costUsd: 0,
|
||||
},
|
||||
],
|
||||
});
|
||||
const tokens = sumFleetTokens([
|
||||
client("install-1", "prov-a", [100, 20, 300, 4]),
|
||||
client("install-2", "prov-a", [1, 2, 3, 4]),
|
||||
client("install-3", "prov-b", [10, 0, 0, 0]),
|
||||
]);
|
||||
expect(tokens?.get("prov-a")).toBe(434);
|
||||
expect(tokens?.get("prov-b")).toBe(10);
|
||||
// No reports must read as "no data" (fall back to local stats), not zero burn.
|
||||
expect(sumFleetTokens([])).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user