feat: implemented client usage tracking and reporting in the auth broker

- Add usage history, reporting, and client summary endpoints to the auth broker.
- Implement SQLite persistence, batching, and periodic flushing for observed client usage.
- Integrate usage reporting into the coding agent session for completed assistant messages.
- Update stats aggregator and dashboard to retrieve usage history snapshots from the broker.
This commit is contained in:
can1357
2026-07-24 15:20:46 +02:00
parent c6dcfa4137
commit 3676b5d97f
15 changed files with 759 additions and 7 deletions
+5
View File
@@ -2,6 +2,11 @@
## [Unreleased]
### Added
- Added `GET /v1/usage/history` to the auth broker (recorded usage-limit snapshots with `sinceMs`/`provider` filters) and `AuthBrokerClient.fetchUsageHistory` — in broker deployments the broker host performs every upstream usage fetch, so its durable history is the only complete utilization record
- Added per-client burn tracking to the auth broker: clients batch observed request usage per (provider, model) and flush it to `POST /v1/usage/observed` every 10 seconds (install id as client key, hostname as display name); the broker persists 5-minute buckets in `client_usage`/`clients` and serves aggregates from `GET /v1/usage/clients`. Brokers without the endpoint disable reporting for the process lifetime
### Changed
- Renamed the Z.AI feature quota row to `ZAI Zread Quota` (tier/id slug `zread`), replacing the 74-char `ZAI Web Search / Reader / Zread Quota (web-search-reader-zread)` title that wrapped `omp usage` rows
+39
View File
@@ -9,6 +9,9 @@ import { readSseEvents } from "@oh-my-pi/pi-utils";
import { type } from "arktype";
import type { AuthCredential } from "../auth-storage";
import type {
ClientUsageReportRequest,
ClientUsageReportResponse,
ClientUsageSummaryResponse,
CredentialBlockRequest,
CredentialBlockResponse,
CredentialBlocksDeleteResponse,
@@ -20,10 +23,13 @@ import type {
HealthzResponse,
SnapshotResponse,
SnapshotStreamEvent,
UsageHistoryResponse,
UsageResponse,
UsageStaleResponse,
} from "./types";
import {
clientUsageReportResponseSchema,
clientUsageSummaryResponseSchema,
credentialBlockResponseSchema,
credentialBlocksDeleteResponseSchema,
credentialDisableResponseSchema,
@@ -32,6 +38,7 @@ import {
healthzResponseSchema,
snapshotResponseSchema,
snapshotStreamEventSchema,
usageHistoryResponseSchema,
usageResponseSchema,
usageStaleResponseSchema,
} from "./wire-schemas";
@@ -244,6 +251,38 @@ export class AuthBrokerClient {
return this.#request<UsageResponse>("GET", "/v1/usage", { schema: usageResponseSchema, signal });
}
/** Recorded usage-limit snapshots from the broker host, oldest first. */
fetchUsageHistory(
query?: { sinceMs?: number; provider?: string },
signal?: AbortSignal,
): Promise<UsageHistoryResponse> {
const params = new URLSearchParams();
if (query?.sinceMs !== undefined) params.set("sinceMs", String(query.sinceMs));
if (query?.provider) params.set("provider", query.provider);
const path = `/v1/usage/history${params.size > 0 ? `?${params.toString()}` : ""}`;
return this.#request<UsageHistoryResponse>("GET", path, { schema: usageHistoryResponseSchema, signal });
}
/** Report this client's batched observed request usage for per-install burn tracking. */
reportClientUsage(report: ClientUsageReportRequest, signal?: AbortSignal): Promise<ClientUsageReportResponse> {
return this.#request<ClientUsageReportResponse>("POST", "/v1/usage/observed", {
body: report,
schema: clientUsageReportResponseSchema,
signal,
});
}
/** Per-client token burn aggregates recorded by the broker host. */
fetchClientUsageSummary(query?: { sinceMs?: number }, signal?: AbortSignal): Promise<ClientUsageSummaryResponse> {
const params = new URLSearchParams();
if (query?.sinceMs !== undefined) params.set("sinceMs", String(query.sinceMs));
const path = `/v1/usage/clients${params.size > 0 ? `?${params.toString()}` : ""}`;
return this.#request<ClientUsageSummaryResponse>("GET", path, {
schema: clientUsageSummaryResponseSchema,
signal,
});
}
notifyUsageStale(signal?: AbortSignal): Promise<UsageStaleResponse> {
return this.#request<UsageStaleResponse>("POST", "/v1/usage/stale", {
schema: usageStaleResponseSchema,
+76 -3
View File
@@ -7,8 +7,9 @@
* usage reports cache TTL is 5 minutes per credential, so durability across
* runs isn't required.
*/
import * as os from "node:os";
import { scheduler } from "node:timers/promises";
import { logger } from "@oh-my-pi/pi-utils";
import { getInstallId, logger } from "@oh-my-pi/pi-utils";
import {
type AuthCredential,
type AuthCredentialSnapshotEntry,
@@ -21,8 +22,8 @@ import {
import * as AIError from "../error";
import type { OAuthCredentials } from "../registry/oauth/types";
import type { Provider } from "../types";
import type { UsageReport } from "../usage";
import { type AuthBrokerClient, AuthBrokerStreamUnsupportedError } from "./client";
import type { ObservedUsageEntry, UsageReport } from "../usage";
import { type AuthBrokerClient, AuthBrokerError, AuthBrokerStreamUnsupportedError } from "./client";
import type {
CredentialBlockSnapshot,
RefresherSchedule,
@@ -232,6 +233,8 @@ export interface RemoteAuthCredentialStoreOptions {
* routing policy, not broker authorization.
*/
accountPool?: AuthBrokerAccountPool;
/** Flush cadence for batched observed-usage reports. Default 10s. */
observedUsageFlushMs?: number;
}
export class RemoteAuthCredentialStore implements AuthCredentialStore {
@@ -259,10 +262,17 @@ export class RemoteAuthCredentialStore implements AuthCredentialStore {
#streamingActive = false;
/** Latched once the broker has answered 404 — never try the stream again. */
#streamingUnsupported = false;
/** Pending observed usage keyed by `provider\u0000model`, merged until flush. */
#observedUsage = new Map<string, ObservedUsageEntry>();
#observedUsageTimer: Timer | undefined;
readonly #observedUsageFlushMs: number;
/** Latched once the broker answered 404 — old broker, never report again. */
#observedUsageUnsupported = false;
constructor(opts: RemoteAuthCredentialStoreOptions) {
this.#client = opts.client;
this.#streamSnapshots = opts.streamSnapshots ?? true;
this.#observedUsageFlushMs = opts.observedUsageFlushMs ?? 10_000;
this.#accountPool = opts.accountPool
? new Map([...opts.accountPool].map(([provider, identities]) => [provider, new Set(identities)]))
: undefined;
@@ -1049,10 +1059,73 @@ export class RemoteAuthCredentialStore implements AuthCredentialStore {
return inflight;
}
/**
* Fold locally observed request usage into the pending report and schedule
* a flush. One `POST /v1/usage/observed` at most per flush interval; on
* failure the batch is retained and retried with the next flush. A 404
* (pre-endpoint broker) disables reporting for the life of this store.
*/
recordObservedUsage(entries: ObservedUsageEntry[]): void {
if (this.#closed || this.#observedUsageUnsupported) return;
for (const entry of entries) {
const key = `${entry.provider}\u0000${entry.model}`;
const pending = this.#observedUsage.get(key);
if (pending) {
pending.at = Math.max(pending.at, entry.at);
pending.requests += entry.requests;
pending.inputTokens += entry.inputTokens;
pending.outputTokens += entry.outputTokens;
pending.cacheReadTokens += entry.cacheReadTokens;
pending.cacheWriteTokens += entry.cacheWriteTokens;
pending.costUsd += entry.costUsd;
} else {
this.#observedUsage.set(key, { ...entry });
}
}
if (this.#observedUsage.size > 0 && this.#observedUsageTimer === undefined) {
this.#observedUsageTimer = setTimeout(() => {
this.#observedUsageTimer = undefined;
void this.#flushObservedUsage();
}, this.#observedUsageFlushMs);
this.#observedUsageTimer.unref?.();
}
}
async #flushObservedUsage(): Promise<void> {
if (this.#observedUsage.size === 0 || this.#observedUsageUnsupported) return;
const batch = [...this.#observedUsage.values()];
this.#observedUsage.clear();
try {
await this.#client.reportClientUsage({
installId: getInstallId(),
hostname: os.hostname(),
entries: batch,
});
} catch (error) {
const status = error instanceof AuthBrokerError ? error.status : undefined;
if (status === 404 || status === 501) {
// Broker predates the endpoint (or store can't persist) — stop trying.
this.#observedUsageUnsupported = true;
logger.debug("auth-broker does not accept observed usage; reporting disabled", { status });
return;
}
logger.debug("auth-broker observed usage flush failed; retrying next flush", { error: String(error) });
// Merge the failed batch back under the (possibly refilled) buffer so
// nothing is lost; bounded because entries are keyed per (provider, model).
if (!this.#closed) this.recordObservedUsage(batch);
}
}
close(): void {
if (this.#closed) return;
this.#closed = true;
this.#backgroundAbort.abort();
if (this.#observedUsageTimer !== undefined) {
clearTimeout(this.#observedUsageTimer);
this.#observedUsageTimer = undefined;
}
// Best-effort final flush; failures are dropped (the process is exiting).
if (this.#observedUsage.size > 0) void this.#flushObservedUsage();
this.#cache.clear();
this.#usageOverlays.clear();
}
+40
View File
@@ -15,6 +15,7 @@ import type { AuthStorage, StoredCredentialBlock } from "../auth-storage";
import { parseBind } from "../utils/parse-bind";
import { AuthBrokerRefresher, type AuthBrokerRefresherSchedule } from "./refresher";
import type {
ClientUsageReportRequest,
CredentialBlockResponse,
CredentialBlockSnapshot,
CredentialBlocksDeleteResponse,
@@ -37,6 +38,7 @@ import {
DEFAULT_STREAM_KEEPALIVE_MS,
} from "./types";
import {
clientUsageReportRequestSchema,
credentialBlockRequestSchema,
credentialDisableRequestSchema,
credentialUploadRequestSchema,
@@ -608,6 +610,44 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
return json(502, { error: message });
}
}
if (req.method === "GET" && pathname === "/v1/usage/history") {
const sinceMsRaw = url.searchParams.get("sinceMs");
const sinceMsParsed = sinceMsRaw === null ? undefined : Number.parseInt(sinceMsRaw, 10);
const sinceMs =
sinceMsParsed !== undefined && Number.isFinite(sinceMsParsed) ? sinceMsParsed : undefined;
const provider = url.searchParams.get("provider") ?? undefined;
const entries = opts.storage.listUsageHistory({ sinceMs, provider });
logger.info("auth-broker usage history served", { peer, entries: entries.length, sinceMs, provider });
return json(200, { generatedAt: Date.now(), entries });
}
if (req.method === "POST" && pathname === "/v1/usage/observed") {
const parsed = await parseBody(req, clientUsageReportRequestSchema);
if (!parsed.ok) return parsed.response;
// Arktype's inferred union collides the `entries` field with
// Array.prototype.entries; the schema already validated the shape.
const report = parsed.data as ClientUsageReportRequest;
try {
const recorded = opts.storage.recordClientUsage(report);
if (!recorded) return json(501, { error: "broker store does not persist client usage" });
logger.debug("auth-broker client usage recorded", {
peer,
installId: report.installId,
hostname: report.hostname,
entries: report.entries.length,
});
return json(200, { ok: true });
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
logger.warn("auth-broker client usage record failed", { peer, error: message });
return json(500, { error: message });
}
}
if (req.method === "GET" && pathname === "/v1/usage/clients") {
const sinceMsRaw = url.searchParams.get("sinceMs");
const sinceMsParsed = sinceMsRaw === null ? Number.NaN : Number.parseInt(sinceMsRaw, 10);
const summary = opts.storage.getClientUsageSummary(Number.isFinite(sinceMsParsed) ? sinceMsParsed : 0);
return json(200, { generatedAt: Date.now(), clients: summary.clients });
}
if (req.method === "POST" && pathname === "/v1/usage/stale") {
try {
opts.storage.invalidateUsageCache?.();
+25 -1
View File
@@ -12,7 +12,7 @@ import type {
AuthCredentialSnapshotEntry,
StoredCredentialBlock,
} from "../auth-storage";
import type { UsageReport } from "../usage";
import type { ClientUsageClientSummary, ClientUsageReport, UsageHistoryEntry, UsageReport } from "../usage";
/** GET /v1/healthz response body. */
export interface HealthzResponse {
@@ -47,6 +47,30 @@ export interface UsageResponse {
reports: UsageReport[];
}
/**
* GET /v1/usage/history response body. Entries come from the broker host's
* durable `usage_history` — the broker performs every upstream usage fetch in
* broker deployments, so this is the only complete utilization record.
*/
export interface UsageHistoryResponse {
generatedAt: number;
entries: UsageHistoryEntry[];
}
/** POST /v1/usage/observed request body — one client's batched observed usage. */
export type ClientUsageReportRequest = ClientUsageReport;
/** POST /v1/usage/observed response body. */
export interface ClientUsageReportResponse {
ok: boolean;
}
/** GET /v1/usage/clients response body — per-client token burn aggregates. */
export interface ClientUsageSummaryResponse {
generatedAt: number;
clients: ClientUsageClientSummary[];
}
/** POST /v1/credential/:id/refresh response body. */
export interface CredentialRefreshResponse {
entry: AuthCredentialSnapshotEntry;
@@ -230,6 +230,77 @@ export const usageResponseSchema = type({
reports: arkUsageReportSchema.array(),
});
const usageHistoryEntrySchema = type({
recordedAt: "number",
provider: "string",
accountKey: "string",
"email?": "string",
"accountId?": "string",
limitId: "string",
label: "string",
"windowLabel?": "string",
"usedFraction?": "number",
"status?": "'ok' | 'warning' | 'exhausted' | 'unknown'",
"resetsAt?": "number",
});
/** Broker `/v1/usage/history` response — recorded usage-limit snapshots, oldest first. */
export const usageHistoryResponseSchema = type({
"+": "reject",
generatedAt: "number",
entries: usageHistoryEntrySchema.array(),
});
const observedUsageEntrySchema = type({
at: "number",
provider: "string",
model: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
/** Broker `POST /v1/usage/observed` request — one client's batched observed usage. */
export const clientUsageReportRequestSchema = type({
"+": "reject",
installId: "string",
"hostname?": "string",
entries: observedUsageEntrySchema.array(),
});
export const clientUsageReportResponseSchema = type({
"+": "reject",
ok: "boolean",
});
const clientProviderUsageSchema = type({
provider: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
const clientUsageClientSummarySchema = type({
installId: "string",
"hostname?": "string",
firstSeen: "number",
lastSeen: "number",
providers: clientProviderUsageSchema.array(),
});
/** Broker `GET /v1/usage/clients` response — per-client token burn aggregates. */
export const clientUsageSummaryResponseSchema = type({
"+": "reject",
generatedAt: "number",
clients: clientUsageClientSummarySchema.array(),
});
// ─── Refresh ───────────────────────────────────────────────────────────────
export const credentialRefreshResponseSchema = type({
+198
View File
@@ -28,8 +28,12 @@ import type {
import { getEnvApiKey, getEnvApiKeyName } from "./stream";
import type { Provider } from "./types";
import type {
ClientProviderUsage,
ClientUsageReport,
ClientUsageSummary,
CredentialRankingContext,
CredentialRankingStrategy,
ObservedUsageEntry,
UsageCostHistoryEntry,
UsageCostHistoryQuery,
UsageCredential,
@@ -399,6 +403,16 @@ export interface AuthCredentialStore {
listUsageCosts?(query?: UsageCostHistoryQuery): UsageCostHistoryEntry[];
/** Read recorded usage-limit snapshots, oldest first. */
listUsageHistory?(query?: UsageHistoryQuery): UsageHistoryEntry[];
/**
* Client hook: forward locally observed request usage. Remote broker stores
* batch these to the broker so it can attribute token burn per install;
* local stores omit it and observation is skipped.
*/
recordObservedUsage?(entries: ObservedUsageEntry[]): void;
/** Broker host: persist one client's observed-usage report. */
recordClientUsage?(report: ClientUsageReport): void;
/** Broker host: aggregate recorded per-client usage since a timestamp. */
getClientUsageSummary?(sinceMs: number): ClientUsageSummary;
/**
* Optional store-supplied OAuth refresh. When present, `AuthStorage` uses
* it before the per-provider local refresh path. `RemoteAuthCredentialStore`
@@ -624,6 +638,12 @@ const USAGE_LAST_GOOD_RETENTION_MS = 24 * 60 * 60_000;
* unnecessary — 1 row/hour is ~9k rows per account window per year.
*/
const USAGE_HISTORY_BUCKET_MS = 60 * 60_000;
/**
* Merge client observed-usage flushes into at most one row per 5 minutes per
* (install, provider, model): ~300 rows/day per active model per client
* instead of one row per 10s flush.
*/
const CLIENT_USAGE_BUCKET_MS = 5 * 60_000;
/**
* Per-credential cool-down after a usage fetch fails. While this window is
* active we serve the last successful value to avoid dropping the credential
@@ -3142,6 +3162,56 @@ export class AuthStorage {
}
}
/**
* Forward one completed request's usage to the store's observer hook.
* Broker-backed stores batch these into per-install reports so the broker
* can track actual token burn per client; local stores have no hook and
* the call is a no-op.
*/
recordObservedUsage(entry: {
provider: Provider;
model: string;
usage: { input: number; output: number; cacheRead: number; cacheWrite: number };
costUsd?: number;
at?: number;
}): void {
const record = this.#store.recordObservedUsage;
if (!record) return;
try {
record.call(this.#store, [
{
at: entry.at ?? Date.now(),
provider: entry.provider,
model: entry.model,
requests: 1,
inputTokens: entry.usage.input,
outputTokens: entry.usage.output,
cacheReadTokens: entry.usage.cacheRead,
cacheWriteTokens: entry.usage.cacheWrite,
costUsd: Number.isFinite(entry.costUsd) ? (entry.costUsd ?? 0) : 0,
},
]);
} catch (error) {
this.#usageLogger?.debug("observed usage record failed", {
provider: entry.provider,
error: String(error),
});
}
}
/** Broker host: persist one client's observed-usage report (per-install token burn). */
recordClientUsage(report: ClientUsageReport): boolean {
const record = this.#store.recordClientUsage;
if (!record) return false;
record.call(this.#store, report);
return true;
}
/** Broker host: aggregate recorded per-client usage since `sinceMs`. */
getClientUsageSummary(sinceMs: number): ClientUsageSummary {
return this.#store.getClientUsageSummary?.(sinceMs) ?? { clients: [] };
}
#resolveObservedUsageCredential(provider: Provider, sessionId?: string): UsageCredential | undefined {
const entries = this.#getStoredCredentials(provider);
const sessionCredential = this.#getSessionCredential(provider, sessionId);
@@ -6603,6 +6673,27 @@ export class SqliteAuthCredentialStore implements AuthCredentialStore {
);
CREATE INDEX IF NOT EXISTS idx_usage_cost_history_lookup ON usage_cost_history(provider, account_key, recorded_at);
CREATE INDEX IF NOT EXISTS idx_usage_history_recorded ON usage_history(recorded_at);
CREATE TABLE IF NOT EXISTS clients (
install_id TEXT PRIMARY KEY,
hostname TEXT,
first_seen INTEGER NOT NULL,
last_seen INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS client_usage (
id INTEGER PRIMARY KEY AUTOINCREMENT,
recorded_at INTEGER NOT NULL,
install_id TEXT NOT NULL,
provider TEXT NOT NULL,
model TEXT NOT NULL,
requests INTEGER NOT NULL,
input_tokens INTEGER NOT NULL,
output_tokens INTEGER NOT NULL,
cache_read_tokens INTEGER NOT NULL,
cache_write_tokens INTEGER NOT NULL,
cost_usd REAL NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS idx_client_usage_series ON client_usage(install_id, provider, model, recorded_at);
CREATE INDEX IF NOT EXISTS idx_client_usage_recorded ON client_usage(recorded_at);
`);
if (!this.#authCredentialsTableExists()) {
@@ -7381,6 +7472,113 @@ export class SqliteAuthCredentialStore implements AuthCredentialStore {
}
}
recordClientUsage(report: ClientUsageReport): void {
const now = Date.now();
this.#db
.query(
`INSERT INTO clients (install_id, hostname, first_seen, last_seen) VALUES (?, ?, ?, ?)
ON CONFLICT(install_id) DO UPDATE SET hostname = COALESCE(excluded.hostname, hostname), last_seen = excluded.last_seen`,
)
.run(report.installId, report.hostname ?? null, now, now);
const findBucket = this.#db.query(
`SELECT id FROM client_usage
WHERE install_id = ? AND provider = ? AND model = ? AND recorded_at >= ?
ORDER BY recorded_at DESC LIMIT 1`,
);
const merge = this.#db.query(
`UPDATE client_usage SET recorded_at = ?, requests = requests + ?, input_tokens = input_tokens + ?,
output_tokens = output_tokens + ?, cache_read_tokens = cache_read_tokens + ?,
cache_write_tokens = cache_write_tokens + ?, cost_usd = cost_usd + ? WHERE id = ?`,
);
const insert = this.#db.query(
`INSERT INTO client_usage (recorded_at, install_id, provider, model, requests, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, cost_usd)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
);
for (const entry of report.entries) {
// Merge into the newest row of the same (install, provider, model)
// bucket so 10s client flushes don't accrete one row apiece forever.
const bucketFloor = entry.at - CLIENT_USAGE_BUCKET_MS;
const existing = findBucket.get(report.installId, entry.provider, entry.model, bucketFloor) as {
id: number;
} | null;
if (existing) {
merge.run(
entry.at,
entry.requests,
entry.inputTokens,
entry.outputTokens,
entry.cacheReadTokens,
entry.cacheWriteTokens,
entry.costUsd,
existing.id,
);
continue;
}
insert.run(
entry.at,
report.installId,
entry.provider,
entry.model,
entry.requests,
entry.inputTokens,
entry.outputTokens,
entry.cacheReadTokens,
entry.cacheWriteTokens,
entry.costUsd,
);
}
}
getClientUsageSummary(sinceMs: number): ClientUsageSummary {
const clients = this.#db
.query("SELECT install_id, hostname, first_seen, last_seen FROM clients ORDER BY last_seen DESC")
.all() as Array<{ install_id: string; hostname: string | null; first_seen: number; last_seen: number }>;
const aggregates = this.#db
.query(
`SELECT install_id, provider, SUM(requests) requests, SUM(input_tokens) input_tokens,
SUM(output_tokens) output_tokens, SUM(cache_read_tokens) cache_read_tokens,
SUM(cache_write_tokens) cache_write_tokens, SUM(cost_usd) cost_usd
FROM client_usage WHERE recorded_at >= ? GROUP BY install_id, provider
ORDER BY install_id, SUM(input_tokens + output_tokens + cache_read_tokens + cache_write_tokens) DESC`,
)
.all(sinceMs) as Array<{
install_id: string;
provider: string;
requests: number;
input_tokens: number;
output_tokens: number;
cache_read_tokens: number;
cache_write_tokens: number;
cost_usd: number;
}>;
const providersByInstall = new Map<string, ClientProviderUsage[]>();
for (const row of aggregates) {
let list = providersByInstall.get(row.install_id);
if (!list) {
list = [];
providersByInstall.set(row.install_id, list);
}
list.push({
provider: row.provider,
requests: row.requests,
inputTokens: row.input_tokens,
outputTokens: row.output_tokens,
cacheReadTokens: row.cache_read_tokens,
cacheWriteTokens: row.cache_write_tokens,
costUsd: row.cost_usd,
});
}
return {
clients: clients.map(client => ({
installId: client.install_id,
hostname: client.hostname ?? undefined,
firstSeen: client.first_seen,
lastSeen: client.last_seen,
providers: providersByInstall.get(client.install_id) ?? [],
})),
};
}
// ─── Convenience methods for CLI ────────────────────────────────────────
/**
+55
View File
@@ -185,6 +185,61 @@ export interface UsageCostHistoryQuery {
sinceMs?: number;
}
/**
* Aggregated request usage a client observed for one (provider, model) pair.
* Clients fold every completed request into per-pair buckets and flush them to
* the auth broker on a short cadence, so the broker can attribute token burn
* to the install that produced it.
*/
export interface ObservedUsageEntry {
/** Epoch ms of the newest request folded into this bucket. */
at: number;
provider: Provider;
model: string;
/** Completed requests folded into this bucket. */
requests: number;
inputTokens: number;
outputTokens: number;
cacheReadTokens: number;
cacheWriteTokens: number;
/** Estimated USD cost of the folded requests (0 when unknown). */
costUsd: number;
}
/** One client's observed-usage report, keyed by its stable install id. */
export interface ClientUsageReport {
/** Stable per-machine install id — the client primary key. */
installId: string;
/** Human-readable machine name for display surfaces. */
hostname?: string;
entries: ObservedUsageEntry[];
}
/** Per-provider aggregate of one client's recorded usage. */
export interface ClientProviderUsage {
provider: string;
requests: number;
inputTokens: number;
outputTokens: number;
cacheReadTokens: number;
cacheWriteTokens: number;
costUsd: number;
}
/** One known client with its usage aggregates over the queried window. */
export interface ClientUsageClientSummary {
installId: string;
hostname?: string;
firstSeen: number;
lastSeen: number;
providers: ClientProviderUsage[];
}
/** Aggregated per-client usage recorded by the broker host. */
export interface ClientUsageSummary {
clients: ClientUsageClientSummary[];
}
// ─── Zod schemas (wire-shape validation for the broker `/v1/usage` endpoint) ─
export const usageUnitSchema = type("'percent' | 'tokens' | 'requests' | 'usd' | 'minutes' | 'bytes' | 'unknown'");
@@ -112,6 +112,87 @@ describe("RemoteAuthCredentialStore SSE integration", () => {
expect(remote!.snapshot.credentials[0].id).not.toBe(bId);
});
test("batches observed usage and reports it to the broker as per-install client usage", async () => {
const client = new AuthBrokerClient({ url: handle!.url, token });
remote = new RemoteAuthCredentialStore({ client, observedUsageFlushMs: 25 });
const at = Date.now();
// Two requests for the same (provider, model) must merge into one entry;
// a second model produces its own entry in the same flush.
remote.recordObservedUsage([
{
at,
provider: "anthropic",
model: "claude-x",
requests: 1,
inputTokens: 100,
outputTokens: 50,
cacheReadTokens: 10,
cacheWriteTokens: 5,
costUsd: 0.5,
},
]);
remote.recordObservedUsage([
{
at: at + 1,
provider: "anthropic",
model: "claude-x",
requests: 1,
inputTokens: 200,
outputTokens: 100,
cacheReadTokens: 20,
cacheWriteTokens: 10,
costUsd: 1.0,
},
{
at: at + 2,
provider: "openai-codex",
model: "gpt-y",
requests: 1,
inputTokens: 30,
outputTokens: 15,
cacheReadTokens: 0,
cacheWriteTokens: 0,
costUsd: 0,
},
]);
await waitUntil(() => storage!.getClientUsageSummary(0).clients.length === 1);
const summary = storage!.getClientUsageSummary(0);
const reported = summary.clients[0];
expect(reported.installId.length).toBeGreaterThan(0);
expect(reported.hostname).toBe(os.hostname());
const anthropic = reported.providers.find(p => p.provider === "anthropic");
expect(anthropic).toMatchObject({
requests: 2,
inputTokens: 300,
outputTokens: 150,
cacheReadTokens: 30,
cacheWriteTokens: 15,
});
expect(anthropic?.costUsd).toBeCloseTo(1.5, 10);
expect(reported.providers.find(p => p.provider === "openai-codex")).toMatchObject({
requests: 1,
inputTokens: 30,
});
// A follow-up report — routed through the client-side AuthStorage facade,
// like the coding-agent does per assistant message — merges into the same
// 5-minute bucket row instead of accreting a new row per flush.
const clientStorage = new AuthStorage(remote);
clientStorage.recordObservedUsage({
provider: "anthropic",
model: "claude-x",
at: at + 3,
usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0 },
});
await waitUntil(() => {
const current = storage!.getClientUsageSummary(0).clients[0];
return current?.providers.find(p => p.provider === "anthropic")?.requests === 3;
});
});
test("calls onSnapshot for broker snapshots but not the constructor snapshot", async () => {
const client = new AuthBrokerClient({ url: handle!.url, token });
const initialResult = await client.fetchSnapshot();
+117
View File
@@ -174,6 +174,123 @@ describe("auth-broker wire surface", () => {
await expect(client.refreshCredential(id)).rejects.toThrow();
});
test("GET /v1/usage/history serves recorded snapshots with sinceMs/provider filters and requires auth", async () => {
const hourMs = 60 * 60 * 1000;
const now = Date.now();
// Hours apart so the store's per-bucket dedupe keeps distinct rows.
store!.recordUsageSnapshots!([
{
recordedAt: now - 3 * hourMs,
provider: "anthropic",
accountKey: "acct-a",
email: "a@example.com",
limitId: "5h",
label: "Claude 5 Hour",
windowLabel: "5h",
usedFraction: 0.2,
status: "ok",
},
{
recordedAt: now - hourMs,
provider: "anthropic",
accountKey: "acct-a",
email: "a@example.com",
limitId: "5h",
label: "Claude 5 Hour",
windowLabel: "5h",
usedFraction: 0.9,
status: "warning",
},
{
recordedAt: now - hourMs,
provider: "openai-codex",
accountKey: "acct-b",
limitId: "weekly",
label: "Weekly",
usedFraction: 0.5,
},
]);
const unauthorized = await fetch(`${handle!.url}/v1/usage/history`);
expect(unauthorized.status).toBe(401);
const client = new AuthBrokerClient({ url: handle!.url, token });
const all = await client.fetchUsageHistory();
expect(all.entries).toHaveLength(3);
expect(all.entries[0].recordedAt).toBeLessThanOrEqual(all.entries[1].recordedAt);
const recent = await client.fetchUsageHistory({ sinceMs: now - 2 * hourMs });
expect(recent.entries.map(e => e.usedFraction).sort()).toEqual([0.5, 0.9]);
const anthropicOnly = await client.fetchUsageHistory({ provider: "anthropic" });
expect(anthropicOnly.entries).toHaveLength(2);
expect(anthropicOnly.entries.every(e => e.provider === "anthropic")).toBe(true);
expect(anthropicOnly.entries[1]).toMatchObject({
accountKey: "acct-a",
email: "a@example.com",
windowLabel: "5h",
usedFraction: 0.9,
status: "warning",
});
});
test("POST /v1/usage/observed persists per-client usage served by GET /v1/usage/clients", async () => {
const unauthorized = await fetch(`${handle!.url}/v1/usage/observed`, { method: "POST" });
expect(unauthorized.status).toBe(401);
const client = new AuthBrokerClient({ url: handle!.url, token });
const now = Date.now();
const entry = {
at: now,
provider: "anthropic",
model: "claude-x",
requests: 2,
inputTokens: 1000,
outputTokens: 400,
cacheReadTokens: 50,
cacheWriteTokens: 25,
costUsd: 3.25,
};
const ack = await client.reportClientUsage({ installId: "install-1", hostname: "mbp.local", entries: [entry] });
expect(ack.ok).toBe(true);
await client.reportClientUsage({
installId: "install-2",
entries: [{ ...entry, provider: "openai-codex", model: "gpt-y", requests: 1, costUsd: 0 }],
});
const summary = await client.fetchClientUsageSummary();
expect(summary.clients).toHaveLength(2);
const first = summary.clients.find(c => c.installId === "install-1");
expect(first).toMatchObject({ hostname: "mbp.local" });
expect(first?.lastSeen).toBeGreaterThan(0);
expect(first?.providers).toEqual([
{
provider: "anthropic",
requests: 2,
inputTokens: 1000,
outputTokens: 400,
cacheReadTokens: 50,
cacheWriteTokens: 25,
costUsd: 3.25,
},
]);
const second = summary.clients.find(c => c.installId === "install-2");
expect(second?.hostname).toBeUndefined();
expect(second?.providers[0]).toMatchObject({ provider: "openai-codex", requests: 1 });
// sinceMs beyond the recorded timestamps returns clients with no aggregates.
const future = await client.fetchClientUsageSummary({ sinceMs: now + 60_000 });
expect(future.clients.every(c => c.providers.length === 0)).toBe(true);
// Malformed body is rejected by schema validation.
const bad = await fetch(`${handle!.url}/v1/usage/observed`, {
method: "POST",
headers: { Authorization: `Bearer ${token}`, "Content-Type": "application/json" },
body: JSON.stringify({ installId: "x", entries: [{ at: "not-a-number" }] }),
});
expect(bad.status).toBe(400);
});
test("Unknown route returns 404", async () => {
const res = await fetch(`${handle!.url}/v1/nope`, {
headers: { Authorization: `Bearer ${token}` },
+1
View File
@@ -6,6 +6,7 @@
- Added in-process moreutils-style shell builtins to the bash tool's embedded shell: `ts` (timestamp lines; `-i`/`-s` elapsed modes, `-m` monotonic clock, `-r` relative rewriting of RFC3339/syslog timestamps, `%.S`/`%.s`/`%.T` subsecond extensions), `sponge` (soak stdin fully before atomically writing the target, so `foo file | ... | sponge file` works; `-a` appends), `ifne` (run a command only when stdin is non-empty; `-n` inverts and passes non-empty stdin through), `isutf8` (streaming UTF-8 validation with line/char/byte diagnostics; `-q`, `-l`, `-i`), `combine` (boolean `and`/`not`/`or`/`xor` on the lines of two files, `-` for stdin), and `errno` (errno name/number/description lookup with `-l` list and `-s` search; unix only). Like the uutils-backed builtins, they run in-process against the command's own stdio, resolve paths against the shell working directory, honor cancellation, and are disabled by `PI_DISABLE_UUTILS_BUILTINS`.
- The bash tool prompt now lists the available shell builtins (`mkdir` through `jq`, `rm`/`mv`/`ln`, and the moreutils set) so the model relies on them without existence checks; the line is dropped when `PI_DISABLE_UUTILS_BUILTINS` disables the builtins and omits unix-only `errno` on Windows.
- Sessions using a broker-backed auth store now report each completed request's token usage and cost to the auth broker (batched, 10s cadence) so the broker can track actual token burn per client install
### Changed
@@ -2406,6 +2406,20 @@ export class AgentSession {
baseUrl: this.#modelRegistry.getProviderBaseUrl?.(assistantMsg.provider),
});
}
// Broker deployments: report this request's burn so the broker can
// attribute token usage per install. No-op with a local auth store.
this.#modelRegistry.authStorage.recordObservedUsage({
provider: assistantMsg.provider,
model: assistantMsg.model,
at: assistantMsg.timestamp,
usage: {
input: assistantMsg.usage.input,
output: assistantMsg.usage.output,
cacheRead: assistantMsg.usage.cacheRead,
cacheWrite: assistantMsg.usage.cacheWrite,
},
costUsd: assistantMsg.usage.cost.total,
});
}
if (event.message.role === "toolResult") {
const { toolName, toolCallId, isError, content } = event.message;
+1 -1
View File
@@ -4,7 +4,7 @@
### Added
- Added a Providers dashboard section: per-provider totals, stacked token/cost burn over time, peak-burn-hours histogram, subscription-window insights (windows burned, estimated tokens per window, peak concurrent utilization, ideal account count, exhaustion events), and latest window utilization per account — window analytics are derived from the auth store's recorded usage-limit snapshots.
- Added a Providers dashboard section: per-provider totals, stacked token/cost burn over time, peak-burn-hours histogram, subscription-window insights (windows burned, estimated tokens per window, peak concurrent utilization, ideal account count, exhaustion events), and latest window utilization per account — window analytics read the auth broker's `/v1/usage/history` when a broker is configured (falling back to the local agent DB), since broker deployments record usage history on the broker host
## [17.1.0] - 2026-07-24
+3 -2
View File
@@ -46,7 +46,7 @@ import type {
RequestDetails,
ToolDashboardStats,
} from "./types";
import { computeUsageWindowStats, readUsageSnapshots } from "./usage-windows";
import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows";
/**
* Apply a freshly parsed result to the database. Runs entirely on the
@@ -526,7 +526,8 @@ export async function getProviderDashboardStats(range?: string | null): Promise<
const { modelSeriesDays, modelSeriesBucketMs, cutoff } = getTimeRangeConfig(range);
const providers = getStatsByProvider(cutoff ?? undefined);
const tokensByProvider = new Map(providers.map(p => [p.provider, p.totalTokens]));
const { usageSeries, windowInsights } = computeUsageWindowStats(readUsageSnapshots(cutoff ?? 0), tokensByProvider);
const snapshots = await fetchUsageSnapshots(cutoff ?? 0);
const { usageSeries, windowInsights } = computeUsageWindowStats(snapshots, tokensByProvider);
return {
providers,
hourly: getProviderHourlyBurn(cutoff ?? undefined),
+33
View File
@@ -13,6 +13,7 @@
* dashboard must keep working for API-key-only setups that never record usage.
*/
import { Database } from "bun:sqlite";
import { AuthBrokerClient, resolveAuthBrokerConfig } from "@oh-my-pi/pi-ai/auth-broker";
import { getAgentDbPath, logger } from "@oh-my-pi/pi-utils";
import type { ProviderWindowInsight, UsageWindowPoint, UsageWindowSeries } from "./shared-types";
@@ -98,6 +99,38 @@ 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.
*/
export async function fetchUsageSnapshots(sinceMs: number): Promise<UsageSnapshotRow[]> {
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,
}));
}
} catch (err) {
logger.debug("broker usage history unavailable, falling back to local", { error: String(err) });
}
return readUsageSnapshots(sinceMs);
}
/** True when a snapshot reports an exhausted window, by status or by fraction. */
function isExhausted(fraction: number | null, status: string | null): boolean {
if (status === "exhausted") return true;