diff --git a/cli/src/codex/codexLocalLauncher.test.ts b/cli/src/codex/codexLocalLauncher.test.ts index e2728de4..10e286fa 100644 --- a/cli/src/codex/codexLocalLauncher.test.ts +++ b/cli/src/codex/codexLocalLauncher.test.ts @@ -645,7 +645,14 @@ describe('codexLocalLauncher', () => { [ JSON.stringify({ type: 'session_meta', payload: { id: 'codex-thread-import' } }), JSON.stringify({ type: 'event_msg', payload: { type: 'user_message', message: 'old imported prompt' } }), - JSON.stringify({ type: 'event_msg', payload: { type: 'agent_message', message: 'old imported message' } }) + JSON.stringify({ type: 'event_msg', payload: { type: 'agent_message', message: 'old imported message' } }), + JSON.stringify({ + type: 'event_msg', + payload: { + type: 'token_count', + info: { total_token_usage: { input_tokens: 100, output_tokens: 10 } } + } + }) ].join('\n') + '\n' ); @@ -694,7 +701,14 @@ describe('codexLocalLauncher', () => { transcriptPath, [ JSON.stringify({ type: 'event_msg', payload: { type: 'user_message', message: 'new local prompt' } }), - JSON.stringify({ type: 'event_msg', payload: { type: 'agent_message', message: 'new local response' } }) + JSON.stringify({ type: 'event_msg', payload: { type: 'agent_message', message: 'new local response' } }), + JSON.stringify({ + type: 'event_msg', + payload: { + type: 'token_count', + info: { total_token_usage: { input_tokens: 120, output_tokens: 12 } } + } + }) ].join('\n') + '\n' ); await wait(700); @@ -713,6 +727,17 @@ describe('codexLocalLauncher', () => { message: 'new local response', id: expect.any(String) }); + const tokenMessages = agentMessages.filter((message) => ( + message as { type?: string } + ).type === 'token_count') as Array>; + expect(tokenMessages).toHaveLength(2); + expect(tokenMessages[0]).toMatchObject({ hapiUsageScope: 'imported-history' }); + expect(tokenMessages[0]).not.toHaveProperty('thread_id'); + expect(tokenMessages[1]).toMatchObject({ + threadId: 'codex-thread-import', + thread_id: 'codex-thread-import', + hapiUsageScope: 'managed' + }); }); it('replays semantic chat and tool events once and keeps a same-turn preface before its plan', async () => { diff --git a/cli/src/codex/codexLocalLauncher.ts b/cli/src/codex/codexLocalLauncher.ts index f09da448..6c5a21b8 100644 --- a/cli/src/codex/codexLocalLauncher.ts +++ b/cli/src/codex/codexLocalLauncher.ts @@ -202,7 +202,7 @@ export async function codexLocalLauncher(session: CodexSession): Promise<'switch } session.onSessionFound(sessionId); }, - onEvent: (event) => { + onEvent: (event, context) => { const observedReasoningEffort = extractTurnContextReasoningEffort(event); if (observedReasoningEffort !== undefined) { session.setModelReasoningEffort(observedReasoningEffort); @@ -242,7 +242,19 @@ export async function codexLocalLauncher(session: CodexSession): Promise<'switch flushPendingExecWrapper(message.callId, message); } } else { - session.sendAgentMessage(message); + const scopedMessage = message.type !== 'token_count' + ? message + : context.replayedHistory + ? { ...message, hapiUsageScope: 'imported-history' } + : primarySessionId + ? { + ...message, + threadId: primarySessionId, + thread_id: primarySessionId, + hapiUsageScope: 'managed' + } + : message; + session.sendAgentMessage(scopedMessage); } } if (converted?.finishedTurnId) { diff --git a/cli/src/codex/utils/codexSessionScanner.test.ts b/cli/src/codex/utils/codexSessionScanner.test.ts index 06c205c1..51362f03 100644 --- a/cli/src/codex/utils/codexSessionScanner.test.ts +++ b/cli/src/codex/utils/codexSessionScanner.test.ts @@ -104,16 +104,28 @@ describe('codexSessionScanner', () => { ].join('\n') + '\n' ); + const replayFlags: boolean[] = []; scanner = await createCodexSessionScanner({ transcriptPath, replayExistingHistory: true, - onEvent: (event) => events.push(event) + onEvent: (event, context) => { + events.push(event); + replayFlags.push(context.replayedHistory); + } }); await wait(300); expect(events).toHaveLength(2); expect(events[0]?.type).toBe('session_meta'); expect(events[1]?.payload).toEqual({ type: 'agent_message', message: 'old' }); + expect(replayFlags).toEqual([true, true]); + + await appendFile( + transcriptPath, + JSON.stringify({ type: 'event_msg', payload: { type: 'agent_message', message: 'new' } }) + '\n' + ); + await scanner.flush(); + expect(replayFlags).toEqual([true, true, false]); }); it('reports session id from the transcript metadata', async () => { diff --git a/cli/src/codex/utils/codexSessionScanner.ts b/cli/src/codex/utils/codexSessionScanner.ts index 4c98af00..73c9c695 100644 --- a/cli/src/codex/utils/codexSessionScanner.ts +++ b/cli/src/codex/utils/codexSessionScanner.ts @@ -5,7 +5,7 @@ import type { CodexSessionEvent } from './codexEventConverter'; interface CodexSessionScannerOptions { transcriptPath: string | null; - onEvent: (event: CodexSessionEvent) => void; + onEvent: (event: CodexSessionEvent, context: { replayedHistory: boolean }) => void; onSessionId?: (sessionId: string) => void; replayExistingHistory?: boolean; } @@ -35,7 +35,7 @@ export async function createCodexSessionScanner(opts: CodexSessionScannerOptions class CodexSessionScannerImpl extends BaseSessionScanner { private transcriptPath: string | null; - private readonly onEvent: (event: CodexSessionEvent) => void; + private readonly onEvent: (event: CodexSessionEvent, context: { replayedHistory: boolean }) => void; private readonly onSessionId?: (sessionId: string) => void; private readonly fileEpochByPath = new Map(); private readonly fileStateByPath = new Map { nextLineIndex: number; }>(); private replayExistingHistoryOnNextAttach: boolean; + private replayingExistingHistory = false; private observedSessionId: string | null = null; constructor(opts: CodexSessionScannerOptions) { @@ -92,8 +93,13 @@ class CodexSessionScannerImpl extends BaseSessionScanner { } protected async handleFileScan(stats: SessionFileScanStats): Promise { - for (const event of stats.events) { - this.onEvent(event); + const replayedHistory = this.replayingExistingHistory; + try { + for (const event of stats.events) { + this.onEvent(event, { replayedHistory }); + } + } finally { + this.replayingExistingHistory = false; } if (stats.newCount > 0) { logger.debug(`[codex-session-scanner] ${stats.newCount} new events from ${stats.filePath}`); @@ -106,9 +112,11 @@ class CodexSessionScannerImpl extends BaseSessionScanner { // 中文注释:导入既有 Codex thread 时,首次挂接 transcript 不能先 prime 到 EOF, // 否则 Hapi 只会看到后续增量,客户端里已经存在的最新消息会被跳过。 this.replayExistingHistoryOnNextAttach = false; + this.replayingExistingHistory = true; return; } + this.replayingExistingHistory = false; await this.primeTranscript(filePath); } diff --git a/hub/README.md b/hub/README.md index c00655fa..8668c852 100644 --- a/hub/README.md +++ b/hub/README.md @@ -115,6 +115,10 @@ See `src/web/routes/` for all endpoints. - `POST /api/machines/:id/spawn` - Spawn new session on machine. - `POST /api/machines/:id/paths/exists` - Check if path exists. +### Usage (`src/web/routes/usage.ts`) + +- `GET /api/usage/summary` - Get cache-aware token usage for the owner namespace (`range=7d|30d|all`). + ### Git/Files (`src/web/routes/git.ts`) - `GET /api/sessions/:id/git-status` - Git status. diff --git a/hub/src/store/index.ts b/hub/src/store/index.ts index 8512ec6e..0f713e17 100644 --- a/hub/src/store/index.ts +++ b/hub/src/store/index.ts @@ -9,6 +9,7 @@ import { FcmStore } from './fcmStore' import { ScratchlistStore } from './scratchlistStore' import { SessionStore } from './sessionStore' import { UserStore } from './userStore' +import { UsageStore } from './usageStore' export type { StoredMachine, @@ -28,8 +29,9 @@ export { FcmStore } from './fcmStore' export { ScratchlistStore } from './scratchlistStore' export { SessionStore } from './sessionStore' export { UserStore } from './userStore' +export { UsageStore } from './usageStore' -const SCHEMA_VERSION: number = 16 +const SCHEMA_VERSION: number = 19 const REQUIRED_TABLES = [ 'sessions', 'machines', @@ -38,7 +40,9 @@ const REQUIRED_TABLES = [ 'users', 'push_subscriptions', 'fcm_devices', - 'session_scratchlist' + 'session_scratchlist', + 'usage_events', + 'usage_scan_state' ] as const export class Store { @@ -53,6 +57,7 @@ export class Store { readonly push: PushStore readonly fcm: FcmStore readonly scratchlist: ScratchlistStore + readonly usage: UsageStore /** * Filesystem path of the underlying SQLite database, or ':memory:' for @@ -105,6 +110,7 @@ export class Store { this.push = new PushStore(this.db) this.fcm = new FcmStore(this.db) this.scratchlist = new ScratchlistStore(this.db) + this.usage = new UsageStore(this.db) } /** @@ -173,6 +179,9 @@ export class Store { 13: () => this.migrateFromV13ToV14(), 14: () => this.migrateFromV14ToV15(), 15: () => this.migrateFromV15ToV16(), + 16: () => this.migrateFromV16ToV17(), + 17: () => this.migrateFromV17ToV18(), + 18: () => this.migrateFromV18ToV19(), }) if (currentVersion === 0) { @@ -333,6 +342,37 @@ export class Store { ); CREATE INDEX IF NOT EXISTS idx_session_scratchlist_session_created ON session_scratchlist(session_id, created_at DESC); + + CREATE TABLE IF NOT EXISTS usage_events ( + session_id TEXT NOT NULL, + source_key TEXT NOT NULL, + source_seq INTEGER NOT NULL, + created_at INTEGER NOT NULL, + agent TEXT NOT NULL, + model TEXT, + kind TEXT NOT NULL CHECK (kind IN ('delta', 'cumulative')), + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + cache_read_tokens INTEGER NOT NULL DEFAULT 0, + cache_creation_tokens INTEGER NOT NULL DEFAULT 0, + last_input_tokens INTEGER, + last_output_tokens INTEGER, + last_cache_read_tokens INTEGER, + last_cache_creation_tokens INTEGER, + PRIMARY KEY (session_id, source_key), + FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE + ); + CREATE INDEX IF NOT EXISTS idx_usage_events_session_created + ON usage_events(session_id, created_at, source_seq); + CREATE INDEX IF NOT EXISTS idx_usage_events_created + ON usage_events(created_at); + + CREATE TABLE IF NOT EXISTS usage_scan_state ( + session_id TEXT PRIMARY KEY, + message_epoch INTEGER NOT NULL DEFAULT 0, + last_seq INTEGER NOT NULL DEFAULT 0, + FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE + ); `) } @@ -605,6 +645,65 @@ export class Store { */ private migrateFromV15ToV16(): void {} + private migrateFromV16ToV17(): void { + this.db.exec(` + CREATE TABLE IF NOT EXISTS usage_events ( + session_id TEXT NOT NULL, + source_key TEXT NOT NULL, + source_seq INTEGER NOT NULL, + created_at INTEGER NOT NULL, + agent TEXT NOT NULL, + model TEXT, + kind TEXT NOT NULL CHECK (kind IN ('delta', 'cumulative')), + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + cache_read_tokens INTEGER NOT NULL DEFAULT 0, + cache_creation_tokens INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (session_id, source_key), + FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE + ); + CREATE INDEX IF NOT EXISTS idx_usage_events_session_created + ON usage_events(session_id, created_at, source_seq); + CREATE INDEX IF NOT EXISTS idx_usage_events_created + ON usage_events(created_at); + `) + } + + private migrateFromV17ToV18(): void { + // Usage events are a rebuildable index; v18 changes their source key + // and baseline semantics, so stale rows must not be mixed with new ones. + const columns = new Set( + (this.db.prepare('PRAGMA table_info(usage_events)').all() as Array<{ name: string }>) + .map((column) => column.name) + ) + for (const name of [ + 'last_input_tokens', + 'last_output_tokens', + 'last_cache_read_tokens', + 'last_cache_creation_tokens' + ]) { + if (!columns.has(name)) { + this.db.exec(`ALTER TABLE usage_events ADD COLUMN ${name} INTEGER`) + } + } + this.db.exec('DELETE FROM usage_events') + } + + private migrateFromV18ToV19(): void { + // Cumulative event keys are stable in v19, so repeated transcript + // imports collapse to one snapshot. Rebuild the derived index once. + this.db.exec(` + CREATE TABLE IF NOT EXISTS usage_scan_state ( + session_id TEXT PRIMARY KEY, + message_epoch INTEGER NOT NULL DEFAULT 0, + last_seq INTEGER NOT NULL DEFAULT 0, + FOREIGN KEY (session_id) REFERENCES sessions(id) ON DELETE CASCADE + ); + DELETE FROM usage_events; + DELETE FROM usage_scan_state; + `) + } + private getSessionColumnNames(): Set { const rows = this.db.prepare('PRAGMA table_info(sessions)').all() as Array<{ name: string }> return new Set(rows.map((row) => row.name)) diff --git a/hub/src/store/messageStore.ts b/hub/src/store/messageStore.ts index 643d6d1c..9c5b3d74 100644 --- a/hub/src/store/messageStore.ts +++ b/hub/src/store/messageStore.ts @@ -27,6 +27,7 @@ import { copyMessageToSession as copyStoredMessageToSession, copyMessagesToSession as copyStoredMessagesToSession, getAllMessages, + getMessagesAfterSeq, truncateMessagesFromLocalId, type CancelQueuedMessageResult, type LookupQueuedMessageResult, @@ -64,6 +65,10 @@ export class MessageStore { return getAllMessages(this.db, sessionId) } + getMessagesAfterSeq(sessionId: string, afterSeq: number): StoredMessage[] { + return getMessagesAfterSeq(this.db, sessionId, afterSeq) + } + getMessages(sessionId: string, limit: number = 200): StoredMessage[] { return getMessages(this.db, sessionId, limit) } diff --git a/hub/src/store/messages.ts b/hub/src/store/messages.ts index 14981a5c..71c98a4b 100644 --- a/hub/src/store/messages.ts +++ b/hub/src/store/messages.ts @@ -240,6 +240,18 @@ export function getAllMessages( return rows.map(toStoredMessage) } +export function getMessagesAfterSeq( + db: Database, + sessionId: string, + afterSeq: number +): StoredMessage[] { + const rows = db.prepare( + 'SELECT * FROM messages WHERE session_id = ? AND seq > ? ORDER BY seq ASC' + ).all(sessionId, afterSeq) as DbMessageRow[] + + return rows.map(toStoredMessage) +} + export function getFirstMessages( db: Database, sessionId: string, diff --git a/hub/src/store/migration-v13.test.ts b/hub/src/store/migration-v13.test.ts index cc0d6e79..d53af0aa 100644 --- a/hub/src/store/migration-v13.test.ts +++ b/hub/src/store/migration-v13.test.ts @@ -10,7 +10,7 @@ describe('Store V12/V13→V14 schema reconciliation', () => { const store = new Store(':memory:') expect(tableExists(store, 'message_epochs')).toBe(true) expect(tableExists(store, 'session_scratchlist')).toBe(true) - expect(getUserVersion(store)).toBe(16) + expect(getUserVersion(store)).toBe(19) store.close() }) @@ -34,7 +34,7 @@ describe('Store V12/V13→V14 schema reconciliation', () => { store = new Store(dbPath) expect(tableExists(store, 'message_epochs')).toBe(true) expect(tableExists(store, 'session_scratchlist')).toBe(true) - expect(getUserVersion(store)).toBe(16) + expect(getUserVersion(store)).toBe(19) expect(store.messages.getMessageEpoch('session-1')).toBe(0) expect(store.messages.getMessages('session-1')).toHaveLength(1) } finally { @@ -72,7 +72,7 @@ describe('Store V12/V13→V14 schema reconciliation', () => { store = new Store(dbPath) expect(tableExists(store, 'message_epochs')).toBe(true) expect(tableExists(store, 'session_scratchlist')).toBe(true) - expect(getUserVersion(store)).toBe(16) + expect(getUserVersion(store)).toBe(19) expect(store.messages.getMessages('session-1')).toHaveLength(1) } finally { store?.close() diff --git a/hub/src/store/migration-v15.test.ts b/hub/src/store/migration-v15.test.ts index 1231d1cf..de05ed27 100644 --- a/hub/src/store/migration-v15.test.ts +++ b/hub/src/store/migration-v15.test.ts @@ -17,7 +17,9 @@ describe('Store V14→V15 migration: scratchlist attachments column', () => { const store = new Store(':memory:') const cols = getColumns(store, 'session_scratchlist') expect(cols).toContain('attachments') - expect(getUserVersion(store)).toBe(16) + expect(getColumns(store, 'usage_events')).toContain('last_input_tokens') + expect(getColumns(store, 'usage_scan_state')).toContain('last_seq') + expect(getUserVersion(store)).toBe(19) store.close() }) @@ -36,7 +38,9 @@ describe('Store V14→V15 migration: scratchlist attachments column', () => { store = new Store(dbPath) const cols = getColumns(store, 'session_scratchlist') expect(cols).toContain('attachments') - expect(getUserVersion(store)).toBe(16) + expect(getColumns(store, 'usage_events')).toContain('last_input_tokens') + expect(getColumns(store, 'usage_scan_state')).toContain('last_seq') + expect(getUserVersion(store)).toBe(19) } finally { store?.close() rmSync(dir, { recursive: true, force: true }) @@ -56,7 +60,7 @@ describe('Store V14→V15 migration: scratchlist attachments column', () => { store2 = new Store(dbPath) const cols2 = getColumns(store2, 'session_scratchlist') expect(cols2).toEqual(cols1) - expect(getUserVersion(store2)).toBe(16) + expect(getUserVersion(store2)).toBe(19) } finally { store2?.close() store1?.close() diff --git a/hub/src/store/migration-v18.test.ts b/hub/src/store/migration-v18.test.ts new file mode 100644 index 00000000..473e8bf8 --- /dev/null +++ b/hub/src/store/migration-v18.test.ts @@ -0,0 +1,48 @@ +import { describe, expect, it } from 'bun:test' +import { Database } from 'bun:sqlite' +import { mkdtempSync, rmSync } from 'node:fs' +import { join } from 'node:path' +import { tmpdir } from 'node:os' +import { Store } from './index' + +describe('Store V18->V19 migration: usage scan state', () => { + it('adds the usage_scan_state table to a V18 database', () => { + const directory = mkdtempSync(join(tmpdir(), 'hapi-migration-v18-to-v19-')) + const dbPath = join(directory, 'test.db') + let store: Store | undefined + try { + store = new Store(dbPath) + store.close() + store = undefined + + const db = new Database(dbPath, { create: true, readwrite: true, strict: true }) + db.exec(` + INSERT INTO sessions (id, created_at, updated_at) + VALUES ('session-1', 1, 1); + INSERT INTO usage_events ( + session_id, source_key, source_seq, created_at, agent, kind + ) VALUES ( + 'session-1', 'old-cumulative-key', 1, 1, 'codex', 'cumulative' + ); + DROP TABLE usage_scan_state; + PRAGMA user_version = 18; + `) + db.close() + + store = new Store(dbPath) + const internalDb = (store as unknown as { db: Database }).db + const table = internalDb.prepare( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'usage_scan_state'" + ).get() as { name: string } | null + const version = internalDb.prepare('PRAGMA user_version').get() as { user_version: number } + const usageRows = internalDb.prepare('SELECT COUNT(*) AS count FROM usage_events').get() as { count: number } + + expect(table?.name).toBe('usage_scan_state') + expect(version.user_version).toBe(19) + expect(usageRows.count).toBe(0) + } finally { + store?.close() + rmSync(directory, { recursive: true, force: true }) + } + }) +}) diff --git a/hub/src/store/usage.ts b/hub/src/store/usage.ts new file mode 100644 index 00000000..8c49ca33 --- /dev/null +++ b/hub/src/store/usage.ts @@ -0,0 +1,257 @@ +import type { Database } from 'bun:sqlite' + +export type UsageEventKind = 'delta' | 'cumulative' + +export type UsageEvent = { + sessionId: string + sourceKey: string + sourceSeq: number + createdAt: number + agent: string + model: string | null + kind: UsageEventKind + inputTokens: number + outputTokens: number + cacheReadTokens: number + cacheCreationTokens: number + lastInputTokens: number | null + lastOutputTokens: number | null + lastCacheReadTokens: number | null + lastCacheCreationTokens: number | null +} + +export type UsageScanState = { + messageEpoch: number + lastSeq: number +} + +type UsageEventRow = { + session_id: string + source_key: string + source_seq: number + created_at: number + agent: string + model: string | null + kind: UsageEventKind + input_tokens: number + output_tokens: number + cache_read_tokens: number + cache_creation_tokens: number + last_input_tokens: number | null + last_output_tokens: number | null + last_cache_read_tokens: number | null + last_cache_creation_tokens: number | null +} + +function toUsageEvent(row: UsageEventRow): UsageEvent { + return { + sessionId: row.session_id, + sourceKey: row.source_key, + sourceSeq: row.source_seq, + createdAt: row.created_at, + agent: row.agent, + model: row.model, + kind: row.kind, + inputTokens: row.input_tokens, + outputTokens: row.output_tokens, + cacheReadTokens: row.cache_read_tokens, + cacheCreationTokens: row.cache_creation_tokens, + lastInputTokens: row.last_input_tokens, + lastOutputTokens: row.last_output_tokens, + lastCacheReadTokens: row.last_cache_read_tokens, + lastCacheCreationTokens: row.last_cache_creation_tokens + } +} + +export function recordUsageScan( + db: Database, + sessionId: string, + messageEpoch: number, + lastSeq: number, + events: UsageEvent[], + replaceEvents: boolean +): void { + db.transaction(() => { + if (replaceEvents) { + db.prepare('DELETE FROM usage_events WHERE session_id = ?').run(sessionId) + } + + if (events.length > 0) { + const statement = db.prepare(` + INSERT INTO usage_events ( + session_id, + source_key, + source_seq, + created_at, + agent, + model, + kind, + input_tokens, + output_tokens, + cache_read_tokens, + cache_creation_tokens, + last_input_tokens, + last_output_tokens, + last_cache_read_tokens, + last_cache_creation_tokens + ) VALUES ( + @session_id, + @source_key, + @source_seq, + @created_at, + @agent, + @model, + @kind, + @input_tokens, + @output_tokens, + @cache_read_tokens, + @cache_creation_tokens, + @last_input_tokens, + @last_output_tokens, + @last_cache_read_tokens, + @last_cache_creation_tokens + ) + ON CONFLICT(session_id, source_key) + DO UPDATE SET + source_seq = excluded.source_seq, + created_at = excluded.created_at, + agent = excluded.agent, + model = excluded.model, + kind = excluded.kind, + input_tokens = excluded.input_tokens, + output_tokens = excluded.output_tokens, + cache_read_tokens = excluded.cache_read_tokens, + cache_creation_tokens = excluded.cache_creation_tokens, + last_input_tokens = excluded.last_input_tokens, + last_output_tokens = excluded.last_output_tokens, + last_cache_read_tokens = excluded.last_cache_read_tokens, + last_cache_creation_tokens = excluded.last_cache_creation_tokens + WHERE usage_events.kind = 'delta' + `) + + for (const event of events) { + statement.run({ + session_id: event.sessionId, + source_key: event.sourceKey, + source_seq: event.sourceSeq, + created_at: event.createdAt, + agent: event.agent, + model: event.model, + kind: event.kind, + input_tokens: event.inputTokens, + output_tokens: event.outputTokens, + cache_read_tokens: event.cacheReadTokens, + cache_creation_tokens: event.cacheCreationTokens, + last_input_tokens: event.lastInputTokens, + last_output_tokens: event.lastOutputTokens, + last_cache_read_tokens: event.lastCacheReadTokens, + last_cache_creation_tokens: event.lastCacheCreationTokens + }) + } + } + + db.prepare(` + INSERT INTO usage_scan_state (session_id, message_epoch, last_seq) + VALUES (?, ?, ?) + ON CONFLICT(session_id) DO UPDATE SET + message_epoch = excluded.message_epoch, + last_seq = CASE + WHEN usage_scan_state.message_epoch = excluded.message_epoch + THEN MAX(usage_scan_state.last_seq, excluded.last_seq) + ELSE excluded.last_seq + END + WHERE excluded.message_epoch >= usage_scan_state.message_epoch + `).run(sessionId, messageEpoch, lastSeq) + })() +} + +export function getUsageEvents(db: Database, sessionIds: string[]): UsageEvent[] { + if (sessionIds.length === 0) return [] + + const placeholders = sessionIds.map(() => '?').join(', ') + const rows = db.prepare(` + SELECT + session_id, + source_key, + source_seq, + created_at, + agent, + model, + kind, + input_tokens, + output_tokens, + cache_read_tokens, + cache_creation_tokens, + last_input_tokens, + last_output_tokens, + last_cache_read_tokens, + last_cache_creation_tokens + FROM usage_events + WHERE session_id IN (${placeholders}) + ORDER BY created_at ASC, source_seq ASC + `).all(...sessionIds) as UsageEventRow[] + + return rows.map(toUsageEvent) +} + +export function getUsageScanStates(db: Database, sessionIds: string[]): Map { + if (sessionIds.length === 0) return new Map() + + const placeholders = sessionIds.map(() => '?').join(', ') + const rows = db.prepare(` + SELECT session_id, message_epoch, last_seq + FROM usage_scan_state + WHERE session_id IN (${placeholders}) + `).all(...sessionIds) as Array<{ session_id: string; message_epoch: number; last_seq: number }> + + return new Map(rows.map((row) => [row.session_id, { + messageEpoch: row.message_epoch, + lastSeq: row.last_seq + }])) +} + +export function transferUsageSession(db: Database, fromSessionId: string, toSessionId: string): void { + if (fromSessionId === toSessionId) return + + db.transaction(() => { + db.prepare(` + INSERT OR IGNORE INTO usage_events ( + session_id, + source_key, + source_seq, + created_at, + agent, + model, + kind, + input_tokens, + output_tokens, + cache_read_tokens, + cache_creation_tokens, + last_input_tokens, + last_output_tokens, + last_cache_read_tokens, + last_cache_creation_tokens + ) + SELECT + ?, + source_key, + source_seq, + created_at, + agent, + model, + kind, + input_tokens, + output_tokens, + cache_read_tokens, + cache_creation_tokens, + last_input_tokens, + last_output_tokens, + last_cache_read_tokens, + last_cache_creation_tokens + FROM usage_events + WHERE session_id = ? + `).run(toSessionId, fromSessionId) + db.prepare('DELETE FROM usage_events WHERE session_id = ?').run(fromSessionId) + db.prepare('DELETE FROM usage_scan_state WHERE session_id IN (?, ?)').run(fromSessionId, toSessionId) + })() +} diff --git a/hub/src/store/usageStore.ts b/hub/src/store/usageStore.ts new file mode 100644 index 00000000..e0831305 --- /dev/null +++ b/hub/src/store/usageStore.ts @@ -0,0 +1,36 @@ +import type { Database } from 'bun:sqlite' + +import { + getUsageEvents, + getUsageScanStates, + recordUsageScan, + transferUsageSession, + type UsageEvent, + type UsageScanState +} from './usage' + +export class UsageStore { + constructor(private readonly db: Database) {} + + recordScan( + sessionId: string, + messageEpoch: number, + lastSeq: number, + events: UsageEvent[], + replaceEvents: boolean + ): void { + recordUsageScan(this.db, sessionId, messageEpoch, lastSeq, events, replaceEvents) + } + + getEvents(sessionIds: string[]): UsageEvent[] { + return getUsageEvents(this.db, sessionIds) + } + + getScanStates(sessionIds: string[]): Map { + return getUsageScanStates(this.db, sessionIds) + } + + transferSession(fromSessionId: string, toSessionId: string): void { + transferUsageSession(this.db, fromSessionId, toSessionId) + } +} diff --git a/hub/src/sync/sessionCache.ts b/hub/src/sync/sessionCache.ts index 23e3c2b5..2ad31a22 100644 --- a/hub/src/sync/sessionCache.ts +++ b/hub/src/sync/sessionCache.ts @@ -947,6 +947,7 @@ export class SessionCache { const movedMessages = this.store.messages.mergeSessionMessages(oldSessionId, newSessionId) if (movedMessages.moved > 0) { + this.store.usage.transferSession(oldSessionId, newSessionId) if (!options.deleteOldSession) { this.publisher.emit({ type: 'messages-invalidated', sessionId: oldSessionId, namespace }) } diff --git a/hub/src/sync/usageService.test.ts b/hub/src/sync/usageService.test.ts new file mode 100644 index 00000000..2b8ee125 --- /dev/null +++ b/hub/src/sync/usageService.test.ts @@ -0,0 +1,440 @@ +import { describe, expect, it } from 'bun:test' +import { Store } from '../store' +import { getUsageSummary } from './usageService' + +function addAgentMessage(store: Store, sessionId: string, content: unknown): void { + store.messages.addMessage(sessionId, { role: 'agent', content }) +} + +describe('usage service', () => { + it('deduplicates Claude stream fragments and normalizes cached input', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'claude-usage-test', + { path: '/tmp', host: 'test', flavor: 'claude' }, + null, + 'default', + 'test-model' + ) + + addAgentMessage(store, session.id, { + type: 'output', + data: { + type: 'assistant', + message: { + id: 'claude-message', + model: 'claude-test', + usage: { input_tokens: 10, output_tokens: 2, cache_read_input_tokens: 80 } + } + } + }) + addAgentMessage(store, session.id, { + type: 'output', + data: { + type: 'assistant', + message: { + id: 'claude-message', + model: 'claude-test', + usage: { input_tokens: 12, output_tokens: 3, cache_read_input_tokens: 90 } + } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(1) + expect(result.totals.inputTokens).toBe(102) + expect(result.totals.outputTokens).toBe(3) + expect(result.totals.cacheReadTokens).toBe(90) + expect(result.totals.totalTokens).toBe(105) + expect(result.totals.uncachedTokens).toBe(15) + expect(result.byModel.find((row) => row.key === 'claude-test')?.totalTokens).toBe(105) + store.close() + }) + + it('uses the latest request as the baseline for a resumed Codex thread', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'codex-usage-test', + { path: '/tmp', host: 'test', flavor: 'codex' }, + null, + 'default', + 'test-model' + ) + + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: 'thread-1', + turn_id: 'turn-1', + scope_role: 'parent', + info: { + total_token_usage: { input_tokens: 1_000, output_tokens: 100, cached_input_tokens: 800 }, + last_token_usage: { input_tokens: 100, output_tokens: 10, cached_input_tokens: 80 } + } + } + }) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: 'thread-2', + turn_id: 'turn-1', + scope_role: 'parent', + info: { + total_token_usage: { input_tokens: 1_000, output_tokens: 100, cached_input_tokens: 800 }, + last_token_usage: { input_tokens: 100, output_tokens: 10, cached_input_tokens: 80 } + } + } + }) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: 'thread-1', + turn_id: 'turn-2', + scope_role: 'parent', + info: { + total_token_usage: { input_tokens: 1_140, output_tokens: 115, cached_input_tokens: 900 }, + last_token_usage: { input_tokens: 140, output_tokens: 15, cached_input_tokens: 100 } + } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(2) + expect(result.totals.inputTokens).toBe(240) + expect(result.totals.outputTokens).toBe(25) + expect(result.totals.cacheReadTokens).toBe(180) + expect(result.totals.totalTokens).toBe(265) + expect(result.totals.uncachedTokens).toBe(85) + store.close() + }) + + it('treats ACP usage totals as per-request deltas', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'kimi-usage-test', + { path: '/tmp', host: 'test', flavor: 'kimi' }, + null, + 'default', + 'kimi-model' + ) + + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + info: { total: { inputTokens: 100, outputTokens: 10, cachedInputTokens: 80 } } + } + }) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + info: { total: { inputTokens: 140, outputTokens: 15, cachedInputTokens: 100 } } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(2) + expect(result.totals.inputTokens).toBe(240) + expect(result.totals.outputTokens).toBe(25) + expect(result.totals.cacheReadTokens).toBe(180) + expect(result.totals.totalTokens).toBe(265) + expect(result.totals.uncachedTokens).toBe(85) + store.close() + }) + + it('accepts Codex events that only contain last_token_usage', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'codex-last-usage-test', + { path: '/tmp', host: 'test', flavor: 'codex' }, + null, + 'default', + 'test-model' + ) + + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: 'thread-1', + turn_id: 'turn-1', + info: { + last_token_usage: { input_tokens: 100, output_tokens: 10, cached_input_tokens: 80 } + } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(1) + expect(result.totals.totalTokens).toBe(110) + expect(result.totals.uncachedTokens).toBe(30) + store.close() + }) + + it('resumes scanning after the last checked message and resets after an epoch change', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'incremental-usage-test', + { path: '/tmp', host: 'test', flavor: 'claude' }, + null, + 'default', + 'test-model' + ) + store.messages.addMessage(session.id, { role: 'user', content: { type: 'text', text: 'hello' } }) + + const afterSeqs: number[] = [] + const getMessagesAfterSeq = store.messages.getMessagesAfterSeq.bind(store.messages) + store.messages.getMessagesAfterSeq = (sessionId, afterSeq) => { + afterSeqs.push(afterSeq) + return getMessagesAfterSeq(sessionId, afterSeq) + } + + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(0) + addAgentMessage(store, session.id, { + type: 'output', + data: { + type: 'assistant', + message: { + id: 'incremental-claude-message', + usage: { input_tokens: 5, output_tokens: 2 } + } + } + }) + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(1) + + store.messages.bumpMessageEpoch(session.id) + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(1) + expect(afterSeqs).toEqual([0, 1, 0]) + store.close() + }) + + it('removes usage from transcript history discarded by a rewind', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'rewound-usage-test', + { path: '/tmp', host: 'test', flavor: 'claude' }, + null, + 'default', + 'test-model' + ) + const claudeUsage = (id: string) => ({ + role: 'agent', + content: { + type: 'output', + data: { + type: 'assistant', + message: { id, usage: { input_tokens: 10, output_tokens: 2 } } + } + } + }) + + store.messages.addMessage(session.id, claudeUsage('kept-message')) + store.messages.addMessage( + session.id, + { role: 'user', content: { type: 'text', text: 'retry this' } }, + 'rewind-point' + ) + store.messages.addMessage(session.id, claudeUsage('discarded-message')) + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(2) + + const result = store.messages.truncateMessagesFromLocalId(session.id, 'rewind-point') + expect(result.deleted).toBe(2) + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(1) + expect(store.usage.getEvents([session.id])).toHaveLength(1) + store.close() + }) + + it('does not deduplicate matching Codex turns from different sessions', () => { + const store = new Store(':memory:') + for (const [sessionId, threadId] of [['codex-session-1', 'thread-1'], ['codex-session-2', 'thread-2']] as const) { + const session = store.sessions.getOrCreateSession( + sessionId, + { path: '/tmp', host: 'test', flavor: 'codex' }, + null, + 'default', + 'test-model' + ) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: threadId, + turn_id: 'matching-turn', + scope_role: 'parent', + info: { + total_token_usage: { input_tokens: 100, output_tokens: 10, cached_input_tokens: 80 }, + last_token_usage: { input_tokens: 100, output_tokens: 10, cached_input_tokens: 80 } + } + } + }) + } + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(2) + expect(result.totals.totalTokens).toBe(220) + store.close() + }) + + it('deduplicates repeated cumulative snapshots without thread metadata', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'imported-codex-usage-test', + { path: '/tmp', host: 'test', flavor: 'codex' }, + null, + 'default', + 'test-model' + ) + const snapshots = [ + { + total: { inputTokens: 100, outputTokens: 10, cachedInputTokens: 80 }, + last: { inputTokens: 100, outputTokens: 10, cachedInputTokens: 80 } + }, + { + total: { inputTokens: 140, outputTokens: 15, cachedInputTokens: 100 }, + last: { inputTokens: 40, outputTokens: 5, cachedInputTokens: 20 } + } + ] + for (const info of [...snapshots, ...snapshots]) { + addAgentMessage(store, session.id, { + type: 'codex', + data: { type: 'token_count', info } + }) + } + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(2) + expect(result.totals.inputTokens).toBe(140) + expect(result.totals.outputTokens).toBe(15) + expect(result.totals.uncachedTokens).toBe(55) + store.close() + }) + + it('excludes pre-HAPI Codex transcript usage from imported sessions', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'resumed-codex-usage-test', + { + path: '/tmp', + host: 'test', + flavor: 'codex', + codexSessionId: 'forked-thread', + codexSourceSessionId: 'imported-thread' + }, + null, + 'default', + 'test-model' + ) + + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + info: { + total_token_usage: { input_tokens: 1_000_000, output_tokens: 10_000, cached_input_tokens: 900_000 }, + last_token_usage: { input_tokens: 100_000, output_tokens: 1_000, cached_input_tokens: 90_000 } + } + } + }) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + thread_id: 'forked-thread', + turn_id: 'managed-turn', + info: { + total_token_usage: { input_tokens: 1_100_000, output_tokens: 11_000, cached_input_tokens: 990_000 }, + last_token_usage: { input_tokens: 100_000, output_tokens: 1_000, cached_input_tokens: 90_000 } + } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(1) + expect(result.totals.totalTokens).toBe(101_000) + expect(result.totals.uncachedTokens).toBe(11_000) + store.close() + }) + + it('excludes transcript replay events explicitly marked by the CLI', () => { + const store = new Store(':memory:') + const session = store.sessions.getOrCreateSession( + 'marked-codex-history-test', + { path: '/tmp', host: 'test', flavor: 'codex' }, + null, + 'default', + 'test-model' + ) + + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + hapiUsageScope: 'imported-history', + info: { + total_token_usage: { input_tokens: 1_000_000, output_tokens: 10_000 }, + last_token_usage: { input_tokens: 100_000, output_tokens: 1_000 } + } + } + }) + addAgentMessage(store, session.id, { + type: 'codex', + data: { + type: 'token_count', + hapiUsageScope: 'managed', + thread_id: 'resumed-thread', + info: { + total_token_usage: { input_tokens: 1_100_000, output_tokens: 11_000 }, + last_token_usage: { input_tokens: 100_000, output_tokens: 1_000 } + } + } + }) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(1) + expect(result.totals.totalTokens).toBe(101_000) + store.close() + }) + + it('keeps usage counted once when session history is merged', () => { + const store = new Store(':memory:') + const source = store.sessions.getOrCreateSession( + 'usage-merge-source', + { path: '/tmp', host: 'test', flavor: 'claude' }, + null, + 'default', + 'test-model' + ) + const target = store.sessions.getOrCreateSession( + 'usage-merge-target', + { path: '/tmp', host: 'test', flavor: 'claude' }, + null, + 'default', + 'test-model' + ) + addAgentMessage(store, source.id, { + type: 'output', + data: { + type: 'assistant', + message: { + id: 'moved-claude-message', + usage: { input_tokens: 10, output_tokens: 2, cache_read_input_tokens: 80 } + } + } + }) + expect(getUsageSummary(store, 'default', 'all').totals.requests).toBe(1) + + const moved = store.messages.mergeSessionMessages(source.id, target.id) + expect(moved.moved).toBe(1) + store.usage.transferSession(source.id, target.id) + + const result = getUsageSummary(store, 'default', 'all') + expect(result.totals.requests).toBe(1) + expect(store.usage.getEvents([source.id])).toEqual([]) + expect(store.usage.getEvents([target.id])).toHaveLength(1) + store.close() + }) +}) diff --git a/hub/src/sync/usageService.ts b/hub/src/sync/usageService.ts new file mode 100644 index 00000000..fcee6398 --- /dev/null +++ b/hub/src/sync/usageService.ts @@ -0,0 +1,319 @@ +import type { UsageSummaryBucket, UsageSummaryResponse } from '@hapi/protocol/apiTypes' +import type { StoredMessage, StoredSession } from '../store' +import type { UsageEvent } from '../store/usage' +import type { Store } from '../store' + +type RecordValue = Record + +function asRecord(value: unknown): RecordValue | null { + return value !== null && typeof value === 'object' && !Array.isArray(value) + ? value as RecordValue + : null +} + +function asCount(value: unknown): number | null { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 + ? Math.floor(value) + : null +} + +function firstCount(record: RecordValue, ...keys: string[]): number { + for (const key of keys) { + const value = asCount(record[key]) + if (value !== null) return value + } + return 0 +} + +function sessionAgent(session: StoredSession): string { + const metadata = asRecord(session.metadata) + const flavor = metadata?.flavor + return typeof flavor === 'string' && flavor.trim() ? flavor.trim() : 'unknown' +} + +function sessionModel(session: StoredSession): string | null { + return typeof session.model === 'string' && session.model.trim() ? session.model.trim() : null +} + +function parseUsageEvent(session: StoredSession, message: StoredMessage): UsageEvent | null { + const envelope = asRecord(message.content) + if (envelope?.role !== 'agent') return null + + const payload = asRecord(envelope.content) + if (!payload) return null + const data = asRecord(payload.data) + if (!data) return null + + // Claude stream-json/SDK messages. A stream emits several updates for one + // assistant message, so the provider's message id is the stable upsert key. + if (payload.type === 'output' && data.type === 'assistant') { + const assistant = asRecord(data.message) + const usage = asRecord(assistant?.usage) + if (!usage) return null + const inputTokens = firstCount(usage, 'input_tokens', 'inputTokens') + const outputTokens = firstCount(usage, 'output_tokens', 'outputTokens') + const cacheReadTokens = firstCount(usage, 'cache_read_input_tokens', 'cacheReadTokens', 'cachedInputTokens') + const cacheCreationTokens = firstCount(usage, 'cache_creation_input_tokens', 'cacheCreationTokens', 'cacheWriteInputTokens') + if (inputTokens + outputTokens + cacheReadTokens + cacheCreationTokens <= 0) return null + const providerId = typeof assistant?.id === 'string' ? assistant.id : message.id + const model = typeof assistant?.model === 'string' && assistant.model.trim() + ? assistant.model.trim() + : sessionModel(session) + return { + sessionId: session.id, + sourceKey: `claude|${providerId}`, + sourceSeq: message.seq, + createdAt: message.createdAt, + agent: 'claude', + model, + kind: 'delta', + inputTokens, + outputTokens, + cacheReadTokens, + cacheCreationTokens, + lastInputTokens: null, + lastOutputTokens: null, + lastCacheReadTokens: null, + lastCacheCreationTokens: null + } + } + + // Codex forwards cumulative thread totals plus the most recent request. + // ACP-compatible backends wrap per-request usage in `total`, so only Codex + // should be diffed as a cumulative stream. + if (data.type === 'token_count' || data.type === 'usage') { + if (data.hapiUsageScope === 'imported-history') return null + const info = asRecord(data.info) ?? data + const agent = sessionAgent(session) + const explicitThreadId = typeof data.threadId === 'string' + ? data.threadId + : typeof data.thread_id === 'string' + ? data.thread_id + : null + const metadata = asRecord(session.metadata) + const hasImportedCodexHistory = typeof metadata?.codexSourceSessionId === 'string' + || metadata?.lifecycleState === 'imported' + if (agent === 'codex' && explicitThreadId === null && hasImportedCodexHistory) { + return null + } + const cumulativeTotal = agent === 'codex' + ? asRecord(info.total) + ?? asRecord(info.total_token_usage) + ?? asRecord(info.totalTokenUsage) + : null + const last = asRecord(info.last) + ?? asRecord(info.last_token_usage) + ?? asRecord(info.lastTokenUsage) + ?? (data.type === 'usage' ? info : null) + const total = cumulativeTotal ?? (agent === 'codex' ? last : asRecord(info.total) ?? info) + if (!total) return null + const inputTokens = firstCount(total, 'inputTokens', 'input_tokens') + const outputTokens = firstCount(total, 'outputTokens', 'output_tokens') + const cacheReadTokens = firstCount(total, 'cachedInputTokens', 'cached_input_tokens', 'cacheReadTokens', 'cache_read_input_tokens') + const cacheCreationTokens = firstCount(total, 'cacheWriteInputTokens', 'cache_write_input_tokens', 'cacheCreationTokens', 'cache_creation_input_tokens') + if (inputTokens + outputTokens + cacheReadTokens + cacheCreationTokens <= 0) return null + const threadId = explicitThreadId ?? session.id + const scope = typeof data.scopeRole === 'string' + ? data.scopeRole + : typeof data.scope_role === 'string' + ? data.scope_role + : 'parent' + const isCumulative = cumulativeTotal !== null + const turnId = typeof data.turnId === 'string' + ? data.turnId + : typeof data.turn_id === 'string' + ? data.turn_id + : '' + return { + sessionId: session.id, + sourceKey: isCumulative + ? [ + 'cumulative', + threadId, + scope, + turnId, + inputTokens, + outputTokens, + cacheReadTokens, + cacheCreationTokens + ].join('|') + : `delta|${message.id}`, + sourceSeq: message.seq, + createdAt: message.createdAt, + agent, + model: sessionModel(session), + kind: isCumulative ? 'cumulative' : 'delta', + inputTokens, + outputTokens, + cacheReadTokens, + cacheCreationTokens, + lastInputTokens: last ? firstCount(last, 'inputTokens', 'input_tokens') : null, + lastOutputTokens: last ? firstCount(last, 'outputTokens', 'output_tokens') : null, + lastCacheReadTokens: last + ? firstCount(last, 'cachedInputTokens', 'cached_input_tokens', 'cacheReadTokens', 'cache_read_input_tokens') + : null, + lastCacheCreationTokens: last + ? firstCount(last, 'cacheWriteInputTokens', 'cache_write_input_tokens', 'cacheCreationTokens', 'cache_creation_input_tokens') + : null + } + } + + return null +} + +function collectUsageEvents(store: Store, sessions: StoredSession[]): void { + const scanStates = store.usage.getScanStates(sessions.map((session) => session.id)) + for (const session of sessions) { + const messageEpoch = store.messages.getMessageEpoch(session.id) + const scanState = scanStates.get(session.id) + const replaceEvents = !scanState || scanState.messageEpoch !== messageEpoch + const afterSeq = replaceEvents ? 0 : scanState.lastSeq + const messages = store.messages.getMessagesAfterSeq(session.id, afterSeq) + const events = new Map() + for (const message of messages) { + const event = parseUsageEvent(session, message) + if (!event) continue + if (event.kind === 'delta' || !events.has(event.sourceKey)) { + events.set(event.sourceKey, event) + } + } + const lastSeq = messages.at(-1)?.seq ?? afterSeq + if (messages.length > 0 || replaceEvents) { + store.usage.recordScan( + session.id, + messageEpoch, + lastSeq, + Array.from(events.values()), + replaceEvents + ) + } + } +} + +type Totals = Omit + +function emptyTotals(): Totals { + return { + inputTokens: 0, + outputTokens: 0, + cacheReadTokens: 0, + cacheCreationTokens: 0, + totalTokens: 0, + uncachedTokens: 0, + requests: 0 + } +} + +function addTotals(target: Totals, inputTokens: number, outputTokens: number, cacheReadTokens: number, cacheCreationTokens: number): void { + target.inputTokens += inputTokens + target.outputTokens += outputTokens + target.cacheReadTokens += cacheReadTokens + target.cacheCreationTokens += cacheCreationTokens + // Codex/Kimi inputTokens already includes cached input. Claude's raw + // input_tokens excludes cache fields and is normalized before this call. + target.totalTokens += inputTokens + outputTokens + target.uncachedTokens += Math.max(0, inputTokens - cacheReadTokens) + outputTokens + target.requests += 1 +} + +function cumulativeDelta(current: number, previous: number | null, last: number | null): number { + if (previous === null) return last ?? current + return current >= previous ? current - previous : last ?? current +} + +function toBucket(key: string, totals: Totals): UsageSummaryBucket { + return { key, ...totals } +} + +function dayKey(timestamp: number): string { + return new Date(timestamp).toISOString().slice(0, 10) +} + +export function getUsageSummary(store: Store, namespace: string, range: string | undefined): UsageSummaryResponse { + const sessions = store.sessions.getSessionsByNamespace(namespace) + // This is intentionally lazy. Existing HAPI databases have no usage table; + // the first dashboard request backfills history, while later requests only + // update the idempotent event rows. + collectUsageEvents(store, sessions) + + const now = Date.now() + const days = range === '30d' ? 30 : range === 'all' ? null : 7 + const from = days === null ? null : now - days * 24 * 60 * 60 * 1000 + const sessionIds = new Set(sessions.map((session) => session.id)) + const events = store.usage.getEvents(Array.from(sessionIds)) + const isInRange = (event: UsageEvent) => (from === null || event.createdAt >= from) && event.createdAt <= now + + const totals = emptyTotals() + const daily = new Map() + const byAgent = new Map() + const byModel = new Map() + const sessionsWithUsage = new Set() + const cumulativePrevious = new Map() + const cumulativeFingerprints = new Set() + + for (const event of events) { + let inputTokens = event.inputTokens + let outputTokens = event.outputTokens + let cacheReadTokens = event.cacheReadTokens + let cacheCreationTokens = event.cacheCreationTokens + let duplicateCumulativeEvent = false + if (event.kind === 'cumulative') { + const sourceParts = event.sourceKey.split('|') + const streamKey = sourceParts.slice(0, 3).join('|') + const previous = cumulativePrevious.get(streamKey) + inputTokens = cumulativeDelta(inputTokens, previous?.[0] ?? null, event.lastInputTokens) + outputTokens = cumulativeDelta(outputTokens, previous?.[1] ?? null, event.lastOutputTokens) + cacheReadTokens = cumulativeDelta(cacheReadTokens, previous?.[2] ?? null, event.lastCacheReadTokens) + cacheCreationTokens = cumulativeDelta(cacheCreationTokens, previous?.[3] ?? null, event.lastCacheCreationTokens) + cumulativePrevious.set(streamKey, [event.inputTokens, event.outputTokens, event.cacheReadTokens, event.cacheCreationTokens]) + const turnId = sourceParts[3] + if (turnId) { + const fingerprint = [ + event.sessionId, + turnId, + event.inputTokens, + event.outputTokens, + event.cacheReadTokens, + event.cacheCreationTokens, + event.lastInputTokens, + event.lastOutputTokens, + event.lastCacheReadTokens, + event.lastCacheCreationTokens + ].join('|') + duplicateCumulativeEvent = cumulativeFingerprints.has(fingerprint) + cumulativeFingerprints.add(fingerprint) + } + } + if (duplicateCumulativeEvent || !isInRange(event) || inputTokens + outputTokens + cacheReadTokens + cacheCreationTokens <= 0) continue + const normalizedInputTokens = event.agent === 'claude' + ? inputTokens + cacheReadTokens + cacheCreationTokens + : inputTokens + addTotals(totals, normalizedInputTokens, outputTokens, cacheReadTokens, cacheCreationTokens) + const dailyTotals = daily.get(dayKey(event.createdAt)) ?? emptyTotals() + addTotals(dailyTotals, normalizedInputTokens, outputTokens, cacheReadTokens, cacheCreationTokens) + daily.set(dayKey(event.createdAt), dailyTotals) + const agentTotals = byAgent.get(event.agent) ?? emptyTotals() + addTotals(agentTotals, normalizedInputTokens, outputTokens, cacheReadTokens, cacheCreationTokens) + byAgent.set(event.agent, agentTotals) + const modelKey = event.model ?? 'unknown' + const modelTotals = byModel.get(modelKey) ?? emptyTotals() + addTotals(modelTotals, normalizedInputTokens, outputTokens, cacheReadTokens, cacheCreationTokens) + byModel.set(modelKey, modelTotals) + sessionsWithUsage.add(event.sessionId) + } + + const sortBuckets = (values: Map): UsageSummaryBucket[] => Array.from(values.entries()) + .map(([key, value]) => toBucket(key, value)) + .sort((a, b) => b.totalTokens - a.totalTokens) + + return { + range: { from, to: now }, + totals: { ...totals, sessions: sessionsWithUsage.size }, + daily: Array.from(daily.entries()) + .map(([key, value]) => toBucket(key, value)) + .sort((a, b) => a.key.localeCompare(b.key)), + byAgent: sortBuckets(byAgent), + byModel: sortBuckets(byModel), + updatedAt: now + } +} diff --git a/hub/src/web/routes/usage.test.ts b/hub/src/web/routes/usage.test.ts new file mode 100644 index 00000000..6a386dee --- /dev/null +++ b/hub/src/web/routes/usage.test.ts @@ -0,0 +1,45 @@ +import { describe, expect, it } from 'bun:test' +import { Hono } from 'hono' +import { Store } from '../../store' +import type { WebAppEnv } from '../middleware/auth' +import { createUsageRoutes } from './usage' + +function createApp(store: Store, namespace: string): Hono { + const app = new Hono() + app.use('*', async (c, next) => { + c.set('namespace', namespace) + await next() + }) + app.route('/api', createUsageRoutes(store)) + return app +} + +describe('GET /api/usage/summary', () => { + it('returns a no-store summary for the hub owner', async () => { + const store = new Store(':memory:') + try { + const response = await createApp(store, 'default').request('/api/usage/summary?range=30d') + + expect(response.status).toBe(200) + expect(response.headers.get('cache-control')).toBe('no-store') + const body = await response.json() as { range: { from: number | null; to: number | null }; totals: { requests: number } } + expect(typeof body.range.from).toBe('number') + expect(typeof body.range.to).toBe('number') + expect(body.totals.requests).toBe(0) + } finally { + store.close() + } + }) + + it('rejects non-default namespaces', async () => { + const store = new Store(':memory:') + try { + const response = await createApp(store, 'tenant').request('/api/usage/summary') + + expect(response.status).toBe(403) + expect(await response.json()).toEqual({ error: 'Usage summary is only available to the hub owner' }) + } finally { + store.close() + } + }) +}) diff --git a/hub/src/web/routes/usage.ts b/hub/src/web/routes/usage.ts new file mode 100644 index 00000000..73c6fc17 --- /dev/null +++ b/hub/src/web/routes/usage.ts @@ -0,0 +1,21 @@ +import { Hono } from 'hono' +import type { UsageSummaryResponse } from '@hapi/protocol/apiTypes' +import type { WebAppEnv } from '../middleware/auth' +import type { Store } from '../../store' +import { getUsageSummary } from '../../sync/usageService' + +export function createUsageRoutes(store: Store): Hono { + const app = new Hono() + + app.get('/usage/summary', (c) => { + if (c.get('namespace') !== 'default') { + return c.json({ error: 'Usage summary is only available to the hub owner' }, 403) + } + const range = c.req.query('range') + const response: UsageSummaryResponse = getUsageSummary(store, c.get('namespace'), range) + c.header('Cache-Control', 'no-store') + return c.json(response) + }) + + return app +} diff --git a/hub/src/web/server.ts b/hub/src/web/server.ts index b9668b27..6ae596b5 100644 --- a/hub/src/web/server.ts +++ b/hub/src/web/server.ts @@ -19,6 +19,7 @@ import { createMessagesRoutes } from './routes/messages' import { createPermissionsRoutes } from './routes/permissions' import { createMachinesRoutes } from './routes/machines' import { createStorageRoutes } from './routes/storage' +import { createUsageRoutes } from './routes/usage' import { createGitRoutes } from './routes/git' import { createCliRoutes } from './routes/cli' import { createCodexDesktopRoutes } from './routes/codexDesktop' @@ -250,6 +251,7 @@ function createWebApp(options: { app.route('/api', createPermissionsRoutes(options.getSyncEngine)) app.route('/api', createMachinesRoutes(options.getSyncEngine)) app.route('/api', createStorageRoutes(configuration.dbPath)) + app.route('/api', createUsageRoutes(options.store)) app.route('/api', createGitRoutes(options.getSyncEngine)) // 中文注释:这里提供两类 Codex 辅助能力:扫描本地 transcript 以导入到 Hapi,以及按需重启 Codex Desktop 客户端。 app.route('/api', createCodexDesktopRoutes({ diff --git a/shared/src/apiTypes.ts b/shared/src/apiTypes.ts index ba51c68d..a80d863f 100644 --- a/shared/src/apiTypes.ts +++ b/shared/src/apiTypes.ts @@ -731,3 +731,35 @@ export type SqliteStorageUsageResponse = { shmBytes: number totalBytes: number } + +export type UsageSummaryBucket = { + key: string + inputTokens: number + outputTokens: number + cacheReadTokens: number + cacheCreationTokens: number + totalTokens: number + uncachedTokens: number + requests: number +} + +export type UsageSummaryResponse = { + range: { + from: number | null + to: number | null + } + totals: { + inputTokens: number + outputTokens: number + cacheReadTokens: number + cacheCreationTokens: number + totalTokens: number + uncachedTokens: number + requests: number + sessions: number + } + daily: Array + byAgent: UsageSummaryBucket[] + byModel: UsageSummaryBucket[] + updatedAt: number +} diff --git a/web/README.md b/web/README.md index 096c02d3..0f706c52 100644 --- a/web/README.md +++ b/web/README.md @@ -37,6 +37,7 @@ See `src/router.tsx` for route definitions. - `/settings/voice` - Everyday voice assistant preferences. - `/settings/voice/voices` - Full-page voice picker. - `/settings/voice/advanced` - Voice persona, tuning, and diagnostics. +- `/settings/usage` - Cache-aware token usage dashboard for the hub owner. - `/settings/about` - Application links and version information. ## Features diff --git a/web/src/api/client.ts b/web/src/api/client.ts index e945c748..d02cbbe1 100644 --- a/web/src/api/client.ts +++ b/web/src/api/client.ts @@ -43,6 +43,7 @@ import type { QueuedStateResponse, ReopenSessionResponse, SqliteStorageUsageResponse, + UsageSummaryResponse, UploadFileResponse } from '@hapi/protocol/apiTypes' import type { AgentFlavor } from '@hapi/protocol' @@ -647,6 +648,10 @@ export class ApiClient { return await this.request('/api/storage/sqlite') } + async getUsageSummary(range: '7d' | '30d' | 'all' = '7d'): Promise { + return await this.request(`/api/usage/summary?range=${encodeURIComponent(range)}`) + } + async listMachineDirectory( machineId: string, path: string, diff --git a/web/src/components/settings/SettingsNav.tsx b/web/src/components/settings/SettingsNav.tsx index 7fb3bf07..2763dd02 100644 --- a/web/src/components/settings/SettingsNav.tsx +++ b/web/src/components/settings/SettingsNav.tsx @@ -7,6 +7,8 @@ import { useAppContext } from '@/lib/app-context' import { settingsCategories } from '@/routes/settings/categories' import { ChevronRightIcon } from './SettingsPrimitives' +const OWNER_ONLY_CATEGORIES = new Set(['storage', 'usage']) + function getNamespace(token: string): string | null { try { const payload = token.split('.')[1] @@ -34,9 +36,10 @@ export function SettingsNav(props: { activeId?: string; mobile?: boolean }) { voice: t('settings.hub.voice.summary'), machines: t('settings.hub.machines.summary'), storage: t('settings.storage.summary'), + usage: t('settings.usage.summary'), about: `v${__APP_VERSION__}`, } - const visibleCategories = settingsCategories.filter((category) => category.id !== 'storage' || getNamespace(token) === 'default') + const visibleCategories = settingsCategories.filter((category) => !OWNER_ONLY_CATEGORIES.has(category.id) || getNamespace(token) === 'default') return (