From 3676b5d97f24ef04ba0accfb5cc179b2f1b5a18c Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 24 Jul 2026 15:20:46 +0200 Subject: [PATCH] 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. --- packages/ai/CHANGELOG.md | 5 + packages/ai/src/auth-broker/client.ts | 39 ++++ packages/ai/src/auth-broker/remote-store.ts | 79 ++++++- packages/ai/src/auth-broker/server.ts | 40 ++++ packages/ai/src/auth-broker/types.ts | 26 ++- packages/ai/src/auth-broker/wire-schemas.ts | 71 +++++++ packages/ai/src/auth-storage.ts | 198 ++++++++++++++++++ packages/ai/src/usage.ts | 55 +++++ .../ai/test/auth-broker-remote-store.test.ts | 81 +++++++ packages/ai/test/auth-broker-wire.test.ts | 117 +++++++++++ packages/coding-agent/CHANGELOG.md | 1 + .../coding-agent/src/session/agent-session.ts | 14 ++ packages/stats/CHANGELOG.md | 2 +- packages/stats/src/aggregator.ts | 5 +- packages/stats/src/usage-windows.ts | 33 +++ 15 files changed, 759 insertions(+), 7 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index e87a735dd..6d8634400 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -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 diff --git a/packages/ai/src/auth-broker/client.ts b/packages/ai/src/auth-broker/client.ts index c90d83f4f..aa03d1ec9 100644 --- a/packages/ai/src/auth-broker/client.ts +++ b/packages/ai/src/auth-broker/client.ts @@ -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("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 { + 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("GET", path, { schema: usageHistoryResponseSchema, signal }); + } + + /** Report this client's batched observed request usage for per-install burn tracking. */ + reportClientUsage(report: ClientUsageReportRequest, signal?: AbortSignal): Promise { + return this.#request("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 { + 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("GET", path, { + schema: clientUsageSummaryResponseSchema, + signal, + }); + } + notifyUsageStale(signal?: AbortSignal): Promise { return this.#request("POST", "/v1/usage/stale", { schema: usageStaleResponseSchema, diff --git a/packages/ai/src/auth-broker/remote-store.ts b/packages/ai/src/auth-broker/remote-store.ts index 7eb249df1..7e5b27ddf 100644 --- a/packages/ai/src/auth-broker/remote-store.ts +++ b/packages/ai/src/auth-broker/remote-store.ts @@ -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(); + #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 { + 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(); } diff --git a/packages/ai/src/auth-broker/server.ts b/packages/ai/src/auth-broker/server.ts index ce8bdf06f..367479e11 100644 --- a/packages/ai/src/auth-broker/server.ts +++ b/packages/ai/src/auth-broker/server.ts @@ -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?.(); diff --git a/packages/ai/src/auth-broker/types.ts b/packages/ai/src/auth-broker/types.ts index d571215cf..1164022d1 100644 --- a/packages/ai/src/auth-broker/types.ts +++ b/packages/ai/src/auth-broker/types.ts @@ -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; diff --git a/packages/ai/src/auth-broker/wire-schemas.ts b/packages/ai/src/auth-broker/wire-schemas.ts index d13284333..d672941f5 100644 --- a/packages/ai/src/auth-broker/wire-schemas.ts +++ b/packages/ai/src/auth-broker/wire-schemas.ts @@ -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({ diff --git a/packages/ai/src/auth-storage.ts b/packages/ai/src/auth-storage.ts index dba774da4..5f87d5365 100644 --- a/packages/ai/src/auth-storage.ts +++ b/packages/ai/src/auth-storage.ts @@ -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(); + 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 ──────────────────────────────────────── /** diff --git a/packages/ai/src/usage.ts b/packages/ai/src/usage.ts index 80ba52e35..fe71d0096 100644 --- a/packages/ai/src/usage.ts +++ b/packages/ai/src/usage.ts @@ -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'"); diff --git a/packages/ai/test/auth-broker-remote-store.test.ts b/packages/ai/test/auth-broker-remote-store.test.ts index ae7b11502..484e7c772 100644 --- a/packages/ai/test/auth-broker-remote-store.test.ts +++ b/packages/ai/test/auth-broker-remote-store.test.ts @@ -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(); diff --git a/packages/ai/test/auth-broker-wire.test.ts b/packages/ai/test/auth-broker-wire.test.ts index 256c3c73d..989d15d94 100644 --- a/packages/ai/test/auth-broker-wire.test.ts +++ b/packages/ai/test/auth-broker-wire.test.ts @@ -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}` }, diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index defbed870..35cc33fed 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 0d9ef27c4..9c6e3caff 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -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; diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index aeefcc4ed..05b5e5953 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -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 diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index 39cadc2b4..2c6628db6 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -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), diff --git a/packages/stats/src/usage-windows.ts b/packages/stats/src/usage-windows.ts index 6b7a297a7..a1cbccadd 100644 --- a/packages/stats/src/usage-windows.ts +++ b/packages/stats/src/usage-windows.ts @@ -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 { + 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;