From ef987e42d77279aaf28e128159341b787c45645f Mon Sep 17 00:00:00 2001 From: "Mr. Zeng" <16125936+blaksmatic@users.noreply.github.com> Date: Sun, 27 Sep 2026 14:29:19 -0700 Subject: [PATCH] feat(analytics): persist local usage session identities (#18759) Co-authored-by: Neil --- docs/site/content/docs/telemetry.mdx | 4 +- src/main/claude-usage/store.ts | 5 + src/main/codex-usage/store.ts | 5 + src/main/ipc/telemetry.test.ts | 1 + src/main/ipc/telemetry.ts | 1 + src/main/opencode-usage/store.ts | 5 + src/main/telemetry/client.ts | 27 ++- .../usage/agent-token-usage-reporter.test.ts | 151 ++++++++++++++ src/main/usage/agent-token-usage-reporter.ts | 122 +++++++++++ src/main/usage/agent-token-usage.test.ts | 111 ++++++++++ src/main/usage/agent-token-usage.ts | 48 +++++ .../usage/analytics-session-id-store.test.ts | 197 ++++++++++++++++++ src/main/usage/analytics-session-id-store.ts | 88 ++++++++ .../usage-provider-store-lifecycle.test.ts | 69 ++++++ .../usage/usage-provider-store-lifecycle.ts | 52 ++++- .../telemetry-agent-token-usage-schema.ts | 23 ++ src/shared/telemetry-event-registry.ts | 2 + 17 files changed, 900 insertions(+), 11 deletions(-) create mode 100644 src/main/usage/agent-token-usage-reporter.test.ts create mode 100644 src/main/usage/agent-token-usage-reporter.ts create mode 100644 src/main/usage/agent-token-usage.test.ts create mode 100644 src/main/usage/agent-token-usage.ts create mode 100644 src/main/usage/analytics-session-id-store.test.ts create mode 100644 src/main/usage/analytics-session-id-store.ts create mode 100644 src/shared/telemetry-agent-token-usage-schema.ts diff --git a/docs/site/content/docs/telemetry.mdx b/docs/site/content/docs/telemetry.mdx index 604202e067e..f7f8568cd95 100644 --- a/docs/site/content/docs/telemetry.mdx +++ b/docs/site/content/docs/telemetry.mdx @@ -21,12 +21,12 @@ The categories of behavior we observe: - **Lifecycle** — when the app opens. Used to estimate daily, weekly, and monthly active users. - **Repos and workspaces** — when you add a repo or create a workspace. We record _how_ you did it (e.g. folder picker vs. clone URL; command palette vs. drag-and-drop), never the repo name, URL, path, branch name, or any free-form text. -- **Agents** — when you start an agent, which agent kind it was (from a fixed list like `claude-code`, `codex`, `gemini`, etc.), and where you launched it from. Never the prompt, model details, or agent output. +- **Agents** — which agents you use and their token counts. Never prompts or agent output. - **Agent errors** — a coarse error category and which agent kind was involved. We never see raw error messages or stack traces; per-incident detail stays in a local diagnostic trace file on your machine and only reaches Orca if you explicitly share a diagnostic bundle. - **Settings** — when you toggle one of a small whitelisted set of feature-flag or UX preferences. We record which preference changed and whether it's a boolean or an enum, never the raw value of any free-form setting. - **Privacy controls** — when you opt in or out of telemetry, so we can tell from aggregate data whether our consent UI is working. -Every field we transmit is either a fixed enum value, a version string, or the anonymous local ID. No free-form strings from any UI input ever leave your machine. +Fields include fixed enum values, version strings, numeric counts and revisions, and random local IDs. Conversation IDs allow repeated token summaries for the same conversation to be associated; they are pseudonymous. No free-form strings from any UI input ever leave your machine. ## What we never send diff --git a/src/main/claude-usage/store.ts b/src/main/claude-usage/store.ts index f5fd8ca3ce3..686b32ab827 100644 --- a/src/main/claude-usage/store.ts +++ b/src/main/claude-usage/store.ts @@ -1,3 +1,4 @@ +import { claudeTokenSessions } from '../usage/agent-token-usage' import { app } from 'electron' import { join } from 'node:path' import type { @@ -79,6 +80,10 @@ export class ClaudeUsageStore extends UsageProviderStoreLifecycle< > { constructor(store: Pick) { super(store, { + tokenUsage: { + provider: 'claude', + selectSessions: (state) => claudeTokenSessions(state.sessions) + }, logTag: '[claude-usage]', resolveCacheFile: getClaudeUsageFile, createDefaultState: getDefaultState, diff --git a/src/main/codex-usage/store.ts b/src/main/codex-usage/store.ts index f11121f5d0d..7e91fb3945a 100644 --- a/src/main/codex-usage/store.ts +++ b/src/main/codex-usage/store.ts @@ -1,3 +1,4 @@ +import { codexOpenCodeTokenSessions } from '../usage/agent-token-usage' import { app } from 'electron' import { join } from 'node:path' import type { @@ -84,6 +85,10 @@ export class CodexUsageStore extends UsageProviderStoreLifecycle< > { constructor(store: Pick) { super(store, { + tokenUsage: { + provider: 'codex', + selectSessions: (state) => codexOpenCodeTokenSessions(state.sessions) + }, logTag: '[codex-usage]', resolveCacheFile: getCodexUsageFile, createDefaultState: getDefaultState, diff --git a/src/main/ipc/telemetry.test.ts b/src/main/ipc/telemetry.test.ts index 9b28270551a..b49e1882e98 100644 --- a/src/main/ipc/telemetry.test.ts +++ b/src/main/ipc/telemetry.test.ts @@ -161,6 +161,7 @@ describe('telemetry IPC handlers', () => { bucket_source: 'crossed_now' }) handler({}, 'daemon_audit_eligibility', {}) + handler({}, 'agent_token_usage', {}) expect(trackMock).not.toHaveBeenCalled() expect(getCohortAtEmitMock).not.toHaveBeenCalled() }) diff --git a/src/main/ipc/telemetry.ts b/src/main/ipc/telemetry.ts index 767ebed598f..ad954b20ce1 100644 --- a/src/main/ipc/telemetry.ts +++ b/src/main/ipc/telemetry.ts @@ -23,6 +23,7 @@ import type { EventName, EventProps, OptInVia } from '../../shared/telemetry-eve let storeRef: Store | null = null const MAIN_OWNED_TELEMETRY_EVENTS = new Set([ + 'agent_token_usage', 'app_starred_orca', 'daemon_adopted', 'daemon_audit_eligibility', diff --git a/src/main/opencode-usage/store.ts b/src/main/opencode-usage/store.ts index 5dc8a87c67b..3bdcefeb0a7 100644 --- a/src/main/opencode-usage/store.ts +++ b/src/main/opencode-usage/store.ts @@ -1,3 +1,4 @@ +import { codexOpenCodeTokenSessions } from '../usage/agent-token-usage' import { app } from 'electron' import { join } from 'node:path' import type { @@ -50,6 +51,10 @@ export class OpenCodeUsageStore extends UsageProviderStoreLifecycle< > { constructor(store: Pick) { super(store, { + tokenUsage: { + provider: 'opencode', + selectSessions: (state) => codexOpenCodeTokenSessions(state.sessions) + }, logTag: '[opencode-usage]', resolveCacheFile: getOpenCodeUsageFile, createDefaultState: getDefaultState, diff --git a/src/main/telemetry/client.ts b/src/main/telemetry/client.ts index dc92cf1461d..50fb8cfaa3b 100644 --- a/src/main/telemetry/client.ts +++ b/src/main/telemetry/client.ts @@ -160,35 +160,47 @@ function waitForCaptureEnqueue(client: PostHog, event: EventName, uuid: string): }) } +/** Lets producers avoid preparing usage payloads when transmission is disabled. */ +export function isTelemetryEnabled(): boolean { + return ( + (testTransportEnabled || (IS_OFFICIAL_BUILD && TELEMETRY_ENABLED)) && + !shuttingDown && + posthog !== null && + commonProps !== null && + storeRef !== null && + resolveConsent(storeRef.getSettings()).effective === 'enabled' + ) +} + // No-op in contributor / non-official builds; only official stable/rc builds (CI-injected `ORCA_BUILD_IDENTITY` + `ORCA_POSTHOG_WRITE_KEY`) transmit. -export function track(name: N, props: EventProps): void { +export function track(name: N, props: EventProps): boolean { if (!testTransportEnabled && (!IS_OFFICIAL_BUILD || !TELEMETRY_ENABLED)) { - return + return false } // (1) Shutdown gate: late IPC arrivals must not enqueue against a flushing client. if (shuttingDown) { - return + return false } if (!posthog || !commonProps || !storeRef) { - return + return false } // (2) Burst cap before consent: the O(1) cap drops floods before the costly settings read, so a compromised opted-out renderer can't burn CPU. if (!consumeBurstToken(name)) { - return + return false } // (3) Consent resolve — reads live settings every call so it can't drift from persisted state / env-var precedence. const consent = resolveConsent(storeRef.getSettings()) if (consent.effective !== 'enabled') { - return + return false } // (4) Validator — single enforcement point for schema, enum, key set, and length caps. const result = validate(name, props) if (!result.ok) { - return + return false } // (5) Capture. `$process_person_profile: false` stops posthog-node creating a person per install_id (no init-time equivalent). @@ -201,6 +213,7 @@ export function track(name: N, props: EventProps): void $process_person_profile: false } }) + return true } export async function setOptIn(via: OptInVia, optedIn: boolean): Promise { diff --git a/src/main/usage/agent-token-usage-reporter.test.ts b/src/main/usage/agent-token-usage-reporter.test.ts new file mode 100644 index 00000000000..b7cae3de3b5 --- /dev/null +++ b/src/main/usage/agent-token-usage-reporter.test.ts @@ -0,0 +1,151 @@ +import { existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { agentTokenUsageSchema } from '../../shared/telemetry-agent-token-usage-schema' +import { _enableTransportForTests, _setShuttingDownForTests } from '../telemetry/client' +import { + cleanupTelemetryClientTest, + setupTelemetryClientTest, + type TelemetryClientTestState +} from '../telemetry/client-test-harness' +import { resetBurstCapsForSession } from '../telemetry/burst-cap' +import { UsageCacheSnapshotWriter } from '../usage-cache-snapshot-writer' +import { AgentTokenUsageReporter } from './agent-token-usage-reporter' + +const ID = '00000000-0000-4000-8000-000000000001' +const row = { + providerSessionId: 'private-provider-session', + input_tokens: 100, + output_tokens: 20, + cached_input_tokens: 50, + cache_write_input_tokens: 10 +} +let directory: string +let file: string +let telemetry: TelemetryClientTestState +let reporters: AgentTokenUsageReporter[] +const identity = vi.fn(async () => ID) +function reporter(): AgentTokenUsageReporter { + const instance = new AgentTokenUsageReporter(file, 'claude', identity) + reporters.push(instance) + return instance +} +function captures() { + return telemetry.mock.capture.mock.calls.map(([message]) => message.properties) +} +beforeEach(() => { + directory = mkdtempSync(join(tmpdir(), 'orca-token-usage-')) + file = join(directory, 'tokens.json') + telemetry = setupTelemetryClientTest() + reporters = [] + identity.mockReset().mockResolvedValue(ID) +}) +afterEach(async () => { + await Promise.all(reporters.map((instance) => instance.flush())) + cleanupTelemetryClientTest(telemetry.envStash) + vi.restoreAllMocks() + rmSync(directory, { recursive: true, force: true }) +}) + +describe('agent token usage', () => { + it('persists before capture, omits provider IDs, and retains revisions across restarts', async () => { + telemetry.mock.capture.mockImplementation(() => { + expect(JSON.parse(readFileSync(file, 'utf8')).snapshots[0].revision).toBeGreaterThan(0) + }) + const original = reporter() + await original.report([row]) + await original.report([row]) + expect(captures()).toHaveLength(1) + expect(captures()[0]).toMatchObject({ + analytics_session_id: ID, + revision: 1, + input_tokens: 100 + }) + expect(JSON.stringify(captures())).not.toContain(row.providerSessionId) + await original.report([{ ...row, input_tokens: 150 }]) + expect(captures()[1]).toMatchObject({ revision: 2, input_tokens: 150 }) + await reporter().report([{ ...row, input_tokens: 150 }]) + expect(captures()[2]).toMatchObject({ revision: 2, input_tokens: 150 }) + await reporter().report([{ ...row, input_tokens: 90 }]) + expect(captures()[3]).toMatchObject({ revision: 3, input_tokens: 90 }) + }) + + it.each(['opt-out', 'pending', 'environment', 'build', 'shutdown'])( + 'does no identity or disk work for %s', + async (gate) => { + if (gate === 'opt-out' && telemetry.settings.telemetry) { + telemetry.settings.telemetry.optedIn = false + } + if (gate === 'pending' && telemetry.settings.telemetry) { + telemetry.settings.telemetry.optedIn = null + } + if (gate === 'environment') { + process.env.DO_NOT_TRACK = '1' + } + if (gate === 'build') { + _enableTransportForTests(false) + } + if (gate === 'shutdown') { + _setShuttingDownForTests(true) + } + await reporter().report([row]) + expect(identity).not.toHaveBeenCalled() + expect(existsSync(file)).toBe(false) + expect(captures()).toEqual([]) + } + ) + + it('rechecks consent after asynchronous identity persistence', async () => { + identity.mockImplementation(async () => { + if (telemetry.settings.telemetry) { + telemetry.settings.telemetry.optedIn = false + } + return ID + }) + await reporter().report([row]) + expect(captures()).toEqual([]) + }) + + it('does not publish a revision whose write failed and retries it', async () => { + const instance = reporter() + const write = vi + .spyOn(UsageCacheSnapshotWriter.prototype, 'write') + .mockRejectedValueOnce(new Error('disk full')) + await expect(instance.report([row])).rejects.toThrow('disk full') + expect(captures()).toEqual([]) + write.mockRestore() + await instance.report([row]) + expect(captures()[0]).toMatchObject({ revision: 1 }) + }) + + it('retries rate-limited snapshots without incrementing their revision', async () => { + const instance = reporter() + for (let count = 1; count <= 31; count++) { + await instance.report([{ ...row, input_tokens: count }]) + } + expect(captures()).toHaveLength(30) + resetBurstCapsForSession() + await instance.report([{ ...row, input_tokens: 31 }]) + expect(captures().at(-1)).toMatchObject({ revision: 31, input_tokens: 31 }) + }) + + it('rejects corrupt saved revisions without replacing them', async () => { + writeFileSync(file, '{broken') + await expect(reporter().report([row])).rejects.toThrow() + expect(readFileSync(file, 'utf8')).toBe('{broken') + expect(captures()).toEqual([]) + }) + + it('rejects invalid counts and unexpected content fields', async () => { + await reporter().report([{ ...row, input_tokens: -1 }]) + expect(identity).not.toHaveBeenCalled() + const payload = { ...row, analytics_session_id: ID, revision: 1, provider: 'claude' } + expect(agentTokenUsageSchema.safeParse(payload).success).toBe(false) + const { providerSessionId: _, ...valid } = payload + expect(agentTokenUsageSchema.safeParse(valid).success).toBe(true) + for (const field of ['prompt', 'path', 'model', 'activity_timestamp', 'estimated_cost']) { + expect(agentTokenUsageSchema.safeParse({ ...valid, [field]: 'private' }).success).toBe(false) + } + }) +}) diff --git a/src/main/usage/agent-token-usage-reporter.ts b/src/main/usage/agent-token-usage-reporter.ts new file mode 100644 index 00000000000..d79b605500a --- /dev/null +++ b/src/main/usage/agent-token-usage-reporter.ts @@ -0,0 +1,122 @@ +import { readFile } from 'node:fs/promises' +import { z } from 'zod' +import { + agentTokenCountsSchema, + agentTokenUsageSchema, + type AgentTokenUsage +} from '../../shared/telemetry-agent-token-usage-schema' +import { isTelemetryEnabled, track } from '../telemetry/client' +import { UsageCacheSnapshotWriter } from '../usage-cache-snapshot-writer' +import type { AgentTokenSession } from './agent-token-usage' + +const stateSchema = z + .object({ + schemaVersion: z.literal(1), + snapshots: z.array(agentTokenUsageSchema) + }) + .strict() + +/** Cumulative snapshots: consumers select the highest revision per provider/session, never sum retries. */ +export class AgentTokenUsageReporter { + private snapshots: Map | null = null + private readonly sent = new Map() + private readonly writer: UsageCacheSnapshotWriter + private pending: Promise = Promise.resolve() + + constructor( + private readonly file: string, + private readonly provider: AgentTokenUsage['provider'], + private readonly getSessionId: (providerSessionId: string) => Promise + ) { + this.writer = new UsageCacheSnapshotWriter('[agent-token-usage]', () => file) + } + + report(sessions: AgentTokenSession[]): Promise { + const operation = this.pending.then(() => this.publish(sessions)) + this.pending = operation.catch(() => {}) + return operation + } + + async flush(): Promise { + await this.pending + await this.writer.flush() + } + + private async publish(sessions: AgentTokenSession[]): Promise { + if (!isTelemetryEnabled()) { + return + } + const previous = this.snapshots ?? (await this.load()) + const next = new Map(previous) + let changed = false + for (const { providerSessionId, ...counts } of sessions) { + if (!isTelemetryEnabled()) { + return + } + const parsed = agentTokenCountsSchema.safeParse(counts) + if (!parsed.success) { + continue + } + const id = await this.getSessionId(providerSessionId) + const existing = next.get(id) + if ( + existing && + existing.input_tokens === counts.input_tokens && + existing.output_tokens === counts.output_tokens && + existing.cached_input_tokens === counts.cached_input_tokens && + existing.cache_write_input_tokens === counts.cache_write_input_tokens + ) { + continue + } + next.set( + id, + agentTokenUsageSchema.parse({ + ...parsed.data, + provider: this.provider, + analytics_session_id: id, + revision: (existing?.revision ?? 0) + 1 + }) + ) + changed = true + } + // Persist revisions before capture so a restart can only retry the same snapshot or a newer one. + if (changed) { + await this.writer.write(() => + JSON.stringify({ schemaVersion: 1, snapshots: [...next.values()] }) + ) + } + this.snapshots = next + for (const snapshot of next.values()) { + if (this.sent.get(snapshot.analytics_session_id) === snapshot.revision) { + continue + } + if (!track('agent_token_usage', snapshot)) { + break + } + this.sent.set(snapshot.analytics_session_id, snapshot.revision) + } + } + + private async load(): Promise> { + let content: string + try { + content = await readFile(this.file, 'utf8') + } catch (error) { + if (error && typeof error === 'object' && 'code' in error && error.code === 'ENOENT') { + return new Map() + } + throw error + } + const state = stateSchema.parse(JSON.parse(content)) + const snapshots = new Map( + state.snapshots.map((snapshot) => [snapshot.analytics_session_id, snapshot]) + ) + if ( + snapshots.size !== state.snapshots.length || + state.snapshots.some((snapshot) => snapshot.provider !== this.provider) + ) { + throw new Error('Invalid agent token usage state') + } + return snapshots + } +} diff --git a/src/main/usage/agent-token-usage.test.ts b/src/main/usage/agent-token-usage.test.ts new file mode 100644 index 00000000000..4503dcabac8 --- /dev/null +++ b/src/main/usage/agent-token-usage.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, it } from 'vitest' +import type { ClaudeUsageSession } from '../claude-usage/types' +import type { CodexUsageSession } from '../codex-usage/types' +import type { OpenCodeUsageSession } from '../opencode-usage/types' +import { claudeTokenSessions, codexOpenCodeTokenSessions } from './agent-token-usage' + +const claudeLocation = { + locationKey: '/private/path', + projectLabel: 'private-repo', + repoId: 'repo', + worktreeId: 'folder', + turnCount: 1, + inputTokens: 100, + outputTokens: 20, + cacheReadTokens: 30, + cacheWriteTokens: 10, + cacheWrite1hTokens: 5 +} +const claude: ClaudeUsageSession = { + sessionId: 'session', + firstTimestamp: 'private-start', + lastTimestamp: 'private-end', + model: 'private-model', + lastCwd: '/private/path', + lastGitBranch: 'private-branch', + primaryWorktreeId: 'folder', + primaryRepoId: 'repo', + turnCount: 2, + totalInputTokens: 200, + totalOutputTokens: 40, + totalCacheReadTokens: 60, + totalCacheWriteTokens: 20, + totalCacheWrite1hTokens: 10, + locationBreakdown: [claudeLocation, { ...claudeLocation, worktreeId: null }] +} +const codexLocation = { + locationKey: '/private/path', + projectLabel: 'private-repo', + repoId: 'repo', + worktreeId: 'folder', + eventCount: 1, + inputTokens: 100, + cachedInputTokens: 30, + outputTokens: 20, + reasoningOutputTokens: 5, + totalTokens: 120, + hasInferredPricing: false, + estimatedCostUsd: 1 +} +const codex: CodexUsageSession = { + sessionId: 'session', + firstTimestamp: 'private-start', + lastTimestamp: 'private-end', + primaryModel: 'private-model', + hasMixedModels: false, + primaryProjectLabel: 'private-repo', + hasMixedLocations: true, + primaryWorktreeId: 'folder', + primaryRepoId: 'repo', + eventCount: 2, + totalInputTokens: 200, + totalCachedInputTokens: 60, + totalOutputTokens: 40, + totalReasoningOutputTokens: 10, + totalTokens: 240, + hasInferredPricing: false, + locationBreakdown: [codexLocation, { ...codexLocation, worktreeId: null }], + modelBreakdown: [], + locationModelBreakdown: [] +} +const opencode: OpenCodeUsageSession = { + ...codex, + estimatedCostUsd: 2, + locationBreakdown: [codexLocation, { ...codexLocation, worktreeId: null }], + modelBreakdown: [], + locationModelBreakdown: [] +} + +describe('Orca token projections', () => { + it('counts Claude cache writes once including the 1-hour subset and excludes outside usage', () => { + expect(claudeTokenSessions([claude])).toEqual([ + { + providerSessionId: 'session', + input_tokens: 100, + output_tokens: 20, + cached_input_tokens: 30, + cache_write_input_tokens: 10 + } + ]) + }) + it.each([codex, opencode])('separates cache reads from inclusive input counts', (session) => { + expect(codexOpenCodeTokenSessions([session])).toEqual([ + { + providerSessionId: 'session', + input_tokens: 70, + output_tokens: 20, + cached_input_tokens: 30, + cache_write_input_tokens: 0 + } + ]) + }) + it('does not fall back to unscoped session totals when attribution is missing', () => { + expect(claudeTokenSessions([{ ...claude, locationBreakdown: [] }])).toEqual([]) + expect(codexOpenCodeTokenSessions([{ ...codex, locationBreakdown: [] }])).toEqual([]) + expect( + codexOpenCodeTokenSessions([ + { ...codex, locationBreakdown: [{ ...codexLocation, worktreeId: null }] } + ]) + ).toEqual([]) + }) +}) diff --git a/src/main/usage/agent-token-usage.ts b/src/main/usage/agent-token-usage.ts new file mode 100644 index 00000000000..36cada6cf01 --- /dev/null +++ b/src/main/usage/agent-token-usage.ts @@ -0,0 +1,48 @@ +import type { AgentTokenCounts } from '../../shared/telemetry-agent-token-usage-schema' +import type { ClaudeUsageSession } from '../claude-usage/types' +import type { CodexUsageSession } from '../codex-usage/types' +import type { OpenCodeUsageSession } from '../opencode-usage/types' + +export type AgentTokenSession = AgentTokenCounts & { providerSessionId: string } + +export function claudeTokenSessions(sessions: ClaudeUsageSession[]): AgentTokenSession[] { + return sessions.flatMap((session) => { + const locations = session.locationBreakdown.filter((entry) => entry.worktreeId !== null) + if (locations.length === 0) { + return [] + } + return [ + { + providerSessionId: session.sessionId, + input_tokens: locations.reduce((sum, entry) => sum + entry.inputTokens, 0), + output_tokens: locations.reduce((sum, entry) => sum + entry.outputTokens, 0), + cached_input_tokens: locations.reduce((sum, entry) => sum + entry.cacheReadTokens, 0), + cache_write_input_tokens: locations.reduce((sum, entry) => sum + entry.cacheWriteTokens, 0) + } + ] + }) +} + +export function codexOpenCodeTokenSessions( + sessions: (CodexUsageSession | OpenCodeUsageSession)[] +): AgentTokenSession[] { + return sessions.flatMap((session) => { + const locations = session.locationBreakdown.filter((entry) => entry.worktreeId !== null) + if (locations.length === 0) { + return [] + } + return [ + { + providerSessionId: session.sessionId, + // Codex/OpenCode include cache hits in input; Claude records them separately. + input_tokens: locations.reduce( + (sum, entry) => sum + entry.inputTokens - entry.cachedInputTokens, + 0 + ), + output_tokens: locations.reduce((sum, entry) => sum + entry.outputTokens, 0), + cached_input_tokens: locations.reduce((sum, entry) => sum + entry.cachedInputTokens, 0), + cache_write_input_tokens: 0 + } + ] + }) +} diff --git a/src/main/usage/analytics-session-id-store.test.ts b/src/main/usage/analytics-session-id-store.test.ts new file mode 100644 index 00000000000..3ec5addfad2 --- /dev/null +++ b/src/main/usage/analytics-session-id-store.test.ts @@ -0,0 +1,197 @@ +import { existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs' +import type * as FsPromises from 'node:fs/promises' +import { join } from 'node:path' +import { tmpdir } from 'node:os' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { AnalyticsSessionIdStore } from './analytics-session-id-store' + +type WriteGate = { entered: (() => void) | null; wait: Promise | null; fail: boolean } + +const { writeGate } = vi.hoisted((): { writeGate: WriteGate } => ({ + writeGate: { + entered: null, + wait: null, + fail: false + } +})) +vi.mock('node:fs/promises', async () => { + const actual = await vi.importActual('node:fs/promises') + return { + ...actual, + open: async (...args: Parameters) => { + if (args[1] === 'w') { + writeGate.entered?.() + if (writeGate.wait) { + await writeGate.wait + } + if (writeGate.fail) { + throw new Error('simulated disk failure') + } + } + return actual.open(...args) + } + } +}) + +const UUID_V4 = /^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/ +let directory: string +let file: string +let stores: AnalyticsSessionIdStore[] +/** Registers every instance for flush-before-cleanup, including simulated restarts. */ +function createStore(path = file): AnalyticsSessionIdStore { + const store = new AnalyticsSessionIdStore(path) + stores.push(store) + return store +} +beforeEach(() => { + directory = mkdtempSync(join(tmpdir(), 'orca-analytics-session-id-')) + file = join(directory, 'identities.json') + stores = [] + writeGate.entered = null + writeGate.wait = null + writeGate.fail = false +}) +afterEach(async () => { + await Promise.all(stores.map((store) => store.flush())) + rmSync(directory, { recursive: true, force: true }) + vi.restoreAllMocks() +}) + +describe('AnalyticsSessionIdStore', () => { + it('is lazy and persists a random ID before returning it', async () => { + const store = createStore() + expect(existsSync(file)).toBe(false) + const id = await store.getOrCreate('provider-session-123') + expect(id).toMatch(UUID_V4) + expect(id).not.toContain('provider-session') + expect(JSON.parse(readFileSync(file, 'utf8'))).toEqual({ + schemaVersion: 1, + entries: [['provider-session-123', id]] + }) + expect(await store.getOrCreate('provider-session-123')).toBe(id) + }) + + it('restores the ID after a restart or session resume', async () => { + const original = createStore() + const id = await original.getOrCreate('resumed-session') + await original.flush() + expect(await createStore().getOrCreate('resumed-session')).toBe(id) + }) + + it('handles concurrent requests for the same and different sessions without lost writes', async () => { + const store = createStore() + const ids = await Promise.all( + Array.from({ length: 20 }, (_, index) => store.getOrCreate(`session-${index % 10}`)) + ) + expect(new Set(ids).size).toBe(10) + expect(ids.slice(0, 10)).toEqual(ids.slice(10)) + const restored = createStore() + for (let index = 0; index < 10; index++) { + expect(await restored.getOrCreate(`session-${index}`)).toBe(ids[index]) + } + }) + + it('isolates identical provider session IDs across provider and host/profile files', async () => { + const ids = await Promise.all([ + createStore(join(directory, 'host-a', 'claude.json')).getOrCreate('same-session'), + createStore(join(directory, 'host-a', 'codex.json')).getOrCreate('same-session'), + createStore(join(directory, 'host-b', 'claude.json')).getOrCreate('same-session') + ]) + expect(new Set(ids).size).toBe(3) + }) + + it.each(['', ' ', 'x'.repeat(1025)])( + 'rejects invalid keys without creating a mapping', + async (key) => { + await expect(createStore().getOrCreate(key)).rejects.toThrow('Invalid provider session ID') + expect(existsSync(file)).toBe(false) + } + ) + + it('handles prototype names and delimiter characters as ordinary opaque keys', async () => { + const store = createStore() + const keys = ['__proto__', 'constructor', 'a::b/c', 'a/b::c'] + const ids = await Promise.all(keys.map((key) => store.getOrCreate(key))) + expect(new Set(ids).size).toBe(keys.length) + const restored = createStore() + expect(await Promise.all(keys.map((key) => restored.getOrCreate(key)))).toEqual(ids) + }) + + it.each([ + '{broken', + JSON.stringify({ schemaVersion: 2, entries: [] }), + JSON.stringify({ schemaVersion: 1, entries: [['session', 'provider-id']] }), + JSON.stringify({ schemaVersion: 1, entries: [], extra: 'secret' }) + ])('fails closed on invalid persisted state without overwriting it', async (content) => { + writeFileSync(file, content) + await expect(createStore().getOrCreate('new-session')).rejects.toThrow( + 'Invalid analytics session identity file' + ) + expect(readFileSync(file, 'utf8')).toBe(content) + }) + + it('rejects conflicting provider keys and analytics IDs', async () => { + const first = '00000000-0000-4000-8000-000000000001' + const second = '00000000-0000-4000-8000-000000000002' + for (const entries of [ + [ + ['a', first], + ['a', second] + ], + [ + ['a', first], + ['b', first] + ] + ]) { + writeFileSync(file, JSON.stringify({ schemaVersion: 1, entries })) + await expect(createStore().getOrCreate('new-session')).rejects.toThrow( + 'Duplicate analytics session identity' + ) + } + }) + + it('keeps all callers and shutdown waiting until the ID write completes', async () => { + let release!: () => void + let entered!: () => void + writeGate.wait = new Promise((resolve) => { + release = resolve + }) + const writing = new Promise((resolve) => { + entered = resolve + }) + writeGate.entered = entered + const store = createStore() + let resolved = false + let flushed = false + const first = store.getOrCreate('session').then((id) => { + resolved = true + return id + }) + await writing + const second = store.getOrCreate('session') + const flush = store.flush().then(() => { + flushed = true + }) + try { + await Promise.resolve() + expect(resolved).toBe(false) + expect(flushed).toBe(false) + } finally { + release() + } + expect(await first).toBe(await second) + await flush + expect(flushed).toBe(true) + }) + + it('does not return an unpersisted ID on write failure and allows retry', async () => { + vi.spyOn(console, 'error').mockImplementation(() => {}) + const store = createStore() + writeGate.fail = true + await expect(store.getOrCreate('session')).rejects.toThrow('simulated disk failure') + expect(existsSync(file)).toBe(false) + writeGate.fail = false + const id = await store.getOrCreate('session') + expect(await createStore().getOrCreate('session')).toBe(id) + }) +}) diff --git a/src/main/usage/analytics-session-id-store.ts b/src/main/usage/analytics-session-id-store.ts new file mode 100644 index 00000000000..0b921e8efaa --- /dev/null +++ b/src/main/usage/analytics-session-id-store.ts @@ -0,0 +1,88 @@ +import { randomUUID } from 'node:crypto' +import { readFile } from 'node:fs/promises' +import { z } from 'zod' +import { UsageCacheSnapshotWriter } from '../usage-cache-snapshot-writer' + +const providerSessionIdSchema = z + .string() + .min(1) + .max(1024) + .refine((id) => id.trim().length > 0) +const analyticsSessionIdSchema = z.uuidv4().brand<'AnalyticsSessionId'>() +export type AnalyticsSessionId = z.infer +const identityFileSchema = z + .object({ + schemaVersion: z.literal(1), + entries: z.array(z.tuple([providerSessionIdSchema, analyticsSessionIdSchema])) + }) + .strict() + +/** One owner per file, scoped to a provider's usage store on its execution host. */ +export class AnalyticsSessionIdStore { + private identities: Map | null = null + private pending: Promise = Promise.resolve() + private writer: UsageCacheSnapshotWriter | null = null + + constructor(private readonly file: string) {} + + /** Resolves only after persistence succeeds, so an uploaded ID survives a restart. */ + async getOrCreate(providerSessionId: string): Promise { + if (!providerSessionIdSchema.safeParse(providerSessionId).success) { + throw new Error('Invalid provider session ID') + } + // Serialize reads as well as writes so no caller sees an ID before it is durable. + const operation = this.pending.then(async () => { + const identities = this.identities ?? (await this.load()) + this.identities = identities + const existing = identities.get(providerSessionId) + if (existing) { + return existing + } + const id = analyticsSessionIdSchema.parse(randomUUID()) + const updated = new Map(identities).set(providerSessionId, id) + this.writer ??= new UsageCacheSnapshotWriter('[analytics-session-id]', () => this.file) + await this.writer.write(() => JSON.stringify({ schemaVersion: 1, entries: [...updated] })) + this.identities = updated + return id + }) + this.pending = operation.then( + () => {}, + () => {} + ) + return operation + } + + /** Includes queued lookups that have not reached the durable writer yet. */ + async flush(): Promise { + await this.pending + await this.writer?.flush() + } + + private async load(): Promise> { + let content: string + try { + content = await readFile(this.file, 'utf8') + } catch (error) { + if (error && typeof error === 'object' && 'code' in error && error.code === 'ENOENT') { + return new Map() + } + throw error + } + let raw: unknown + try { + raw = JSON.parse(content) + } catch { + throw new Error('Invalid analytics session identity file') + } + const parsed = identityFileSchema.safeParse(raw) + if (!parsed.success) { + throw new Error('Invalid analytics session identity file') + } + const identities = new Map(parsed.data.entries) + const analyticsIds = new Set(parsed.data.entries.map(([, id]) => id)) + if (identities.size !== parsed.data.entries.length || analyticsIds.size !== identities.size) { + throw new Error('Duplicate analytics session identity') + } + return identities + } +} diff --git a/src/main/usage/usage-provider-store-lifecycle.test.ts b/src/main/usage/usage-provider-store-lifecycle.test.ts index 02a5571363d..ef8c820052a 100644 --- a/src/main/usage/usage-provider-store-lifecycle.test.ts +++ b/src/main/usage/usage-provider-store-lifecycle.test.ts @@ -1,3 +1,7 @@ +import { + setupTelemetryClientTest, + cleanupTelemetryClientTest +} from '../telemetry/client-test-harness' import { existsSync, mkdtempSync, readFileSync, readdirSync, rmSync, writeFileSync } from 'node:fs' import type * as FsPromises from 'node:fs/promises' import { tmpdir } from 'node:os' @@ -106,6 +110,17 @@ class TestUsageStore extends UsageProviderStoreLifecycle< getAllWorktreeMeta: () => ({}) }, { + tokenUsage: { + provider: 'claude', + selectSessions: (state) => + state.sessions.map((session) => ({ + providerSessionId: session.id, + input_tokens: 10, + output_tokens: 2, + cached_input_tokens: 3, + cache_write_input_tokens: 1 + })) + }, logTag: '[test-usage]', resolveCacheFile: () => cacheFile, createDefaultState: makeState, @@ -156,6 +171,60 @@ describe('UsageProviderStoreLifecycle', () => { vi.restoreAllMocks() }) + it('reports enabled scans and preserves revisions when the usage cache is rebuilt', async () => { + const telemetry = setupTelemetryClientTest() + try { + const cacheFile = join(tempDirectory, 'provider.json') + scan.mockResolvedValue({ ...emptyScanResult(), sessions: [{ id: 'provider-session' }] }) + const original = createStore(cacheFile) + await original.refresh(true) + expect(telemetry.mock.capture).not.toHaveBeenCalled() + await original.setEnabled(true) + await original.refresh(true) + const first = telemetry.mock.capture.mock.calls[0]?.[0] + expect(first).toMatchObject({ + event: 'agent_token_usage', + properties: { revision: 1, input_tokens: 10 } + }) + await original.flush() + rmSync(cacheFile) + const rebuilt = createStore(cacheFile) + await rebuilt.setEnabled(true) + await rebuilt.refresh(true) + expect(telemetry.mock.capture.mock.calls[1]?.[0]).toEqual(first) + } finally { + cleanupTelemetryClientTest(telemetry.envStash) + } + }) + + it('keeps analytics identity separate from usage cache rebuilds and never puts it in snapshots', async () => { + const cacheFile = join(tempDirectory, 'provider.json') + const identityFile = join(tempDirectory, 'provider-analytics-session-ids.json') + const original = createStore(cacheFile) + expect(existsSync(identityFile)).toBe(false) + const id = await original.getAnalyticsSessionId('provider-session') + expect(existsSync(identityFile)).toBe(true) + await original.setEnabled(true) + await original.refresh(true) + await original.flush() + expect(JSON.stringify(original.getState())).not.toContain(id) + expect(readFileSync(cacheFile, 'utf8')).not.toContain(id) + rmSync(cacheFile) + const rebuilt = createStore(cacheFile) + expect(await rebuilt.getAnalyticsSessionId('provider-session')).toBe(id) + expect(await createStore().getAnalyticsSessionId('provider-session')).not.toBe(id) + }) + + it('flush waits for queued analytics identity creation', async () => { + const store = createStore() + const identity = store.getAnalyticsSessionId('session') + await store.flush() + const id = await identity + expect( + readFileSync(join(tempDirectory, 'usage-0-analytics-session-ids.json'), 'utf8') + ).toContain(id) + }) + it('skips disabled and fresh matching states', async () => { const store = createStore() diff --git a/src/main/usage/usage-provider-store-lifecycle.ts b/src/main/usage/usage-provider-store-lifecycle.ts index 850d9e61bda..bbfd4be8e13 100644 --- a/src/main/usage/usage-provider-store-lifecycle.ts +++ b/src/main/usage/usage-provider-store-lifecycle.ts @@ -1,4 +1,10 @@ import { existsSync, readFileSync } from 'node:fs' +import { AgentTokenUsageReporter } from './agent-token-usage-reporter' +import type { AgentTokenSession } from './agent-token-usage' +import type { AgentTokenUsage } from '../../shared/telemetry-agent-token-usage-schema' +import { isTelemetryEnabled } from '../telemetry/client' +import { join, parse } from 'node:path' +import { AnalyticsSessionIdStore, type AnalyticsSessionId } from './analytics-session-id-store' import type { Store } from '../persistence' import { UsageCacheSnapshotWriter } from '../usage-cache-snapshot-writer' import { loadKnownUsageWorktreesByRepo } from '../usage-worktree-metadata' @@ -38,6 +44,10 @@ type UsageProviderStoreLifecycleConfig< sourceKey: SourceKey dataPresenceKey: DataPresenceKey jsonIndent?: number + tokenUsage?: { + provider: AgentTokenUsage['provider'] + selectSessions: (state: State) => AgentTokenSession[] + } scan: ( worktrees: UsageScanWorktreeRef[], previous: State[SourceKey] @@ -55,6 +65,8 @@ export abstract class UsageProviderStoreLifecycle< > { protected state: State private scanPromise: Promise | null = null + private tokenReporter: AgentTokenUsageReporter | null = null + private analyticsSessionIds: AnalyticsSessionIdStore | null = null private readonly writer: UsageCacheSnapshotWriter constructor( @@ -74,9 +86,21 @@ export abstract class UsageProviderStoreLifecycle< } as PublicUsageProviderScanState } + /** Local identity only; callers must use the usage store on the execution host. */ + getAnalyticsSessionId(providerSessionId: string): Promise { + if (!this.analyticsSessionIds) { + const { dir, name } = parse(this.config.resolveCacheFile()) + this.analyticsSessionIds = new AnalyticsSessionIdStore( + join(dir, `${name}-analytics-session-ids.json`) + ) + } + return this.analyticsSessionIds.getOrCreate(providerSessionId) + } + /** Await queued cache writes so quit does not drop the final snapshot. */ - flush(): Promise { - return this.writer.flush() + async flush(): Promise { + await this.tokenReporter?.flush() + await Promise.all([this.writer.flush(), this.analyticsSessionIds?.flush()]) } async setEnabled(enabled: boolean): Promise> { @@ -152,6 +176,7 @@ export abstract class UsageProviderStoreLifecycle< this.state.scanState.lastScanError = null // Persistence failures do not turn a successful source scan into a scan failure. await this.writeToDisk().catch(() => {}) + await this.reportTokenUsage() } catch (error) { this.state.scanState.lastScanError = error instanceof Error ? error.message : String(error) await this.writeToDisk().catch(() => {}) @@ -163,6 +188,29 @@ export abstract class UsageProviderStoreLifecycle< await this.scanPromise } + private async reportTokenUsage(): Promise { + const config = this.config.tokenUsage + if (!config || !isTelemetryEnabled() || !this.state.scanState.enabled) { + return + } + try { + if (!this.tokenReporter) { + const { dir, name } = parse(this.config.resolveCacheFile()) + this.tokenReporter = new AgentTokenUsageReporter( + join(dir, `${name}-token-usage.json`), + config.provider, + (id) => this.getAnalyticsSessionId(id) + ) + } + await this.tokenReporter.report(config.selectSessions(this.state)) + } catch { + // Reporting failures must not invalidate a successful local usage scan. + console.warn( + '[agent-token-usage] Could not report token usage; will retry after the next scan' + ) + } + } + private async getCurrentWorktreeFingerprint(): Promise { const repos = this.store.getRepos() return getUsageWorktreeFingerprint(loadKnownUsageWorktreesByRepo(this.store, repos)) diff --git a/src/shared/telemetry-agent-token-usage-schema.ts b/src/shared/telemetry-agent-token-usage-schema.ts new file mode 100644 index 00000000000..18a3d4a08d5 --- /dev/null +++ b/src/shared/telemetry-agent-token-usage-schema.ts @@ -0,0 +1,23 @@ +import { z } from 'zod' + +const tokenCount = z.number().int().nonnegative().max(Number.MAX_SAFE_INTEGER) + +export const agentTokenCountsSchema = z + .object({ + input_tokens: tokenCount, + output_tokens: tokenCount, + cached_input_tokens: tokenCount, + cache_write_input_tokens: tokenCount + }) + .strict() + +export const agentTokenUsageSchema = agentTokenCountsSchema + .extend({ + provider: z.enum(['claude', 'codex', 'opencode']), + analytics_session_id: z.uuidv4(), + revision: z.number().int().positive().max(Number.MAX_SAFE_INTEGER) + }) + .strict() + +export type AgentTokenCounts = z.infer +export type AgentTokenUsage = z.infer diff --git a/src/shared/telemetry-event-registry.ts b/src/shared/telemetry-event-registry.ts index cba4ad9a208..f24e60eccd2 100644 --- a/src/shared/telemetry-event-registry.ts +++ b/src/shared/telemetry-event-registry.ts @@ -1,3 +1,4 @@ +import { agentTokenUsageSchema } from './telemetry-agent-token-usage-schema' import { agentErrorSchema, agentPromptSentSchema, @@ -116,6 +117,7 @@ export const eventSchemas = { setup_script_prompt_shown: setupScriptPromptShownSchema, setup_script_prompt_action: setupScriptPromptActionSchema, + agent_token_usage: agentTokenUsageSchema, agent_started: agentStartedSchema, agent_prompt_sent: agentPromptSentSchema, agent_error: agentErrorSchema,