From e8960e1d59ae2ac223111ef0bab3f87a7559ab98 Mon Sep 17 00:00:00 2001 From: liuxiaocs7 Date: Sun, 20 Sep 2026 05:15:39 +0800 Subject: [PATCH] perf(storage): reuse Usage range statistics across filter changes Fixes #5516 Refs #5038 Generated-by: Codex --- .../src/__tests__/usage-screen.test.ts | 16 + .../src/__tests__/usage-stores.test.ts | 310 +++++++++++++++++- packages/storage/src/usage-screen.ts | 276 +++++++++------- scripts/perf/usage-statistics.mjs | 258 +++++++++++++++ 4 files changed, 742 insertions(+), 118 deletions(-) create mode 100644 scripts/perf/usage-statistics.mjs diff --git a/packages/runtime-host/src/__tests__/usage-screen.test.ts b/packages/runtime-host/src/__tests__/usage-screen.test.ts index 7adaede9e9..f239bc2a16 100644 --- a/packages/runtime-host/src/__tests__/usage-screen.test.ts +++ b/packages/runtime-host/src/__tests__/usage-screen.test.ts @@ -296,6 +296,15 @@ test('real Host returns bounded failures, stays usable, and fences a replacement assert.ok(loaded.ok && loaded.result.kind === 'screen'); assert.doesNotThrow(() => decodeUsageScreenResult(loaded.result)); const value = loaded.result.screen; + assert.ok(Object.isFrozen(value.byTool[0])); + const filtered = await host.handlers['usage.query']( + { kind: 'screen', query: { ...query, search: 'tool' } }, + context, + ); + assert.ok(filtered.ok && filtered.result.kind === 'screen'); + assert.deepEqual(filtered.result.screen.byTool, value.byTool); + assert.notEqual(filtered.result.screen.byTool[0], value.byTool[0]); + assert.equal(filtered.result.screen.revision, value.revision); assert.ok(value.nextCursor); const input = { kind: 'activity' as const, @@ -310,6 +319,13 @@ test('real Host returns bounded failures, stays usable, and fences a replacement ok: true, result: { kind: 'revision_changed' }, }); + const replacement = await coordinator().handlers['usage.query']( + { kind: 'screen', query }, + context, + ); + assert.ok(replacement.ok && replacement.result.kind === 'screen'); + assert.notEqual(replacement.result.screen.revision, value.revision); + assert.deepEqual(replacement.result.screen.byTool, value.byTool); } finally { lease.close(); await stores.close(); diff --git a/packages/storage/src/__tests__/usage-stores.test.ts b/packages/storage/src/__tests__/usage-stores.test.ts index 3df3181327..546f03db89 100644 --- a/packages/storage/src/__tests__/usage-stores.test.ts +++ b/packages/storage/src/__tests__/usage-stores.test.ts @@ -19,7 +19,7 @@ import assert from 'node:assert/strict'; import { DatabaseSync } from 'node:sqlite'; -import type { UsageScreen, UsageScreenRequest } from '@maka/core/settings'; +import type { UsageScreen, UsageScreenQuery, UsageScreenRequest } from '@maka/core/settings'; import { copyFile, mkdir, @@ -816,7 +816,7 @@ async function withScreenStores( } async function initialScreen( stores: Awaited>, - query = screenQuery, + query: UsageScreenQuery = screenQuery, ) { const result = await stores.readUsageScreen({ kind: 'screen', query }); assert.equal(result.kind, 'screen'); @@ -1101,3 +1101,309 @@ describe('revision-consistent Usage screen', () => { }); }); }); + +describe('Usage range statistics reuse', () => { + test('search and status changes reuse complete statistics for the same revision and range', async () => { + await withScreenStores(async (stores, root) => { + await seedScreen(stores); + await stores.telemetry.recordLlmCall( + llmRecord({ id: 'old-error', ts: 1, modelId: 'ÄModel', status: 'error' }), + ); + await stores.telemetry.recordToolInvocation(toolRecord({ ts: 2 })); + const lease = acquireOperationalStateDatabase(await realpath(root)); + const prepare = lease.database.prepare; + let aggregateQueries = 0; + // Observe database work without replacing its results. Value assertions + // below independently protect the range/filter contract. + lease.database.prepare = function (sql) { + if (/\bGROUP\s+BY\b|\bSUM\s*\(/i.test(sql)) aggregateQueries++; + return prepare.call(this, sql); + }; + try { + const first = await initialScreen(stores); + assert.equal(first.summary.totalRequests, 62); + assert.equal(first.byModel.length, 2); + assert.equal(first.byTool[0]?.calls, 1); + assert.ok(aggregateQueries >= 4, 'the cold screen computes range statistics'); + for (const [search, status, total] of [ + ['gpt', 'success', 61], + ['ÄMODEL', 'error', 1], + ['absentword', 'all', 0], + ['', 'all', 63], + ] as const) { + aggregateQueries = 0; + const filtered = await initialScreen(stores, { ...screenQuery, search, status }); + assert.equal(filtered.revision, first.revision); + assert.deepEqual(rangeStatistics(filtered), rangeStatistics(first)); + assert.equal(filtered.activityTotal, total); + assert.equal(filtered.logs.length, Math.min(50, total)); + assert.equal(aggregateQueries, 0, 'filter edits do not reaggregate unchanged history'); + } + } finally { + lease.database.prepare = prepare; + lease.close(); + } + }); + }); + + test('mutating or freezing a response cannot affect a later screen', async () => { + await withScreenStores(async (stores) => { + await seedScreen(stores); + await stores.telemetry.recordToolInvocation(toolRecord()); + await stores.pricing.upsert(0, { + modelKey: 'openai:gpt-5', + inputUsdPer1M: 1, + outputUsdPer1M: 2, + }); + const first = await initialScreen(stores, structuredClone(screenQuery)); + const expected = structuredClone(rangeStatistics(first)); + first.summary.totalRequests = 999; + first.byProvider[0]!.provider = 'projected provider'; + first.byModel[0]!.model = 'projected model'; + first.byTool[0]!.tool = 'projected tool'; + first.pricing[0]!.inputPerMTokUsd = 999; + first.provenance.coverage.attempts = 999; + first.provenance.pendingRepairs = 999; + first.query.range.to = 0; + const freeze = (value: unknown): void => { + if (!value || typeof value !== 'object') return; + for (const child of Object.values(value)) freeze(child); + Object.freeze(value); + }; + freeze(first); + const second = await initialScreen(stores, { ...screenQuery, search: 'gpt' }); + assert.deepEqual(rangeStatistics(second), expected); + assert.equal(second.activityTotal, 61); + assert.equal(Object.isFrozen(second.byProvider[0]), false); + second.byProvider[0]!.provider = 'second projection'; + freeze(second); + assert.deepEqual(rangeStatistics(await initialScreen(stores)), expected); + }); + }); + + test('failed transaction commits cannot publish statistics under a reusable revision', async () => { + await withScreenStores(async (stores, root) => { + await stores.telemetry.recordToolInvocation(toolRecord({ toolName: 'Original' })); + const lease = acquireOperationalStateDatabase(await realpath(root)); + const exec = lease.database.exec; + const prepare = lease.database.prepare; + let rolledBackRevision: string | undefined; + let aggregateQueries = 0; + let commitFailed = false; + // Force the reader to join a write transaction, then fail its outer + // commit. Both the data and trigger-owned revision roll back together. + lease.database.exec = function (sql) { + if (sql === 'BEGIN') { + exec.call(this, 'BEGIN IMMEDIATE'); + exec.call( + this, + "UPDATE usage_tool_invocations SET record_json = json_set(record_json, '$.toolName', 'Rolled back')", + ); + rolledBackRevision = String( + prepare.call(this, 'SELECT revision FROM usage_screen_revision').get()?.revision, + ); + return; + } + if (sql === 'COMMIT') { + commitFailed = true; + throw new Error('simulated commit failure'); + } + return exec.call(this, sql); + }; + try { + await assert.rejects(initialScreen(stores), /simulated commit failure/); + assert.ok(commitFailed); + lease.database.exec = exec; + lease.transaction('write', () => { + lease.database.exec( + "UPDATE usage_tool_invocations SET record_json = json_set(record_json, '$.toolName', 'Committed')", + ); + }); + assert.equal( + String( + lease.database.prepare('SELECT revision FROM usage_screen_revision').get()?.revision, + ), + rolledBackRevision, + 'the committed write can reuse the counter observed before rollback', + ); + lease.database.prepare = function (sql) { + if (/\bGROUP\s+BY\b|\bSUM\s*\(/i.test(sql)) aggregateQueries++; + return prepare.call(this, sql); + }; + const committed = await initialScreen(stores); + assert.equal(committed.byTool[0]?.tool, 'Committed'); + assert.equal(committed.logs[0]?.toolName, 'Committed'); + assert.ok(aggregateQueries >= 4, 'rolled-back statistics are recomputed'); + aggregateQueries = 0; + assert.deepEqual( + rangeStatistics(await initialScreen(stores, { ...screenQuery, search: 'committed' })), + rangeStatistics(committed), + ); + assert.equal(aggregateQueries, 0, 'a successfully committed read can be reused'); + } finally { + lease.database.exec = exec; + lease.database.prepare = prepare; + lease.close(); + } + }); + }); + + test('range boundaries, accounting writes, pricing and repair changes refresh statistics', async () => { + await withScreenStores(async (stores, root) => { + await stores.telemetry.recordLlmCall(llmRecord({ id: 'early', ts: 10, costUsd: 1 })); + await stores.telemetry.recordLlmCall(llmRecord({ id: 'late', ts: 20, costUsd: 2 })); + await stores.telemetry.recordToolInvocation(toolRecord({ ts: 15 })); + assert.equal((await initialScreen(stores)).summary.totalCostUsd, 3); + const narrow = { ...screenQuery, range: { from: 10, to: 19 } }; + assert.equal((await initialScreen(stores, narrow)).summary.totalCostUsd, 1); + const later = { ...screenQuery, range: { from: 11, to: 20 } }; + assert.equal((await initialScreen(stores, later)).summary.totalCostUsd, 2); + const empty = await initialScreen(stores, { ...screenQuery, range: { from: 11, to: 14 } }); + assert.equal(empty.summary.totalRequests, 0); + assert.equal(empty.byTool.length, 0); + assert.equal(empty.activityTotal, 0); + let previous = await initialScreen(stores); + const refreshed = async () => { + const screen = await initialScreen(stores); + assert.notEqual(screen.revision, previous.revision); + previous = screen; + return screen; + }; + await stores.telemetry.recordLlmCall(llmRecord({ id: 'early', ts: 10, costUsd: 4 })); + assert.equal((await refreshed()).summary.totalCostUsd, 6); + await stores.telemetry.recordToolInvocation(toolRecord({ id: 'second', status: 'error' })); + const tools = await refreshed(); + assert.equal(tools.byTool[0]?.calls, 2); + assert.equal(tools.byTool[0]?.errors, 1); + await stores.pricing.upsert(0, { + modelKey: 'openai:gpt-5', + inputUsdPer1M: 3, + outputUsdPer1M: 4, + }); + assert.deepEqual((await refreshed()).pricing, [ + { provider: 'openai', model: 'gpt-5', inputPerMTokUsd: 3, outputPerMTokUsd: 4 }, + ]); + appendModelCallAuthorityEvent(root, modelCallAttempt('cache-source')); + assert.equal((await refreshed()).provenance.pendingRepairs, 1); + await stores.modelCalls.catchUpModelCallProjection(); + const repaired = await refreshed(); + assert.equal(repaired.provenance.pendingRepairs, 0); + assert.equal(repaired.provenance.coverage.attempts, 1); + assert.equal(repaired.summary.totalRequests, 3); + const lease = acquireOperationalStateDatabase(await realpath(root)); + try { + lease.transaction('write', () => { + lease.database.exec( + 'UPDATE usage_model_call_projection_checkpoints SET unreadable_events = 2', + ); + }); + assert.equal((await refreshed()).provenance.unreadableRecords, 2); + lease.transaction('write', () => lease.database.exec('DELETE FROM usage_llm_calls')); + assert.equal((await refreshed()).summary.totalRequests, 1); + } finally { + lease.close(); + } + }); + }); + + test('a warm cache cannot bypass capacity failures or keep a failed range result', async () => { + await withScreenStores(async (stores) => { + await stores.telemetry.recordToolInvocation(toolRecord({ id: 'first', ts: 1 })); + const first = await initialScreen(stores); + assert.equal(first.byTool.length, 1); + for (let i = 0; i < 100; i++) { + await stores.telemetry.recordToolInvocation( + toolRecord({ id: `tool-${i}`, toolName: `tool-${i}`, ts: 2 }), + ); + } + for (const search of ['', 'absentword']) { + assert.deepEqual( + await stores.readUsageScreen({ kind: 'screen', query: { ...screenQuery, search } }), + { kind: 'screen_response_too_large', section: 'tool_breakdown' }, + ); + } + const small = await initialScreen(stores, { ...screenQuery, range: { from: 1, to: 1 } }); + assert.deepEqual(small.byTool, first.byTool); + assert.deepEqual(await stores.readUsageScreen({ kind: 'screen', query: screenQuery }), { + kind: 'screen_response_too_large', + section: 'tool_breakdown', + }); + await stores.telemetry.recordToolInvocation( + toolRecord({ id: 'tool-0', toolName: 'Bash', ts: 2 }), + ); + const recovered = await initialScreen(stores); + assert.equal(recovered.byTool.length, 100); + assert.equal(recovered.activityTotal, 101); + const filtered = await initialScreen(stores, { ...screenQuery, search: 'absentword' }); + assert.deepEqual(filtered.byTool, recovered.byTool); + assert.equal(filtered.activityTotal, 0); + }); + }); + + test('a cache hit and its activity stay in one snapshot during an external WAL commit', async () => { + await withScreenStores(async (stores, root) => { + await seedScreen(stores); + const first = await initialScreen(stores); + const lease = acquireOperationalStateDatabase(await realpath(root)); + const external = new DatabaseSync(lease.databasePath); + let committed = false; + lease.database.function('usage_screen_lower', { deterministic: true }, (value) => { + if (!committed) { + committed = true; + external.exec( + "UPDATE usage_llm_calls SET record_json = json_set(record_json, '$.costUsd', 2)", + ); + } + return String(value ?? '').toLowerCase(); + }); + try { + const filtered = await initialScreen(stores, { ...screenQuery, search: 'gpt' }); + assert.ok(committed); + assert.equal(filtered.revision, first.revision); + assert.deepEqual(rangeStatistics(filtered), rangeStatistics(first)); + assert.ok(filtered.logs.every((row) => row.costUsd === 0.001)); + assert.deepEqual(await stores.readUsageScreen(continuation(filtered)), { + kind: 'revision_changed', + }); + const refreshed = await initialScreen(stores); + assert.notEqual(refreshed.revision, first.revision); + assert.equal(refreshed.summary.totalCostUsd, 122); + assert.ok(refreshed.logs.every((row) => row.costUsd === 2)); + } finally { + lease.database.function('usage_screen_lower', { deterministic: true }, (value) => + String(value ?? '').toLowerCase(), + ); + external.close(); + lease.close(); + } + }); + }); + + test('different roots with identical ranges retain independent statistics', async () => { + await withScreenStores(async (first) => { + await withScreenStores(async (second) => { + await first.telemetry.recordLlmCall(llmRecord({ costUsd: 1, modelId: 'first' })); + await second.telemetry.recordLlmCall(llmRecord({ costUsd: 9, modelId: 'second' })); + const firstScreen = await initialScreen(first); + const secondScreen = await initialScreen(second); + assert.notEqual(firstScreen.revision, secondScreen.revision); + assert.equal(firstScreen.summary.totalCostUsd, 1); + assert.equal(secondScreen.summary.totalCostUsd, 9); + assert.equal( + (await initialScreen(first, { ...screenQuery, search: 'second' })).activityTotal, + 0, + ); + assert.deepEqual(rangeStatistics(await initialScreen(first)), rangeStatistics(firstScreen)); + assert.deepEqual( + rangeStatistics(await initialScreen(second)), + rangeStatistics(secondScreen), + ); + }); + }); + }); +}); + +function rangeStatistics(screen: UsageScreen) { + const { summary, byProvider, byModel, byTool, pricing, provenance } = screen; + return { summary, byProvider, byModel, byTool, pricing, provenance }; +} diff --git a/packages/storage/src/usage-screen.ts b/packages/storage/src/usage-screen.ts index d66a752404..c313b9d584 100644 --- a/packages/storage/src/usage-screen.ts +++ b/packages/storage/src/usage-screen.ts @@ -70,14 +70,33 @@ const n = (value: unknown): number => Number(value ?? 0); const hash = (value: unknown): string => createHash('sha256').update(JSON.stringify(value)).digest('hex'); +type RangeStatistics = Pick< + UsageScreen, + 'summary' | 'byProvider' | 'byModel' | 'byTool' | 'pricing' | 'provenance' +>; +interface CachedRangeStatistics { + revision: string; + from: number; + to: number; + statistics: RangeStatistics; +} + export function createUsageScreenReader(root: string) { const lease = acquireOperationalStateDatabase(root); // A restore/reopen can repeat both the durable counter and incarnation from a // backup. This lifecycle fence prevents tokens from surviving that reopen. const generation = randomUUID(); const db = lease.database; + // One complete range per reader; changing filters does not retain history or + // grow the cache. Revision includes the database incarnation and lifecycle. + let cached: CachedRangeStatistics | undefined; + let closed = false; return { - close: () => lease.close(), + close: () => { + closed = true; + cached = undefined; + lease.close(); + }, read: (input: UsageScreenRequest): UsageScreenResult => lease.transaction('read', () => { validateQuery(input.query); @@ -99,133 +118,158 @@ export function createUsageScreenReader(root: string) { page: { revision, queryIdentity, ...activity(db, input.query, input.cursor) }, }; } - const rangeArgs = [input.query.range.from, input.query.range.to]; - const modelArgs = [...rangeArgs, ...rangeArgs]; - const aggregate = db - .prepare(`WITH rows AS (${MODEL_ROWS}) SELECT - COUNT(*) AS totalRequests, COALESCE(SUM(costUsd), 0) AS totalCostUsd, - ${['totalTokens', 'inputTokens', 'outputTokens', 'cacheMiss', 'cacheRead', 'cacheCreation', 'reasoning'].map((k) => `COALESCE(SUM(${k}), 0) AS ${k}`).join(', ')}, - SUM(source = 'legacy') AS legacyRecords, - SUM(source = 'canonical') AS attempts, SUM(costBasis = 'priced') AS pricedAttempts, - SUM(costBasis = 'unpriced') AS unpricedAttempts, SUM(usageBasis = 'reported') AS usageReportedAttempts, - SUM(usageBasis = 'partial') AS usagePartialAttempts, SUM(usageBasis = 'missing') AS usageMissingAttempts - FROM rows`) - .get(...modelArgs) as Row; - const breakdown = (column: 'connection' | 'model') => - db - .prepare(`WITH rows AS (${MODEL_ROWS}) - SELECT ${column} AS name, COUNT(*) AS requests, SUM(inputTokens + outputTokens) AS tokens, - COALESCE(SUM(costUsd), 0) AS costUsd FROM rows GROUP BY ${column} ORDER BY requests DESC, name LIMIT 101`) - .all(...modelArgs) as Row[]; - const tools = db - .prepare(`WITH rows AS (${TOOL_ROWS}) SELECT toolName AS tool, - COUNT(*) AS calls, SUM(status = 'success') AS success, SUM(status = 'error') AS errors, - ROUND(AVG(latencyMs)) AS avgDurationMs FROM rows GROUP BY toolName ORDER BY calls DESC, toolName LIMIT 101`) - .all(...rangeArgs) as Row[]; - const providers = breakdown('connection'); - const models = breakdown('model'); - // Read one sentinel past each wire count. It can only become a whole-screen - // failure, never a successful truncated collection. No history array is built. - for (const [rows, section] of [ - [providers, 'provider_breakdown'], - [models, 'model_breakdown'], - [tools, 'tool_breakdown'], - ] as const) { - if (rows.length > 100) return { kind: 'screen_response_too_large', section }; - } - const pricing = db - .prepare('SELECT record_json FROM usage_pricing_overrides ORDER BY model_key LIMIT 129') - .all() - .map((row) => { - const value = JSON.parse(String(row.record_json)); - const separator = value.modelKey.indexOf(':'); - return { - provider: separator < 0 ? '' : value.modelKey.slice(0, separator), - model: separator < 0 ? value.modelKey : value.modelKey.slice(separator + 1), - inputPerMTokUsd: value.inputUsdPer1M, - outputPerMTokUsd: value.outputUsdPer1M, - }; - }); - if (pricing.length > 128) return { kind: 'screen_response_too_large', section: 'pricing' }; - const unreadable = - n( - db - .prepare(`SELECT COUNT(*) AS count FROM usage_model_call_attempts - WHERE completed_at >= ? AND completed_at <= ? AND cost_basis IS NULL`) - .get(...rangeArgs)?.count, - ) + - n( - db - .prepare( - 'SELECT SUM(unreadable_events) AS count FROM usage_model_call_projection_checkpoints', - ) - .get()?.count, - ); - const pending = n( - db - .prepare(`SELECT COUNT(*) AS count FROM core_agent_runs AS source - LEFT JOIN usage_model_call_projection_checkpoints AS checkpoint - ON checkpoint.session_id = source.session_id AND checkpoint.run_id = source.run_id - WHERE source.latest_model_call_sequence > COALESCE(checkpoint.applied_through_sequence, -1)`) - .get()?.count, - ); + const { from, to } = input.query.range; + const hit = + cached?.revision === revision && cached.from === from && cached.to === to + ? cached + : undefined; + const statistics = hit?.statistics ?? readRangeStatistics(db, input.query.range); + if ('kind' in statistics) return statistics; const screen: UsageScreen = { revision, queryIdentity, query: input.query, activityTotal: activityCount(db, input.query), - summary: { - totalRequests: n(aggregate.totalRequests), - totalCostUsd: n(aggregate.totalCostUsd), - totalTokens: n(aggregate.totalTokens), - inputTokens: n(aggregate.inputTokens), - outputTokens: n(aggregate.outputTokens), - cacheTokens: n(aggregate.cacheRead) + n(aggregate.cacheCreation), - cacheMiss: n(aggregate.cacheMiss), - cacheRead: n(aggregate.cacheRead), - cacheCreation: n(aggregate.cacheCreation), - reasoning: n(aggregate.reasoning), - }, - byProvider: providers.map((row) => ({ - provider: String(row.name), - requests: n(row.requests), - tokens: n(row.tokens), - costUsd: n(row.costUsd), - })), - byModel: models.map((row) => ({ - model: String(row.name), - requests: n(row.requests), - tokens: n(row.tokens), - costUsd: n(row.costUsd), - })), - byTool: tools.map((row) => ({ - tool: String(row.tool), - calls: n(row.calls), - success: n(row.success), - errors: n(row.errors), - avgDurationMs: n(row.avgDurationMs), - })), - pricing, - provenance: { - coverage: { - attempts: n(aggregate.attempts), - pricedAttempts: n(aggregate.pricedAttempts), - unpricedAttempts: n(aggregate.unpricedAttempts), - usageReportedAttempts: n(aggregate.usageReportedAttempts), - usagePartialAttempts: n(aggregate.usagePartialAttempts), - usageMissingAttempts: n(aggregate.usageMissingAttempts), - }, - legacyRecords: n(aggregate.legacyRecords), - unreadableRecords: unreadable, - pendingRepairs: pending, - }, + // Host projects identities and freezes replies; returned objects must + // never share mutable rows or arrays with another cached response. + ...structuredClone(statistics), ...activity(db, input.query), }; + if (!hit) { + // A nested read can observe writes that later roll back and reuse the + // same revision counter. Publish only after the outermost commit. + lease.onTransactionSettled((committed) => { + if (committed && !closed) cached = { revision, from, to, statistics }; + }); + } return { kind: 'screen', screen }; }), }; } +function readRangeStatistics( + db: DatabaseSync, + range: UsageScreenQuery['range'], +): RangeStatistics | Extract { + const rangeArgs = [range.from, range.to]; + const modelArgs = [...rangeArgs, ...rangeArgs]; + const aggregate = db + .prepare(`WITH rows AS (${MODEL_ROWS}) SELECT + COUNT(*) AS totalRequests, COALESCE(SUM(costUsd), 0) AS totalCostUsd, + ${['totalTokens', 'inputTokens', 'outputTokens', 'cacheMiss', 'cacheRead', 'cacheCreation', 'reasoning'].map((k) => `COALESCE(SUM(${k}), 0) AS ${k}`).join(', ')}, + SUM(source = 'legacy') AS legacyRecords, + SUM(source = 'canonical') AS attempts, SUM(costBasis = 'priced') AS pricedAttempts, + SUM(costBasis = 'unpriced') AS unpricedAttempts, SUM(usageBasis = 'reported') AS usageReportedAttempts, + SUM(usageBasis = 'partial') AS usagePartialAttempts, SUM(usageBasis = 'missing') AS usageMissingAttempts + FROM rows`) + .get(...modelArgs) as Row; + const breakdown = (column: 'connection' | 'model') => + db + .prepare(`WITH rows AS (${MODEL_ROWS}) + SELECT ${column} AS name, COUNT(*) AS requests, SUM(inputTokens + outputTokens) AS tokens, + COALESCE(SUM(costUsd), 0) AS costUsd FROM rows GROUP BY ${column} ORDER BY requests DESC, name LIMIT 101`) + .all(...modelArgs) as Row[]; + const tools = db + .prepare(`WITH rows AS (${TOOL_ROWS}) SELECT toolName AS tool, + COUNT(*) AS calls, SUM(status = 'success') AS success, SUM(status = 'error') AS errors, + ROUND(AVG(latencyMs)) AS avgDurationMs FROM rows GROUP BY toolName ORDER BY calls DESC, toolName LIMIT 101`) + .all(...rangeArgs) as Row[]; + const providers = breakdown('connection'); + const models = breakdown('model'); + // Read one sentinel past each wire count. It can only become a whole-screen + // failure, never a successful truncated collection. No history array is built. + for (const [rows, section] of [ + [providers, 'provider_breakdown'], + [models, 'model_breakdown'], + [tools, 'tool_breakdown'], + ] as const) { + if (rows.length > 100) return { kind: 'screen_response_too_large', section }; + } + const pricing = db + .prepare('SELECT record_json FROM usage_pricing_overrides ORDER BY model_key LIMIT 129') + .all() + .map((row) => { + const value = JSON.parse(String(row.record_json)); + const separator = value.modelKey.indexOf(':'); + return { + provider: separator < 0 ? '' : value.modelKey.slice(0, separator), + model: separator < 0 ? value.modelKey : value.modelKey.slice(separator + 1), + inputPerMTokUsd: value.inputUsdPer1M, + outputPerMTokUsd: value.outputUsdPer1M, + }; + }); + if (pricing.length > 128) return { kind: 'screen_response_too_large', section: 'pricing' }; + const unreadable = + n( + db + .prepare(`SELECT COUNT(*) AS count FROM usage_model_call_attempts + WHERE completed_at >= ? AND completed_at <= ? AND cost_basis IS NULL`) + .get(...rangeArgs)?.count, + ) + + n( + db + .prepare( + 'SELECT SUM(unreadable_events) AS count FROM usage_model_call_projection_checkpoints', + ) + .get()?.count, + ); + const pending = n( + db + .prepare(`SELECT COUNT(*) AS count FROM core_agent_runs AS source + LEFT JOIN usage_model_call_projection_checkpoints AS checkpoint + ON checkpoint.session_id = source.session_id AND checkpoint.run_id = source.run_id + WHERE source.latest_model_call_sequence > COALESCE(checkpoint.applied_through_sequence, -1)`) + .get()?.count, + ); + return { + summary: { + totalRequests: n(aggregate.totalRequests), + totalCostUsd: n(aggregate.totalCostUsd), + totalTokens: n(aggregate.totalTokens), + inputTokens: n(aggregate.inputTokens), + outputTokens: n(aggregate.outputTokens), + cacheTokens: n(aggregate.cacheRead) + n(aggregate.cacheCreation), + cacheMiss: n(aggregate.cacheMiss), + cacheRead: n(aggregate.cacheRead), + cacheCreation: n(aggregate.cacheCreation), + reasoning: n(aggregate.reasoning), + }, + byProvider: providers.map((row) => ({ + provider: String(row.name), + requests: n(row.requests), + tokens: n(row.tokens), + costUsd: n(row.costUsd), + })), + byModel: models.map((row) => ({ + model: String(row.name), + requests: n(row.requests), + tokens: n(row.tokens), + costUsd: n(row.costUsd), + })), + byTool: tools.map((row) => ({ + tool: String(row.tool), + calls: n(row.calls), + success: n(row.success), + errors: n(row.errors), + avgDurationMs: n(row.avgDurationMs), + })), + pricing, + provenance: { + coverage: { + attempts: n(aggregate.attempts), + pricedAttempts: n(aggregate.pricedAttempts), + unpricedAttempts: n(aggregate.unpricedAttempts), + usageReportedAttempts: n(aggregate.usageReportedAttempts), + usagePartialAttempts: n(aggregate.usagePartialAttempts), + usageMissingAttempts: n(aggregate.usageMissingAttempts), + }, + legacyRecords: n(aggregate.legacyRecords), + unreadableRecords: unreadable, + pendingRepairs: pending, + }, + }; +} + function validateQuery(query: UsageScreenQuery): void { if ( !Number.isSafeInteger(query.range.from) || diff --git a/scripts/perf/usage-statistics.mjs b/scripts/perf/usage-statistics.mjs new file mode 100644 index 0000000000..663385245a --- /dev/null +++ b/scripts/perf/usage-statistics.mjs @@ -0,0 +1,258 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +// Build core/storage at both revisions and pass the baseline checkout. Timed +// requests use the public facade; query probes run outside the timing loop. +import assert from 'node:assert/strict'; +import { execFileSync } from 'node:child_process'; +import { mkdtemp, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { performance } from 'node:perf_hooks'; +import { pathToFileURL } from 'node:url'; +import { acquireOperationalStateDatabase } from '../../packages/storage/dist/operational-state-store.js'; +import { + resolveStorageRoot, + tryAcquireInteractiveRootOwner, + tryAcquireInteractiveRootReader, +} from '../../packages/storage/dist/root-authority.js'; +import { + openInteractiveUsageStoresForRead, + openInteractiveUsageStoresForWrite, +} from '../../packages/storage/dist/usage-stores.js'; +import { removeControlDirectory } from '../../packages/storage/dist/__tests__/fixtures/control-directory-hygiene.js'; +import { report, summarize } from './report.mjs'; + +assert.ok(process.argv[2], 'Pass a built baseline checkout'); +const baselinePath = resolve(process.argv[2]); +const baseline = await import( + pathToFileURL(join(baselinePath, 'packages/storage/dist/usage-stores.js')).href +); +const baselineAuthority = await import( + pathToFileURL(join(baselinePath, 'packages/storage/dist/root-authority.js')).href +); +const baselineDatabase = await import( + pathToFileURL(join(baselinePath, 'packages/storage/dist/operational-state-store.js')).href +); +const rows = []; +const profiles = []; +let sqlite; + +function seed(lease, count, mixed) { + const db = lease.database; + const canonical = db.prepare(`INSERT INTO usage_model_call_attempts + (attempt_id, completed_at, session_id, logical_call_id, turn_id, call_kind, + connection_slug, provider_id, model_id, latency_ms, status, usage_basis, + input_tokens, output_tokens, cost_basis, cost_usd) + VALUES (?, ?, 'session', 'logical', 'turn', 'main', ?, 'provider', ?, 1, + 'completed', 'reported', 10, 2, 'priced', 0.001)`); + const legacy = db.prepare( + 'INSERT INTO usage_llm_calls(storage_key, id, ts, record_json) VALUES (?, ?, ?, ?)', + ); + const tool = db.prepare( + 'INSERT INTO usage_tool_invocations(storage_key, id, ts, record_json) VALUES (?, ?, ?, ?)', + ); + lease.transaction('write', () => { + for (let i = 1; i <= count; i++) { + const id = `row-${String(i).padStart(8, '0')}`; + const model = `model-${i % 20}`; + const connection = `connection-${i % 10}`; + if (mixed && i % 3 === 0) { + canonical.run(id, i, connection, model); + } else { + (mixed && i % 3 === 1 ? legacy : tool).run( + id, + id, + i, + JSON.stringify({ + providerId: 'provider', + connectionSlug: connection, + modelId: model, + toolName: `tool-${i % 8}`, + status: i % 7 === 0 ? 'error' : 'success', + inputTokens: 10, + outputTokens: 2, + totalTokens: 12, + costUsd: 0.001, + latencyMs: 1, + durationMs: 1, + }), + ); + } + } + }); +} + +function capture(db) { + const prepare = db.prepare; + const statements = []; + db.prepare = function (sql) { + const statement = prepare.call(this, sql); + for (const method of ['all', 'get']) { + const run = statement[method].bind(statement); + statement[method] = (...args) => { + statements.push({ sql, args }); + return run(...args); + }; + } + return statement; + }; + return () => { + db.prepare = prepare; + return { + aggregateQueries: statements.filter(({ sql }) => /\bGROUP\s+BY\b|\bSUM\s*\(/i.test(sql)) + .length, + statements, + }; + }; +} + +function comparable(result) { + assert.equal(result.kind, 'screen'); + return { ...result.screen, revision: '' }; +} + +for (const [count, mixed] of [ + [1_000, false], + [10_000, false], + [50_000, false], + [50_000, true], +]) { + const directory = await mkdtemp(join(tmpdir(), 'maka-perf-statistics-')); + const capability = await resolveStorageRoot({ + path: join(directory, 'root'), + kind: 'interactive', + }); + const cleanups = []; + try { + const owner = await tryAcquireInteractiveRootOwner(capability); + assert.ok(owner); + cleanups.push(() => owner.close()); + const writer = await openInteractiveUsageStoresForWrite(owner.lease); + cleanups.push(() => writer.close()); + const seedLease = acquireOperationalStateDatabase(capability.canonicalPath); + try { + seed(seedLease, count, mixed); + } finally { + seedLease.close(); + } + await writer.close(); + await owner.close(); + + const oldCapability = await baselineAuthority.resolveStorageRoot({ + path: capability.canonicalPath, + kind: 'interactive', + }); + const readerOwner = await tryAcquireInteractiveRootReader(capability); + assert.ok(readerOwner); + cleanups.push(() => readerOwner.close()); + const oldOwner = await baselineAuthority.tryAcquireInteractiveRootReader(oldCapability); + assert.ok(oldOwner); + cleanups.push(() => oldOwner.close()); + const stores = await openInteractiveUsageStoresForRead(readerOwner.lease); + cleanups.push(() => stores.close()); + const oldStores = await baseline.openInteractiveUsageStoresForRead(oldOwner.lease); + cleanups.push(() => oldStores.close()); + const lease = acquireOperationalStateDatabase(capability.canonicalPath); + cleanups.push(() => lease.close()); + const oldLease = baselineDatabase.acquireOperationalStateDatabase(capability.canonicalPath); + cleanups.push(() => oldLease.close()); + sqlite = lease.database.prepare('SELECT sqlite_version() AS version').get().version; + + const versions = [ + { name: 'before', read: (input) => oldStores.readUsageScreen(input), db: oldLease.database }, + { name: 'after', read: (input) => stores.readUsageScreen(input), db: lease.database }, + ]; + const scenario = `${count}/${mixed ? 'mixed' : 'tools'}`; + const query = { range: { from: 0, to: count }, search: '', status: 'all' }; + const inputs = [ + { name: 'matching-search', query: { ...query, search: 'model' } }, + { name: 'no-match-search', query: { ...query, search: 'absentword' } }, + { name: 'status-change', query: { ...query, status: 'error' } }, + { name: 'range-change', query: { ...query, range: { from: 1, to: count } } }, + { name: 'revision-change', query, mutate: true }, + ]; + for (const input of inputs) { + const samples = [[], []]; + for (let repetition = -3; repetition < 15; repetition++) { + const results = []; + const order = repetition % 2 === 0 ? [0, 1] : [1, 0]; + for (const i of order) { + await versions[i].read({ kind: 'screen', query }); + } + if (input.mutate) { + // Trigger the existing invalidation authority with a real write, + // outside timing, so both readers observe the same new snapshot. + lease.transaction('write', () => { + lease.database.exec(`UPDATE usage_tool_invocations + SET record_json = json_set(record_json, '$.durationMs', ${repetition + 5}) + WHERE storage_key = 'row-00000002'`); + }); + } + for (const i of order) { + const start = performance.now(); + results[i] = await versions[i].read({ kind: 'screen', query: input.query }); + if (repetition >= 0) samples[i].push(performance.now() - start); + } + assert.deepEqual(comparable(results[0]), comparable(results[1])); + } + for (let i = 0; i < versions.length; i++) { + const timing = summarize(samples[i]); + rows.push({ scenario, metric: `${versions[i].name}/${input.name}-ms`, ...timing }); + console.log( + `${scenario} ${versions[i].name}/${input.name}: ${timing.median.toFixed(3)} ms median, ${timing.p95.toFixed(3)} ms p95`, + ); + } + } + for (const version of versions) { + for (const search of ['model', 'absentword']) { + await version.read({ kind: 'screen', query }); + const finish = capture(version.db); + try { + await version.read({ kind: 'screen', query: { ...query, search } }); + } finally { + profiles.push({ scenario, version: version.name, search, ...finish() }); + } + } + } + } finally { + for (const cleanup of cleanups.reverse()) await cleanup(); + await removeControlDirectory(capability.rootId); + await rm(directory, { recursive: true, force: true }); + } +} +await report( + 'usage-statistics', + { + baselineCommit: execFileSync('git', ['-C', baselinePath, 'rev-parse', 'HEAD'], { + encoding: 'utf8', + }).trim(), + currentDiff: execFileSync('git', ['diff', '--stat'], { encoding: 'utf8' }).trim(), + sqlite, + electron: process.versions.electron ?? null, + warmup: 3, + repetitions: 15, + conditions: + 'Same persisted SQLite fixture, alternating public readUsageScreen facades in one process. Untimed preparation loads the original range, then a timed search/status edit reuses that range; range-change and revision-change requests measure statistics cache misses. SQL probes and fixture writes are outside timing. No cold OS page-cache claim.', + limits: + 'Synthetic local facade measurements exclude Host repair, wire/IPC and rendering. Exact activity count and page SQL are unchanged from main; first-range aggregation and sparse substring searches remain history-sized. A global revision means ongoing writes may prevent cache hits.', + profiles, + }, + rows, +);