diff --git a/apps/mobile/src/features/usage/UsageRouteScreen.tsx b/apps/mobile/src/features/usage/UsageRouteScreen.tsx index e2022058616b..99ac66db0fd2 100644 --- a/apps/mobile/src/features/usage/UsageRouteScreen.tsx +++ b/apps/mobile/src/features/usage/UsageRouteScreen.tsx @@ -1,8 +1,13 @@ import { ChatGptUsageSummary } from "./ChatGptUsageSummary"; import { ScreenScrollView as ScrollView } from "../../components/ScreenScrollView"; -import { EnvironmentId, USAGE_CONTRACT_VERSION } from "@t3tools/contracts"; +import { EnvironmentId, USAGE_CONTRACT_VERSION, type UsageProviderKind } from "@t3tools/contracts"; import { type RouteProp, useIsFocused, useNavigation, useRoute } from "@react-navigation/native"; import { cursorKeychainAccessEnvironments } from "@t3tools/client-runtime/state/usage"; +import { + updatingProvidersLabel, + usageEnvironmentProgress, + usageProgress, +} from "@t3tools/client-runtime/state/usage-progress"; import { isCompatibleUsageContractVersion, isModelCostUnknown, @@ -21,9 +26,16 @@ import { formatUsd, makeWindow, } from "@t3tools/shared/usageFormat"; -import { useCallback, useLayoutEffect, useMemo, useRef, useState } from "react"; -import { Platform, Pressable, RefreshControl, View } from "react-native"; -import Animated, { FadeIn, ReduceMotion } from "react-native-reanimated"; +import { useCallback, useLayoutEffect, useMemo, useRef, useState, type ReactNode } from "react"; +import { ActivityIndicator, Platform, Pressable, RefreshControl, View } from "react-native"; +import Animated, { + FadeIn, + ReduceMotion, + useAnimatedStyle, + useSharedValue, + withDelay, + withTiming, +} from "react-native-reanimated"; import { useSafeAreaInsets } from "react-native-safe-area-context"; import { SegmentedControl } from "../../components/SegmentedControl"; @@ -65,6 +77,7 @@ const METRIC_OPTIONS = [ ] as const satisfies readonly { value: UsageChartMetric; label: string }[]; const CHART_HEIGHT = 180; +const providerLabel = (provider: UsageProviderKind) => PROVIDER_LABEL[provider]; const CURSOR_KEYCHAIN_COPY = "Requires access to your Cursor login in macOS Keychain."; /** @@ -156,6 +169,10 @@ export function UsageRouteScreen() { const [refreshingUsage, setRefreshingUsage] = useState(false); const refreshingRef = useRef(false); const showingLimits = tab === "limits"; + const progress = usageProgress(selectedEnvironments, { + refreshing: refreshingUsage, + providerLabel, + }); const selectWindow = (days: number) => { setWindowSelection({ days, @@ -182,10 +199,6 @@ export function UsageRouteScreen() { }; const showEnvironmentFilter = environments.length > 0 || selectedEnvironmentIds !== null; - const hasLoadingEnvironments = selectedEnvironments.some(isUsageLoading); - const filterAccessibilityLabel = hasLoadingEnvironments - ? "Filter usage environments, some environments are loading" - : "Filter usage environments"; const filterIcon = selectedEnvironmentIds === null ? "line.3.horizontal.decrease" @@ -201,14 +214,14 @@ export function UsageRouteScreen() { ...environments.map((environment) => ({ id: environment.environmentId, title: environment.label, - subtitle: usageEnvironmentStatus(environment), + subtitle: usageEnvironmentStatus(environment, refreshingUsage), state: selectedEnvironmentIds === null || selectedEnvironmentIds.has(environment.environmentId) ? ("on" as const) : ("off" as const), })), ], - [environments, selectedEnvironmentIds], + [environments, refreshingUsage, selectedEnvironmentIds], ); const selectEnvironment = useCallback( (value: string) => { @@ -227,37 +240,24 @@ export function UsageRouteScreen() { selectEnvironment(nativeEvent.event)} > - {hasLoadingEnvironments ? ( - - ) : null} ) : null, - [ - showEnvironmentFilter, - environmentActions, - selectEnvironment, - filterAccessibilityLabel, - filterIcon, - hasLoadingEnvironments, - ], + [showEnvironmentFilter, environmentActions, selectEnvironment, filterIcon], ); useLayoutEffect(() => { @@ -361,26 +361,28 @@ export function UsageRouteScreen() { {message} ))} - - 1} - onCursorEnabled={refreshAfterCursorEnable} - /> - - - + + + 1} + onCursorEnabled={refreshAfterCursorEnable} + /> + + + + )} @@ -391,6 +393,53 @@ export function UsageRouteScreen() { ); } +/** + * Dims totals that are about to change and says what is still updating. The + * status overlays the first child's top-right corner, the chart card's label + * row, so appearing never moves anything. + */ +function UsageUpdating({ + dimmed, + label, + children, +}: { + readonly dimmed: boolean; + readonly label: string | null; + readonly children: ReactNode; +}) { + const opacity = useSharedValue(1); + useLayoutEffect(() => { + // The delay keeps a quick cached answer from flashing the dim. + opacity.set( + dimmed + ? withDelay(150, withTiming(0.5, { duration: 150, reduceMotion: ReduceMotion.System })) + : withTiming(1, { duration: 150, reduceMotion: ReduceMotion.System }), + ); + }, [opacity, dimmed]); + const dimStyle = useAnimatedStyle(() => ({ opacity: opacity.get() })); + + return ( + + + {children} + + {label !== null ? ( + + + + {label} + + + ) : null} + + ); +} + function CursorEnableAction({ environmentId, label, @@ -860,11 +909,7 @@ function ModelsSection(props: { readonly merged: MergedUsage; readonly metric: U * one that failed, or one whose transcripts another environment already * reported. */ -function isUsageLoading(environment: EnvironmentUsageStatus) { - return environment.isPending || (environment.summary === null && environment.error === null); -} - -function usageEnvironmentStatus(environment: EnvironmentUsageStatus): string { +function usageEnvironmentStatus(environment: EnvironmentUsageStatus, refreshing: boolean): string { if ( environment.summary && !isCompatibleUsageContractVersion(environment.summary.contractVersion, USAGE_CONTRACT_VERSION) @@ -881,7 +926,10 @@ function usageEnvironmentStatus(environment: EnvironmentUsageStatus): string { return environment.summary ? `${environment.error} Showing saved totals.` : environment.error; if (!environment.isConnected) return environment.summary ? "Disconnected · showing saved usage" : "Waiting for connection…"; - if (isUsageLoading(environment)) - return environment.summary ? "Updating usage…" : "Loading usage…"; + const progress = usageEnvironmentProgress(environment, refreshing); + if (progress.phase === "loading") return "Loading usage…"; + if (progress.phase === "stale") return "Updating usage…"; + if (progress.phase === "partway") + return updatingProvidersLabel(progress.providers, providerLabel); return "Usage up to date"; } diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 192b17204aea..8c50ab6258ff 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -153,6 +153,7 @@ import * as NativeTelemetryClient from "./resourceTelemetry/NativeTelemetryClien import * as ResourceAttribution from "./resourceTelemetry/ResourceAttribution.ts"; import * as ResourceMonitorBinary from "./resourceTelemetry/ResourceMonitorBinary.ts"; import * as ResourceTelemetry from "./resourceTelemetry/ResourceTelemetry.ts"; +import * as CursorUsageReader from "./usage/cursorUsageReader.ts"; import * as UsageService from "./usage/UsageService.ts"; import * as RuntimeLayer from "./orchestration-v2/runtimeLayer.ts"; import * as ProjectStore from "./orchestration-v2/ProjectStore.ts"; @@ -225,7 +226,10 @@ const layerBackground = BackgroundPolicy.layer.pipe( Layer.provideMerge(layerServerSettings), ); -const layerUsage = UsageService.layer.pipe(Layer.provide(layerServerSettings)); +const layerUsage = UsageService.layer.pipe( + Layer.provide(layerServerSettings), + Layer.provide(CursorUsageReader.layer), +); const layerResourceDiagnostics = Layer.mergeAll( HostResources.layer, diff --git a/apps/server/src/usage/UsageService.test.ts b/apps/server/src/usage/UsageService.test.ts index 037eb5cc8ef4..ab5607df41f0 100644 --- a/apps/server/src/usage/UsageService.test.ts +++ b/apps/server/src/usage/UsageService.test.ts @@ -15,6 +15,7 @@ import { ProviderDriverKind, ProviderInstanceId, UsageDay, + type UsageSource, type UsageSummaryInput, } from "@t3tools/contracts"; import * as Duration from "effect/Duration"; @@ -31,6 +32,8 @@ import { HttpClient, HttpClientResponse } from "effect/http"; import * as ServerConfig from "../config.ts"; import * as ServerSettings from "../serverSettings.ts"; +import * as CursorUsageReader from "./cursorUsageReader.ts"; +import type { UsageRecord } from "./usageTranscripts.ts"; import * as UsageService from "./UsageService.ts"; const encodeUnknownJson = Schema.encodeEffect(Schema.fromJsonString(Schema.Unknown)); @@ -90,6 +93,7 @@ const layerService = (input: { }) => ServerConfig.layerTest(process.cwd(), { prefix: input.prefix }).pipe( Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(CursorUsageReader.layer), Layer.provideMerge(Layer.succeed(HostProcessPlatform, input.platform ?? "linux")), Layer.provideMerge(ServerSettings.layerTest(input.settings)), Layer.provideMerge( @@ -157,6 +161,71 @@ function totalOutputTokens(summary: { buckets: readonly { totals: { outputTokens return summary.buckets.reduce((sum, bucket) => sum + bucket.totals.outputTokens, 0); } +/** Inside `WINDOW`: the time Cursor tests run at. */ +const CURSOR_NOW = Date.parse("2026-08-02T12:00:00Z"); +const HOUR_MS = 60 * 60 * 1000; +/** The service's cache retention, which the Cursor account cache always covers. */ +const CURSOR_RETENTION_MS = 92 * 24 * HOUR_MS; + +/** + * Stands in for Cursor's dashboard API. Each read returns the account's events + * in its range, keyed like the real reader: occurrences count within one read. + * While `gate` is set, reads wait for it. + */ +function makeFakeCursor() { + const state = { + accountKey: "account-a", + events: [] as { readonly timestampMs: number; readonly outputTokens: number }[], + error: null as string | null, + gate: undefined as Deferred.Deferred | undefined, + calls: [] as { readonly sinceMs: number; readonly untilMs: number }[], + }; + const read = (_credential: unknown, sinceMs: number, untilMs: number) => + Effect.gen(function* () { + state.calls.push({ sinceMs, untilMs }); + if (state.gate !== undefined) yield* Deferred.await(state.gate); + const { accountKey, error } = state; + if (error !== null) return { accountKey, records: [], missing: false, error }; + const occurrences = new Map(); + const records = state.events + .filter(({ timestampMs }) => timestampMs >= sinceMs && timestampMs <= untilMs) + .map(({ timestampMs, outputTokens }): UsageRecord => { + const key = `${timestampMs}:${outputTokens}`; + const occurrence = occurrences.get(key) ?? 0; + occurrences.set(key, occurrence + 1); + return { + provider: "cursor", + timestampMs, + model: "claude-fable-5", + sessionId: "conversation-1", + totals: { + uncachedInputTokens: 0, + cachedInputTokens: 0, + cacheCreationTokens: 0, + outputTokens, + reasoningTokens: 0, + }, + reportedCostUsd: null, + speed: "standard", + dedupeKey: `cursor-account:${accountKey}:${key}:${occurrence}`, + }; + }); + return { accountKey, records, missing: false, error: null }; + }); + return { state, read }; +} + +const writeCursorLogin = (home: string) => + Effect.promise(async () => { + const authPath = NodePath.join(home, "config", "cursor", "auth.json"); + await NodeFSP.mkdir(NodePath.dirname(authPath), { recursive: true }); + await NodeFSP.writeFile(authPath, "{}"); + }); + +function cursorSource(summary: { readonly sources: readonly UsageSource[] }) { + return summary.sources.find((source) => source.fingerprint.provider === "cursor"); +} + describe("UsageService", () => { it.live.each([ { explicitDefault: true, label: "explicit" }, @@ -196,6 +265,9 @@ describe("UsageService", () => { const service = yield* UsageService.make; return yield* service.readSummary(WINDOW); }).pipe( + // Scoped inside the state directory, so pending cache writes land + // before it is removed. + Effect.scoped, Effect.provide( layerService({ prefix: "usage-managed-accounts", @@ -282,12 +354,201 @@ describe("UsageService", () => { layerService({ prefix: "usage-service-cursor-invalid-login", home, settings }), ), ); - const summary = yield* service.readSummary(WINDOW); + const summary = yield* service.readSummary({ ...WINDOW, awaitRefresh: true }); const cursor = summary.sources.find((source) => source.fingerprint.provider === "cursor"); assert.strictEqual(cursor?.message, "Cursor credentials could not be read."); }).pipe(Effect.scoped), ); + it.live("answers Cursor from its cache while one shared refresh runs", () => + Effect.gen(function* () { + const { settings, home } = yield* setup; + yield* writeCursorLogin(home); + yield* TestClock.setTime(CURSOR_NOW); + const cursor = makeFakeCursor(); + cursor.state.events = [{ timestampMs: CURSOR_NOW - 2 * HOUR_MS, outputTokens: 5 }]; + const gate = yield* Deferred.make(); + cursor.state.gate = gate; + yield* Effect.gen(function* () { + const service = yield* UsageService.make.pipe( + Effect.provideService(CursorUsageReader.CursorAccountReader, { read: cursor.read }), + ); + // Cold: nothing cached yet, so Cursor answers empty while it refreshes. + const cold = yield* service.readSummary(WINDOW); + assert.strictEqual(cursorSource(cold)?.refreshing, true); + assert.strictEqual(cursorSource(cold)?.status, "ok"); + assert.strictEqual(totalOutputTokens(cold), 0); + // A second reader joins the refresh in flight. + const joined = yield* service.readSummary(WINDOW); + assert.strictEqual(cursorSource(joined)?.refreshing, true); + + const waited = yield* service + .readSummary({ ...WINDOW, awaitRefresh: true }) + .pipe(Effect.forkChild); + yield* Deferred.succeed(gate, undefined); + const refreshed = yield* Fiber.join(waited); + assert.isFalse(refreshed.sources.some((source) => source.refreshing)); + assert.strictEqual(cursorSource(refreshed)?.status, "ok"); + assert.strictEqual(totalOutputTokens(refreshed), 5); + assert.strictEqual(cursor.state.calls.length, 1); + + // Inside the TTL a plain read answers from the cache without refreshing. + yield* TestClock.adjust(Duration.seconds(30)); + const cached = yield* service.readSummary(WINDOW); + assert.isUndefined(cursorSource(cached)?.refreshing); + assert.strictEqual(totalOutputTokens(cached), 5); + assert.strictEqual(cursor.state.calls.length, 1); + + // Past it, a plain read answers from the stale cache and refreshes. + yield* TestClock.adjust(Duration.seconds(31)); + cursor.state.events.push({ timestampMs: CURSOR_NOW, outputTokens: 7 }); + const stale = yield* service.readSummary(WINDOW); + assert.strictEqual(cursorSource(stale)?.refreshing, true); + assert.strictEqual(totalOutputTokens(stale), 5); + const fresh = yield* service.readSummary({ ...WINDOW, awaitRefresh: true }); + assert.isUndefined(cursorSource(fresh)?.refreshing); + assert.strictEqual(totalOutputTokens(fresh), 12); + assert.strictEqual(cursor.state.calls.length, 2); + }).pipe(Effect.provide(layerService({ prefix: "usage-service-cursor-swr", home, settings }))); + }).pipe(Effect.scoped, Effect.provide(TestClock.layer())), + ); + + it.live("caches Cursor's whole retention and then refetches only the newest edge", () => + Effect.gen(function* () { + const { settings, home } = yield* setup; + yield* writeCursorLogin(home); + yield* TestClock.setTime(CURSOR_NOW); + const cursor = makeFakeCursor(); + // Two identical billed rows inside the refetched hour must both survive it, once each. + cursor.state.events = [ + { timestampMs: Date.parse("2026-07-10T10:00:00Z"), outputTokens: 17 }, + { timestampMs: CURSOR_NOW - 24 * HOUR_MS, outputTokens: 5 }, + { timestampMs: CURSOR_NOW - HOUR_MS / 2, outputTokens: 7 }, + { timestampMs: CURSOR_NOW - HOUR_MS / 2, outputTokens: 7 }, + ]; + yield* Effect.gen(function* () { + const service = yield* UsageService.make.pipe( + Effect.provideService(CursorUsageReader.CursorAccountReader, { read: cursor.read }), + ); + const read = (input: UsageSummaryInput) => + service.readSummary({ ...input, awaitRefresh: true }); + assert.strictEqual(totalOutputTokens(yield* read(WINDOW)), 19); + assert.deepStrictEqual(cursor.state.calls, [ + { sinceMs: CURSOR_NOW - CURSOR_RETENTION_MS, untilMs: CURSOR_NOW }, + ]); + + // A wider window is answered from the same cache. + const wide = { ...WINDOW, sinceDay: UsageDay.make("2026-07-01") }; + assert.strictEqual(totalOutputTokens(yield* read(wide)), 36); + assert.strictEqual(cursor.state.calls.length, 1); + + // An event finalized late, inside the overlap, and a new one. + cursor.state.events.push( + { timestampMs: CURSOR_NOW - HOUR_MS / 6, outputTokens: 11 }, + { timestampMs: CURSOR_NOW + HOUR_MS, outputTokens: 13 }, + ); + yield* TestClock.adjust(Duration.hours(2)); + assert.strictEqual(totalOutputTokens(yield* read(WINDOW)), 43); + assert.deepStrictEqual(cursor.state.calls[1], { + sinceMs: CURSOR_NOW - HOUR_MS, + untilMs: CURSOR_NOW + 2 * HOUR_MS, + }); + + // Another login replaces the cached account's history instead of adding to it. + cursor.state.accountKey = "account-b"; + yield* TestClock.adjust(Duration.minutes(2)); + assert.strictEqual(totalOutputTokens(yield* read(wide)), 60); + assert.strictEqual(cursor.state.calls.length, 4); + assert.strictEqual(cursorSource(yield* read(wide))?.fingerprint.volumeId, "account-b"); + }).pipe( + Effect.provide( + layerService({ prefix: "usage-service-cursor-incremental", home, settings }), + ), + ); + }).pipe(Effect.scoped, Effect.provide(TestClock.layer())), + ); + + it.live("reports a failed Cursor refresh from its cache without refreshing", () => + Effect.gen(function* () { + const { settings, home } = yield* setup; + yield* writeCursorLogin(home); + yield* TestClock.setTime(CURSOR_NOW); + const cursor = makeFakeCursor(); + cursor.state.events = [{ timestampMs: CURSOR_NOW - HOUR_MS * 3, outputTokens: 5 }]; + yield* Effect.gen(function* () { + const service = yield* UsageService.make.pipe( + Effect.provideService(CursorUsageReader.CursorAccountReader, { read: cursor.read }), + ); + yield* service.readSummary({ ...WINDOW, awaitRefresh: true }); + + yield* TestClock.adjust(Duration.minutes(2)); + cursor.state.error = "Sign in to Cursor again to read account usage."; + const failed = yield* service.readSummary({ ...WINDOW, awaitRefresh: true }); + assert.strictEqual(cursorSource(failed)?.status, "partial"); + assert.strictEqual(cursorSource(failed)?.message, cursor.state.error); + assert.isUndefined(cursorSource(failed)?.refreshing); + assert.strictEqual(totalOutputTokens(failed), 5); + + // The failure stands for the TTL; a plain read does not retry it. + const again = yield* service.readSummary(WINDOW); + assert.strictEqual(cursorSource(again)?.status, "partial"); + assert.isUndefined(cursorSource(again)?.refreshing); + assert.strictEqual(cursor.state.calls.length, 2); + + // A failing read of another login drops the previous account's history. + yield* TestClock.adjust(Duration.minutes(2)); + cursor.state.accountKey = "account-b"; + const switched = yield* service.readSummary({ ...WINDOW, awaitRefresh: true }); + assert.strictEqual(cursorSource(switched)?.status, "missing"); + assert.strictEqual(cursorSource(switched)?.message, cursor.state.error); + assert.strictEqual(totalOutputTokens(switched), 0); + }).pipe( + Effect.provide(layerService({ prefix: "usage-service-cursor-failure", home, settings })), + ); + }).pipe(Effect.scoped, Effect.provide(TestClock.layer())), + ); + + it.live("restores the Cursor account cache after a restart", () => + Effect.gen(function* () { + const { settings, home } = yield* setup; + yield* writeCursorLogin(home); + yield* TestClock.setTime(CURSOR_NOW); + const before = makeFakeCursor(); + before.state.events = [ + { timestampMs: CURSOR_NOW - 30 * HOUR_MS, outputTokens: 5 }, + { timestampMs: CURSOR_NOW - 3 * HOUR_MS, outputTokens: 7 }, + ]; + yield* Effect.gen(function* () { + const first = yield* UsageService.make.pipe( + Effect.provideService(CursorUsageReader.CursorAccountReader, { read: before.read }), + ); + const original = yield* first.readSummary({ ...WINDOW, awaitRefresh: true }); + yield* first.awaitPersisted; + + const after = makeFakeCursor(); + after.state.events = before.state.events; + const restarted = yield* UsageService.make.pipe( + Effect.provideService(CursorUsageReader.CursorAccountReader, { read: after.read }), + ); + const restored = yield* restarted.readSummary(WINDOW); + assert.isUndefined(cursorSource(restored)?.refreshing); + assert.deepStrictEqual(restored.buckets, original.buckets); + assert.deepStrictEqual(cursorSource(restored), cursorSource(original)); + assert.strictEqual(after.state.calls.length, 0); + + // Once stale, the restored cache refreshes only its newest edge. + yield* TestClock.adjust(Duration.minutes(2)); + const refreshed = yield* restarted.readSummary({ ...WINDOW, awaitRefresh: true }); + assert.deepStrictEqual(refreshed.buckets, original.buckets); + assert.deepStrictEqual(after.state.calls, [ + { sinceMs: CURSOR_NOW - HOUR_MS, untilMs: CURSOR_NOW + 2 * 60 * 1000 }, + ]); + }).pipe( + Effect.provide(layerService({ prefix: "usage-service-cursor-restart", home, settings })), + ); + }).pipe(Effect.scoped, Effect.provide(TestClock.layer())), + ); + it.live("does not read the macOS Cursor Keychain before account usage is enabled", () => Effect.gen(function* () { const { settings, home } = yield* setup; @@ -642,6 +903,7 @@ describe("UsageService", () => { environmentProjects, ); }).pipe( + Effect.scoped, Effect.provide( layerService({ prefix: "usage-service-home-refresh-test", @@ -747,6 +1009,7 @@ describe("UsageService", () => { const restored = yield* service.readSummary(WINDOW); assert.deepStrictEqual(restored.buckets, original.buckets); }).pipe( + Effect.scoped, Effect.provide( layerService({ prefix: "usage-service-price-overrides-test", home, settings }), ), @@ -801,9 +1064,11 @@ describe("UsageService", () => { appended.buckets.reduce((sum, bucket) => sum + bucket.totals.uncachedInputTokens, 0), 20, ); + yield* service.awaitPersisted; const restarted = yield* UsageService.make; const restored = yield* restarted.readSummary(WINDOW); assert.deepStrictEqual(restored.buckets, appended.buckets); + yield* restarted.awaitPersisted; yield* Effect.promise(() => NodeFSP.rm(transcript)); const afterCleanup = yield* UsageService.make; assert.deepStrictEqual( @@ -811,6 +1076,7 @@ describe("UsageService", () => { appended.buckets, ); }).pipe( + Effect.scoped, Effect.provide( layerService({ prefix: "usage-service-large-record-test", @@ -865,7 +1131,9 @@ describe("UsageService", () => { const { stateDir } = yield* ServerConfig.ServerConfig; const cachePath = NodePath.join(stateDir, "usage-scan-cache-v5.json"); const legacyPath = NodePath.join(stateDir, "usage-scan-cache.json"); - yield* (yield* UsageService.make).readSummary(WINDOW); + const first = yield* UsageService.make; + yield* first.readSummary(WINDOW); + yield* first.awaitPersisted; // Rewrite the cache as a v4 server left it: every Codex record at // speed 0 (standard), and no tier in the reducer state. @@ -937,6 +1205,7 @@ describe("UsageService", () => { assert.deepStrictEqual(deleted.buckets, first.buckets); assert.deepStrictEqual(deleted.sources, first.sources); + yield* service.awaitPersisted; const restarted = yield* UsageService.make; const restored = yield* restarted.readSummary(WINDOW); assert.deepStrictEqual(restored.buckets, first.buckets); @@ -947,6 +1216,7 @@ describe("UsageService", () => { const moved = yield* restarted.readSummary(WINDOW); assert.deepStrictEqual(moved.buckets, first.buckets); assert.strictEqual(moved.sources[0]?.distinctSessions, 1); + yield* restarted.awaitPersisted; const replacementProjects = NodePath.join(home, "replacement-projects"); yield* Effect.promise(() => NodeFSP.mkdir(replacementProjects)); @@ -994,6 +1264,9 @@ describe("UsageService", () => { assert.deepStrictEqual(outsideWindow.buckets, []); assert.strictEqual(outsideWindow.sources[0]?.distinctSessions, 0); }).pipe( + // Scoped inside the state directory, so pending cache writes land + // before it is removed. + Effect.scoped, Effect.provide( layerService({ prefix: "usage-service-cleanup-test", @@ -1039,6 +1312,7 @@ describe("UsageService", () => { const saved = yield* service.readSummary(WINDOW); assert.deepStrictEqual(saved.buckets, live.buckets); assert.deepStrictEqual(saved.sources, live.sources); + yield* service.awaitPersisted; const restored = yield* (yield* UsageService.make).readSummary(WINDOW); assert.deepStrictEqual(restored.buckets, live.buckets); }).pipe( diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index 059d4e176bca..c06c02515898 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -12,6 +12,10 @@ * remain visible. Antigravity databases are memoised in memory while the * database and its WAL keep the same `(size, mtime, ctime)`. * + * Cursor's account API is slow, so its source answers from a cache and marks + * itself `refreshing` while a background refresh runs; `awaitRefresh` waits + * for that refresh instead. See `cursorAccountCache`. + * * @module UsageService */ import * as NodeOS from "node:os"; @@ -38,6 +42,8 @@ import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -55,7 +61,18 @@ import { resolveAntigravityInstanceDirectories } from "../provider/antigravityAu import { mergeProviderInstanceEnvironment } from "../provider/ProviderInstanceEnvironment.ts"; import { readOpenCodeUsage } from "./opencodeUsageReader.ts"; import { makeAntigravityUsageCache, readAntigravityUsage } from "./antigravityUsageReader.ts"; -import { readCursorAccountUsage } from "./cursorUsageReader.ts"; +import { + CURSOR_ACCOUNT_CACHE_FILE_NAME, + CURSOR_ACCOUNT_TTL_MS, + cursorFetchRange, + decodeCursorAccountCaches, + encodeCursorAccountCaches, + isCursorCacheFresh, + mergeCursorFetch, + type CursorAccountCache, + type CursorCredentialSource, +} from "./cursorAccountCache.ts"; +import * as CursorUsageReader from "./cursorUsageReader.ts"; import { resolveModelAliases, UsageAggregator } from "./usageAggregation.ts"; import { createOverrideRateTable, parseRateTable, type RateTable } from "./usagePricing.ts"; import { @@ -91,8 +108,13 @@ const RATES_REFRESH_FLOOR_MS = 60 * 1000; const MTIME_SLACK_MS = 36 * 60 * 60 * 1000; const MAX_HOURLY_WINDOW_MS = 24 * 60 * 60 * 1000; -/** Longest window the UI offers, plus slack. Older entries are pruned. */ -const CACHE_RETENTION_DAYS = 90; +/** + * The longest window the UI offers, 90 days, plus its `MTIME_SLACK_MS`, rounded + * up. Older entries are pruned. + */ +const CACHE_RETENTION_DAYS = 92; + +const CURSOR_ACCOUNT_READ_ERROR = "Cursor account usage could not be read."; /** Transcripts parsed at once. More gains little once the disk stays busy. */ const TRANSCRIPT_READ_CONCURRENCY = 4; @@ -195,11 +217,26 @@ export const make = Effect.gen(function* () { const httpClient = yield* HttpClient.HttpClient; const hostEnvironment = yield* HostProcessEnvironment; const platform = yield* HostProcessPlatform; + const cursorAccountReader = yield* CursorUsageReader.CursorAccountReader; const fileCache: ScanCache = new Map(); const antigravityCache = makeAntigravityUsageCache(); const sourceCache = new Map(); let cacheDirty = false; + /** Cursor account caches by credential source. */ + const cursorCaches = new Map(); + let cursorCacheDirty = false; + /** + * The last failed refresh per credential source, standing for a TTL so a + * client refetching a broken login does not refetch Cursor each time. A + * `null` error means there is no login: no source to report. + */ + const cursorFailures = new Map< + string, + { readonly atMs: number; readonly error: string | null } + >(); + /** The refresh in flight per credential source, which every read joins. */ + const cursorRefreshes = new Map>(); const isWithinDirectory = (filePath: string, dir: string) => { const relative = path.relative(dir, filePath); return relative !== ".." && !relative.startsWith(".." + path.sep) && !path.isAbsolute(relative); @@ -208,6 +245,7 @@ export const make = Effect.gen(function* () { const ratesCachePath = path.join(config.stateDir, "usage-model-rates.json"); const scanCachePath = path.join(config.stateDir, SCAN_CACHE_FILE_NAME); const legacyScanCachePath = path.join(config.stateDir, LEGACY_SCAN_CACHE_FILE_NAME); + const cursorCachePath = path.join(config.stateDir, CURSOR_ACCOUNT_CACHE_FILE_NAME); const writeCacheFile = (filePath: string, contents: string) => writeFileStringAtomically({ filePath, contents }).pipe( Effect.provideService(FileSystem.FileSystem, fileSystem), @@ -398,7 +436,7 @@ export const make = Effect.gen(function* () { }); /** - * Loads the persisted scan cache exactly once per process. + * Loads the persisted scan and Cursor account caches exactly once per process. * * `Effect.cached` makes concurrent first readers await the same load rather * than each seeing a "loaded" flag set before the read finished and cold @@ -411,6 +449,9 @@ export const make = Effect.gen(function* () { Effect.flatMap((raw) => decodeScanCacheFile(raw)), Effect.catchCause(() => Effect.succeed(null)), ); + for (const [key, cache] of decodeCursorAccountCaches(yield* readDocument(cursorCachePath))) { + cursorCaches.set(key, cache); + } let document = yield* readDocument(scanCachePath); if (document === null) { document = yield* readDocument(legacyScanCachePath); @@ -432,25 +473,52 @@ export const make = Effect.gen(function* () { // keeps an older snapshot from landing after a newer one. const persistLock = yield* Semaphore.make(1); - const persistScanCache = Effect.fn("UsageService.persistScanCache")(function* () { - if (!cacheDirty) return; - // Cleared before encoding, so a scan that changes the cache while this - // write is in flight marks it dirty again. A failed write restores the - // flag, so the next scan retries instead of leaving disk stale. - cacheDirty = false; - yield* Effect.sync(() => - writeScanCache(fileCache, { sources: Object.fromEntries(sourceCache) }), - ).pipe( - Effect.flatMap((contents) => writeCacheFile(scanCachePath, contents)), - // A cache we cannot write is a slower next start, not a failed read. - Effect.catchCause(() => - Effect.sync(() => { - cacheDirty = true; - }), - ), - persistLock.withPermit, - ); - }); + // Each dirty flag is cleared before encoding, so a change while the write is + // in flight marks it dirty again. A failed write restores the flag, so the + // next persist retries instead of leaving disk stale. A cache we cannot + // write is a slower next start, not a failed read. + const persistCaches = Effect.gen(function* () { + if (cacheDirty) { + cacheDirty = false; + yield* Effect.sync(() => + writeScanCache(fileCache, { sources: Object.fromEntries(sourceCache) }), + ).pipe( + Effect.flatMap((contents) => writeCacheFile(scanCachePath, contents)), + Effect.catchCause(() => + Effect.sync(() => { + cacheDirty = true; + }), + ), + ); + } + if (cursorCacheDirty) { + cursorCacheDirty = false; + yield* Effect.sync(() => encodeCursorAccountCaches(cursorCaches)).pipe( + Effect.flatMap((contents) => writeCacheFile(cursorCachePath, contents)), + Effect.catchCause(() => + Effect.sync(() => { + cursorCacheDirty = true; + }), + ), + ); + } + }).pipe(persistLock.withPermit, Effect.withSpan("UsageService.persistCaches")); + + const pendingPersists = new Set>(); + /** Writes dirty caches in the background, after the summary that dirtied them answers. */ + const schedulePersist = Effect.forkDetach(persistCaches).pipe( + Effect.map((fiber) => { + pendingPersists.add(fiber); + fiber.addObserver(() => pendingPersists.delete(fiber)); + }), + ); + /** Waits for every write scheduled so far, as a restart would need. */ + const awaitPersisted = Effect.suspend(() => Fiber.awaitAll([...pendingPersists])).pipe( + Effect.asVoid, + ); + // A write still running when the service shuts down finishes first, so the + // next start does not lose the last scan. + yield* Effect.addFinalizer(() => awaitPersisted); /** * Parses one transcript, reusing the cached result when it is unchanged. @@ -541,58 +609,223 @@ export const make = Effect.gen(function* () { readonly files: | readonly { readonly path: string; readonly records: readonly UsageRecord[] }[] | null; + /** Answered from a cache while a refresh runs. */ + readonly refreshing?: true; } + const scanTranscriptDir = Effect.fn("UsageService.scanTranscriptDir")(function* ( + source: { + readonly provider: UsageProviderKind; + readonly dir: string; + readonly volumeId: string; + readonly fileName?: string; + }, + windowStartMs: number, + ) { + const { provider, dir, volumeId, fileName } = source; + const exists = yield* fileSystem + .exists(dir) + .pipe(Effect.catchCause(() => Effect.succeed(false))); + if (!exists) return { provider, dir, volumeId, files: null } satisfies ScannedDir; + const files = yield* Effect.promise(() => + listTranscriptFiles(dir, windowStartMs, fileName === undefined ? undefined : { fileName }), + ); + // A cold parse waits on disk reads, so a few files in flight read + // close to twice as fast. Results keep walk order. + const read = yield* Effect.forEach( + files, + (file) => + readFileRecords(file.path, file.size, file.mtimeMs, provider).pipe( + Effect.map((result) => ({ path: file.path, ...result })), + ), + { concurrency: TRANSCRIPT_READ_CONCURRENCY }, + ); + const parsedFiles = read.map(({ path, records, update }) => { + if (update === undefined) return { path, records }; + // A scan of another window may have cached its own read of this file + // meanwhile. Then keep whichever read saw the later file, so a slower + // scan never replaces newer usage with older. + const current = fileCache.get(path); + if ( + current === update.replaces || + current === undefined || + !isLaterRead(current, update.entry) + ) { + fileCache.set(path, update.entry); + cacheDirty = true; + } + return { path, records }; + }); + return { provider, dir, volumeId, files: parsedFiles } satisfies ScannedDir; + }); + + /** Fetches what one account cache is missing, then persists it if anything changed. */ + const refreshCursorAccount = Effect.fn("UsageService.refreshCursorAccount")(function* ( + credential: CursorCredentialSource, + credentialKey: string, + retentionStartMs: number, + ) { + const nowMs = yield* Clock.currentTimeMillis; + const fetchMissing = (cache: CursorAccountCache | undefined) => { + const range = cursorFetchRange(cache, retentionStartMs, nowMs); + return cursorAccountReader + .read(credential, range.sinceMs, range.untilMs) + .pipe(Effect.map((result) => ({ range, result }))); + }; + let base = cursorCaches.get(credentialKey); + let fetched = yield* fetchMissing(base); + // Another login's history replaces the cached account's, even if reading it fails. + if ( + base !== undefined && + fetched.result.accountKey !== null && + fetched.result.accountKey !== base.accountKey + ) { + cursorCaches.delete(credentialKey); + cursorCacheDirty = true; + base = undefined; + fetched = yield* fetchMissing(base); + } + const { range, result } = fetched; + if (result.missing || result.error !== null || result.accountKey === null) { + cursorFailures.set(credentialKey, { + atMs: nowMs, + error: result.missing || result.error !== null ? result.error : CURSOR_ACCOUNT_READ_ERROR, + }); + yield* schedulePersist; + return; + } + const merged = mergeCursorFetch( + base, + result.accountKey, + range, + result.records, + nowMs, + retentionStartMs, + ); + cursorCaches.set(credentialKey, merged.cache); + cursorFailures.delete(credentialKey); + // An unchanged edge is not worth rewriting the file for: after a restart + // the cache just refetches a slightly wider edge. + if (merged.changed) { + cursorCacheDirty = true; + yield* schedulePersist; + } + }); + + /** Joins the refresh in flight, or starts one. */ + const startCursorRefresh = ( + credential: CursorCredentialSource, + credentialKey: string, + retentionStartMs: number, + ) => + // Enrollment and fork are atomic, so an interrupted caller cannot leave a + // registered refresh that nothing will finish. + Effect.uninterruptible( + Effect.gen(function* () { + const current = cursorRefreshes.get(credentialKey); + if (current !== undefined) return current; + const done = Deferred.makeUnsafe(); + cursorRefreshes.set(credentialKey, done); + // Detached: a departing client must not cancel a fetch later reads reuse. + yield* refreshCursorAccount(credential, credentialKey, retentionStartMs).pipe( + Effect.catchCause(() => + Clock.currentTimeMillis.pipe( + Effect.map((atMs) => + cursorFailures.set(credentialKey, { atMs, error: CURSOR_ACCOUNT_READ_ERROR }), + ), + ), + ), + Effect.ensuring( + Effect.suspend(() => { + cursorRefreshes.delete(credentialKey); + return Deferred.succeed(done, undefined); + }), + ), + Effect.forkDetach, + ); + return done; + }), + ); + + /** + * The Cursor account source. A cache inside its TTL, or a refresh that failed + * inside it, answers directly. Otherwise a refresh starts: `awaitRefresh` + * waits for it, and anything else answers from the cache marked `refreshing`. + */ + const cursorAccountSource = Effect.fn("UsageService.cursorAccountSource")(function* ( + credential: CursorCredentialSource, + authPath: string, + windowStartMs: number, + retentionStartMs: number, + awaitRefresh: boolean, + ) { + // No saved login means there is no account source to report, not a setup error. + if ( + typeof credential === "string" && + !(yield* fileSystem.exists(credential).pipe(Effect.orElseSucceed(() => true))) + ) { + return null; + } + const credentialKey = typeof credential === "string" ? credential : "keychain"; + const nowMs = yield* Clock.currentTimeMillis; + const recentFailure = cursorFailures.get(credentialKey); + let refreshing = false; + if ( + (recentFailure === undefined || nowMs - recentFailure.atMs >= CURSOR_ACCOUNT_TTL_MS) && + !isCursorCacheFresh(cursorCaches.get(credentialKey), nowMs) + ) { + const refresh = yield* startCursorRefresh(credential, credentialKey, retentionStartMs); + if (awaitRefresh) yield* Deferred.await(refresh); + else refreshing = true; + } + + const cache = cursorCaches.get(credentialKey); + const failure = refreshing ? undefined : cursorFailures.get(credentialKey); + const failureMessage = failure === undefined ? undefined : failure.error; + if (failureMessage === null) return null; + if (cache === undefined) { + return { + provider: "cursor", + dir: authPath, + volumeId: yield* Effect.promise(() => readDirectoryVolumeId(authPath)), + // Never combine a local fallback with another server's account-wide history. + ...(refreshing + ? { files: [], refreshing: true } + : { files: null, message: failureMessage ?? CURSOR_ACCOUNT_READ_ERROR }), + } satisfies ScannedDir; + } + // The same account includes CLI and desktop history from every machine. + // A stable remote fingerprint prevents connected environments counting it twice. + const source = `cursor-account:${cache.accountKey}`; + return { + provider: "cursor", + dir: source, + hostId: "cursor.com", + volumeId: cache.accountKey, + files: [ + { + path: source, + records: cache.records.filter((record) => record.timestampMs >= windowStartMs), + }, + ], + ...(failureMessage === undefined + ? { status: "ok" } + : { status: "partial", message: failureMessage }), + ...(refreshing ? { refreshing: true } : {}), + } satisfies ScannedDir; + }); + const collectDirs = Effect.fn("UsageService.collectDirs")(function* ( windowStartMs: number, settings: ServerSettingsValue, retentionCutoffMs: number, + awaitRefresh: boolean, ) { // The home resolvers ask for `Path` themselves; satisfy them from the // instance we already hold so the scan stays context-free. const dirs = yield* resolveTranscriptDirs(settings, retentionCutoffMs).pipe( Effect.provideService(Path.Path, path), ); - const scanned: ScannedDir[] = []; - for (const { provider, dir, volumeId, fileName } of dirs) { - const exists = yield* fileSystem - .exists(dir) - .pipe(Effect.catchCause(() => Effect.succeed(false))); - if (!exists) { - scanned.push({ provider, dir, volumeId, files: null }); - continue; - } - const files = yield* Effect.promise(() => - listTranscriptFiles(dir, windowStartMs, fileName === undefined ? undefined : { fileName }), - ); - // A cold parse waits on disk reads, so a few files in flight read - // close to twice as fast. Results keep walk order. - const read = yield* Effect.forEach( - files, - (file) => - readFileRecords(file.path, file.size, file.mtimeMs, provider).pipe( - Effect.map((result) => ({ path: file.path, ...result })), - ), - { concurrency: TRANSCRIPT_READ_CONCURRENCY }, - ); - const parsedFiles = read.map(({ path, records, update }) => { - if (update === undefined) return { path, records }; - // A scan of another window may have cached its own read of this file - // meanwhile. Then keep whichever read saw the later file, so a slower - // scan never replaces newer usage with older. - const current = fileCache.get(path); - if ( - current === update.replaces || - current === undefined || - !isLaterRead(current, update.entry) - ) { - fileCache.set(path, update.entry); - cacheDirty = true; - } - return { path, records }; - }); - scanned.push({ provider, dir, volumeId, files: parsedFiles }); - } const home = NodeOS.homedir(); const envRoots = Effect.fnUntraced(function* (key: string, defaults: readonly string[]) { @@ -610,156 +843,168 @@ export const make = Effect.gen(function* () { return [...canonical]; }); const dataHome = hostEnvironment["XDG_DATA_HOME"]?.trim(); - for (const dir of yield* envRoots("OPENCODE_DATA_DIR", [ - path.join( - dataHome && path.isAbsolute(dataHome) ? dataHome : path.join(home, ".local", "share"), - "opencode", - ), - ])) { - const result = yield* Effect.promise(() => readOpenCodeUsage(dir, windowStartMs)); - scanned.push({ - provider: "opencode", - dir, - volumeId: yield* Effect.promise(() => readDirectoryVolumeId(dir)), - files: result.missing && !result.error ? null : result.files, - status: result.error ? "partial" : "ok", - ...(result.error ? { message: "Some OpenCode history could not be read." } : {}), - }); - } - const antigravityRoots = yield* envRoots("ANTIGRAVITY_DATA_DIR", [ - ...["antigravity", "antigravity-cli", "antigravity-ide", "antigravity-backup"].map((name) => - path.join(home, ".gemini", name), - ), - path.join(home, ".config", "antigravity"), - ]); - for (const [instanceId, instance] of Object.entries(settings.providerInstances)) { - if (instance.driver === "antigravity") { - const directories = yield* resolveAntigravityInstanceDirectories( - config.stateDir, - ProviderInstanceId.make(instanceId), - ).pipe( - Effect.provideService(Crypto.Crypto, crypto), - Effect.provideService(Path.Path, path), - Effect.mapError( - (cause) => - new UsageReadError({ - reason: "scanFailed", - detail: "Antigravity profile directory could not be resolved.", - cause, - }), - ), - ); - antigravityRoots.push(path.join(directories.profile, "antigravity-acp")); - } - } - const antigravityDirs = new Set(); - for (const root of antigravityRoots) { - const resolvedRoot = yield* fileSystem.realPath(root).pipe(Effect.orElseSucceed(() => root)); - const nested = path.join(resolvedRoot, "conversations"); - const dir = (yield* fileSystem - .exists(nested) - .pipe(Effect.catchCause(() => Effect.succeed(false)))) - ? nested - : resolvedRoot; - antigravityDirs.add(yield* fileSystem.realPath(dir).pipe(Effect.orElseSucceed(() => dir))); - } - const antigravity = yield* Effect.promise(() => - readAntigravityUsage([...antigravityDirs], windowStartMs, antigravityCache), - ); - for (const dir of antigravityDirs) { - const exists = yield* fileSystem - .exists(dir) - .pipe(Effect.catchCause(() => Effect.succeed(false))); - const failed = antigravity.errors.some( - (error) => error === dir || error.startsWith(`${dir}${path.sep}`), + + const openCode = Effect.gen(function* () { + const roots = yield* envRoots("OPENCODE_DATA_DIR", [ + path.join( + dataHome && path.isAbsolute(dataHome) ? dataHome : path.join(home, ".local", "share"), + "opencode", + ), + ]); + return yield* Effect.forEach( + roots, + (dir) => + Effect.gen(function* () { + const result = yield* Effect.promise(() => readOpenCodeUsage(dir, windowStartMs)); + return { + provider: "opencode", + dir, + volumeId: yield* Effect.promise(() => readDirectoryVolumeId(dir)), + files: result.missing && !result.error ? null : result.files, + status: result.error ? "partial" : "ok", + ...(result.error ? { message: "Some OpenCode history could not be read." } : {}), + } satisfies ScannedDir; + }), + { concurrency: "unbounded" }, ); - scanned.push({ - provider: "antigravity", - dir, - volumeId: yield* Effect.promise(() => readDirectoryVolumeId(dir)), - files: !exists && !failed ? null : antigravity.files.filter((file) => file.root === dir), - status: failed ? "partial" : "ok", - ...(failed ? { message: "Some Antigravity history could not be read." } : {}), - }); - } - const cursorUserHome = - (platform === "win32" ? hostEnvironment["USERPROFILE"] : hostEnvironment["HOME"]) || home; - const configHome = hostEnvironment["XDG_CONFIG_HOME"]?.trim(); - const cursorHome = - platform === "darwin" - ? path.join(cursorUserHome, "Library", "Application Support") - : platform === "win32" - ? hostEnvironment["APPDATA"] || path.join(cursorUserHome, "AppData", "Roaming") - : configHome && path.isAbsolute(configHome) - ? configHome - : path.join(cursorUserHome, ".config"); - const cursorAuthPath = - platform === "darwin" - ? path.join(cursorUserHome, ".cursor", "auth.json") - : path.join(cursorHome, platform === "win32" ? "Cursor" : "cursor", "auth.json"); - const credentialStore = hostEnvironment["AGENT_CLI_CREDENTIAL_STORE"]; - const loginUnavailable = - Boolean(hostEnvironment["CURSOR_AUTH_TOKEN"]?.trim()) || - Boolean(hostEnvironment["CURSOR_API_KEY"]?.trim()) || - credentialStore === "memory"; - if ( - platform === "darwin" && - credentialStore !== "file" && - !loginUnavailable && - !settings.cursorKeychainUsageEnabled - ) { - scanned.push({ - provider: "cursor", - dir: cursorAuthPath, - volumeId: "", - files: null, - message: "Cursor account usage is off on this environment.", - action: "enableCursorKeychain", - }); - return scanned; - } - const cursorUntilMs = yield* Clock.currentTimeMillis; - const account = loginUnavailable - ? { - accountKey: null, - records: [], - missing: true, - error: "Cursor account history needs a Cursor CLI login on this server.", + }); + + const antigravity = Effect.gen(function* () { + const antigravityRoots = yield* envRoots("ANTIGRAVITY_DATA_DIR", [ + ...["antigravity", "antigravity-cli", "antigravity-ide", "antigravity-backup"].map((name) => + path.join(home, ".gemini", name), + ), + path.join(home, ".config", "antigravity"), + ]); + for (const [instanceId, instance] of Object.entries(settings.providerInstances)) { + if (instance.driver === "antigravity") { + const directories = yield* resolveAntigravityInstanceDirectories( + config.stateDir, + ProviderInstanceId.make(instanceId), + ).pipe( + Effect.provideService(Crypto.Crypto, crypto), + Effect.provideService(Path.Path, path), + Effect.mapError( + (cause) => + new UsageReadError({ + reason: "scanFailed", + detail: "Antigravity profile directory could not be resolved.", + cause, + }), + ), + ); + antigravityRoots.push(path.join(directories.profile, "antigravity-acp")); } - : yield* Effect.promise(() => - readCursorAccountUsage( - platform === "darwin" && credentialStore !== "file" - ? { kind: "keychain" } - : cursorAuthPath, - windowStartMs, - cursorUntilMs, - ), + } + const antigravityDirs = new Set(); + for (const root of antigravityRoots) { + const resolvedRoot = yield* fileSystem + .realPath(root) + .pipe(Effect.orElseSucceed(() => root)); + const nested = path.join(resolvedRoot, "conversations"); + const dir = (yield* fileSystem + .exists(nested) + .pipe(Effect.catchCause(() => Effect.succeed(false)))) + ? nested + : resolvedRoot; + antigravityDirs.add(yield* fileSystem.realPath(dir).pipe(Effect.orElseSucceed(() => dir))); + } + const result = yield* Effect.promise(() => + readAntigravityUsage([...antigravityDirs], windowStartMs, antigravityCache), + ); + const scanned: ScannedDir[] = []; + for (const dir of antigravityDirs) { + const exists = yield* fileSystem + .exists(dir) + .pipe(Effect.catchCause(() => Effect.succeed(false))); + const failed = result.errors.some( + (error) => error === dir || error.startsWith(`${dir}${path.sep}`), ); - // No saved login means there is no account source to report, not a setup error. - if (account.missing && account.error === null) return scanned; - if (account.accountKey !== null && account.error === null && !account.missing) { - // The same account includes CLI and desktop history from every machine. - // A stable remote fingerprint prevents connected environments counting it twice. - const source = `cursor-account:${account.accountKey}`; - scanned.push({ - provider: "cursor", - dir: source, - hostId: "cursor.com", - volumeId: account.accountKey, - files: [{ path: source, records: account.records }], - status: "ok", - }); + scanned.push({ + provider: "antigravity", + dir, + volumeId: yield* Effect.promise(() => readDirectoryVolumeId(dir)), + files: !exists && !failed ? null : result.files.filter((file) => file.root === dir), + status: failed ? "partial" : "ok", + ...(failed ? { message: "Some Antigravity history could not be read." } : {}), + }); + } return scanned; - } - scanned.push({ - provider: "cursor", - dir: cursorAuthPath, - volumeId: yield* Effect.promise(() => readDirectoryVolumeId(cursorAuthPath)), - // Never combine a local fallback with another server's account-wide history. - files: null, - message: - account.error ?? "Cursor account history needs a Cursor CLI login saved on this server.", }); + + const cursor = Effect.gen(function* () { + const cursorUserHome = + (platform === "win32" ? hostEnvironment["USERPROFILE"] : hostEnvironment["HOME"]) || home; + const configHome = hostEnvironment["XDG_CONFIG_HOME"]?.trim(); + const cursorHome = + platform === "darwin" + ? path.join(cursorUserHome, "Library", "Application Support") + : platform === "win32" + ? hostEnvironment["APPDATA"] || path.join(cursorUserHome, "AppData", "Roaming") + : configHome && path.isAbsolute(configHome) + ? configHome + : path.join(cursorUserHome, ".config"); + const cursorAuthPath = + platform === "darwin" + ? path.join(cursorUserHome, ".cursor", "auth.json") + : path.join(cursorHome, platform === "win32" ? "Cursor" : "cursor", "auth.json"); + const credentialStore = hostEnvironment["AGENT_CLI_CREDENTIAL_STORE"]; + const loginUnavailable = + Boolean(hostEnvironment["CURSOR_AUTH_TOKEN"]?.trim()) || + Boolean(hostEnvironment["CURSOR_API_KEY"]?.trim()) || + credentialStore === "memory"; + const useKeychain = platform === "darwin" && credentialStore !== "file"; + if (useKeychain && !loginUnavailable && !settings.cursorKeychainUsageEnabled) { + return [ + { + provider: "cursor", + dir: cursorAuthPath, + volumeId: "", + files: null, + message: "Cursor account usage is off on this environment.", + action: "enableCursorKeychain", + } satisfies ScannedDir, + ]; + } + if (loginUnavailable) { + return [ + { + provider: "cursor", + dir: cursorAuthPath, + volumeId: yield* Effect.promise(() => readDirectoryVolumeId(cursorAuthPath)), + files: null, + message: "Cursor account history needs a Cursor CLI login on this server.", + } satisfies ScannedDir, + ]; + } + const source = yield* cursorAccountSource( + useKeychain ? { kind: "keychain" } : cursorAuthPath, + cursorAuthPath, + windowStartMs, + retentionCutoffMs, + awaitRefresh, + ); + return source === null ? [] : [source]; + }); + + // Independent sources scan together. Transcript directories go one at a + // time, so open files stay at `TRANSCRIPT_READ_CONCURRENCY`. The result + // keeps this order, since aggregation keeps the first copy of a duplicate. + const [transcripts, openCodeDirs, antigravityDirs, cursorDirs] = yield* Effect.all( + [ + Effect.forEach(dirs, (dir) => scanTranscriptDir(dir, windowStartMs)), + openCode, + antigravity, + cursor, + ], + { concurrency: "unbounded" }, + ); + const scanned: readonly ScannedDir[] = [ + ...transcripts, + ...openCodeDirs, + ...antigravityDirs, + ...cursorDirs, + ]; return scanned; }); @@ -818,7 +1063,10 @@ export const make = Effect.gen(function* () { // loads while transcripts stream instead of gating them: a cold rates // fetch on a slow network no longer delays the scan by its own timeout. const [, scannedDirs] = yield* Effect.all( - [ensureRates(false), collectDirs(windowStartMs, settings, retentionCutoffMs)], + [ + ensureRates(false), + collectDirs(windowStartMs, settings, retentionCutoffMs, input.awaitRefresh === true), + ], { concurrency: 2 }, ); @@ -859,7 +1107,7 @@ export const make = Effect.gen(function* () { for (const [ index, - { provider, dir, volumeId, files, status, message, action, hostId: sourceHostId }, + { provider, dir, volumeId, files, status, message, action, refreshing, hostId: sourceHostId }, ] of scannedDirs.entries()) { let scannedFiles = 0; let skippedFiles = 0; @@ -910,12 +1158,12 @@ export const make = Effect.gen(function* () { message: message ?? (files === null ? "No transcript directory on this environment." : null), ...(action ? { action } : {}), + ...(refreshing ? { refreshing } : {}), }); } const pruned = pruneScanCache(fileCache, retentionCutoffMs); if (pruned > 0) cacheDirty = true; - yield* persistScanCache(); const aggregated = aggregator.finish(); const readAt = yield* DateTime.now; @@ -952,6 +1200,8 @@ export const make = Effect.gen(function* () { settings.usagePriceOverrides, settings.usageModelAliases, settings.cursorKeychainUsageEnabled, + // A waiting read must never share a scan that answers with `refreshing`. + input.awaitRefresh === true, ]); const readSummary = Effect.fn("UsageService.readSummary")(function* (input: UsageSummaryInput) { @@ -968,9 +1218,12 @@ export const make = Effect.gen(function* () { inflightScans.set(key, created); // Detached so one departing client cannot tear the scan out from under // the fibers awaiting it; a finished scan warms the cache either way. + // The cache write is registered before the waiters resume, so they + // can await it, but its fiber starts after they have the summary. yield* scanSummary(input, settings).pipe( Effect.onExit((exit) => Effect.sync(() => inflightScans.delete(key)).pipe( + Effect.andThen(Exit.isSuccess(exit) ? schedulePersist : Effect.void), Effect.andThen(Deferred.done(created, exit)), ), ), @@ -984,7 +1237,9 @@ export const make = Effect.gen(function* () { return yield* Deferred.await(deferred); }); - return { readSummary, refreshRates } as const; + // `awaitPersisted` is outside the service interface: tests use it to restart + // against what a previous instance wrote. + return { readSummary, refreshRates, awaitPersisted } as const; }); export const layer = Layer.effect(UsageService, make); diff --git a/apps/server/src/usage/cursorAccountCache.test.ts b/apps/server/src/usage/cursorAccountCache.test.ts new file mode 100644 index 000000000000..b6d20d36cf6a --- /dev/null +++ b/apps/server/src/usage/cursorAccountCache.test.ts @@ -0,0 +1,80 @@ +import { assert, describe, it } from "@effect/vitest"; + +import { cursorFetchRange, mergeCursorFetch } from "./cursorAccountCache.ts"; +import type { UsageRecord } from "./usageTranscripts.ts"; + +const DAY_MS = 24 * 60 * 60 * 1000; + +const record = (timestampMs: number): UsageRecord => ({ + provider: "cursor", + timestampMs, + model: "claude-fable-5", + sessionId: "conversation-1", + totals: { + uncachedInputTokens: 0, + cachedInputTokens: 0, + cacheCreationTokens: 0, + outputTokens: 1, + reasoningTokens: 0, + }, + reportedCostUsd: null, + speed: "standard", + dedupeKey: `cursor-account:a:${timestampMs}:0`, +}); + +describe("cursorAccountCache", () => { + it("refetches everything once the cache no longer reaches the retention start", () => { + const cache = { + accountKey: "a", + sinceMs: 0, + untilMs: 10 * DAY_MS, + fetchedAtMs: 10 * DAY_MS, + records: [record(5 * DAY_MS)], + }; + const nowMs = 100 * DAY_MS; + const range = cursorFetchRange(cache, 10 * DAY_MS, nowMs); + assert.deepStrictEqual(range, { sinceMs: 10 * DAY_MS, untilMs: nowMs }); + + const { cache: merged } = mergeCursorFetch( + cache, + "a", + range, + [record(50 * DAY_MS)], + nowMs, + 10 * DAY_MS, + ); + assert.deepStrictEqual( + merged.records.map((entry) => entry.timestampMs), + [50 * DAY_MS], + ); + // Then only the newest edge. + assert.deepStrictEqual(cursorFetchRange(merged, 11 * DAY_MS, nowMs + DAY_MS), { + sinceMs: nowMs - 60 * 60 * 1000, + untilMs: nowMs + DAY_MS, + }); + }); + + it("reports an edge that only confirms the cache as unchanged", () => { + const cache = { + accountKey: "a", + sinceMs: 0, + untilMs: 10 * DAY_MS, + fetchedAtMs: 10 * DAY_MS, + records: [record(DAY_MS), record(10 * DAY_MS)], + }; + const range = cursorFetchRange(cache, 0, 11 * DAY_MS); + assert.isFalse( + mergeCursorFetch(cache, "a", range, [record(10 * DAY_MS)], 11 * DAY_MS, 0).changed, + ); + assert.isTrue( + mergeCursorFetch( + cache, + "a", + range, + [record(10 * DAY_MS), record(10.5 * DAY_MS)], + 11 * DAY_MS, + 0, + ).changed, + ); + }); +}); diff --git a/apps/server/src/usage/cursorAccountCache.ts b/apps/server/src/usage/cursorAccountCache.ts new file mode 100644 index 000000000000..4f3a127ec959 --- /dev/null +++ b/apps/server/src/usage/cursorAccountCache.ts @@ -0,0 +1,271 @@ +/** + * Cursor account usage, cached per credential source. + * + * Cursor's dashboard API takes about 20 seconds for a 30-day window, so fetched + * records are kept for the whole retention period, whatever window was asked + * for. After one full fetch, a refresh fetches only the newest edge. + * + * A record's dedupe key carries its timestamp and an occurrence index counted + * within one fetch. A refetched range therefore replaces every cached record in + * that range rather than adding to it, so each timestamp's records always come + * from a single fetch and their occurrence indexes cannot collide. + * + * @module cursorAccountCache + */ +import { cursorRateModel } from "./cursorUsageReader.ts"; +import type { UsageRecord } from "./usageTranscripts.ts"; + +export type CursorCredentialSource = string | { readonly kind: "keychain" }; + +/** Matches the clients' stale time, so a page refetching on focus reuses one fetch. */ +export const CURSOR_ACCOUNT_TTL_MS = 60 * 1000; + +/** Dashboard events can be finalized after they were first read. */ +const REFETCH_OVERLAP_MS = 60 * 60 * 1000; + +export interface CursorAccountCache { + readonly accountKey: string; + /** Inclusive range of event timestamps the records are complete for. */ + readonly sinceMs: number; + readonly untilMs: number; + /** When the newest edge was fetched. */ + readonly fetchedAtMs: number; + readonly records: readonly UsageRecord[]; +} + +interface CursorFetchRange { + readonly sinceMs: number; + readonly untilMs: number; +} + +/** Whether the cache answers without a refresh. */ +export function isCursorCacheFresh( + cache: CursorAccountCache | undefined, + nowMs: number, +): cache is CursorAccountCache { + return cache !== undefined && nowMs - cache.fetchedAtMs < CURSOR_ACCOUNT_TTL_MS; +} + +/** + * The range to fetch so the cache covers `[retentionStartMs, nowMs]`: the + * newest edge, or everything when the cache does not reach back that far. + */ +export function cursorFetchRange( + cache: CursorAccountCache | undefined, + retentionStartMs: number, + nowMs: number, +): CursorFetchRange { + return cache === undefined || + cache.sinceMs > retentionStartMs || + cache.untilMs - REFETCH_OVERLAP_MS < retentionStartMs + ? { sinceMs: retentionStartMs, untilMs: nowMs } + : { sinceMs: cache.untilMs - REFETCH_OVERLAP_MS, untilMs: nowMs }; +} + +/** + * Applies a fetched range to the cache, replacing the cached records it covers + * and dropping those before the retention start. Pass `undefined` for a cache + * that belongs to another account. `changed` is false when the fetch only + * confirmed the cached records, so there is nothing new to persist. + */ +export function mergeCursorFetch( + cache: CursorAccountCache | undefined, + accountKey: string, + range: CursorFetchRange, + fetched: readonly UsageRecord[], + nowMs: number, + retentionStartMs: number, +) { + const inRange = (record: UsageRecord) => + record.timestampMs >= range.sinceMs && record.timestampMs <= range.untilMs; + const previous = cache?.records ?? []; + const added = fetched.filter(inRange); + const replacedKeys = new Set(previous.filter(inRange).map((record) => record.dedupeKey)); + const changed = + cache === undefined || + replacedKeys.size !== added.length || + added.some((record) => !replacedKeys.has(record.dedupeKey)); + return { + changed, + cache: { + accountKey, + sinceMs: retentionStartMs, + untilMs: range.untilMs, + fetchedAtMs: nowMs, + records: [ + ...previous.filter((record) => !inRange(record) && record.timestampMs >= retentionStartMs), + ...added, + ], + } satisfies CursorAccountCache, + }; +} + +/** + * Each version writes its own file, as the scan cache does, so servers of + * different versions sharing a state directory keep their own caches. + */ +export const CURSOR_ACCOUNT_CACHE_FILE_NAME = "usage-cursor-account-cache-v1.json"; +const CURSOR_ACCOUNT_CACHE_VERSION = 1; + +/** Positional rows with interned strings, as in the scan cache. */ +type SerializedRecord = readonly [ + timestampMs: number, + modelIndex: number, + sessionIndex: number, + uncachedInputTokens: number, + cachedInputTokens: number, + cacheCreationTokens: number, + outputTokens: number, + reportedCostUsd: number | null, + dedupeKey: string | null, +]; + +interface SerializedAccount { + readonly accountKey: string; + readonly sinceMs: number; + readonly untilMs: number; + readonly fetchedAtMs: number; + readonly records: readonly SerializedRecord[]; +} + +/** Serialises the caches by credential source. */ +export function encodeCursorAccountCaches(caches: ReadonlyMap): string { + const strings: string[] = []; + const indexes = new Map(); + const intern = (value: string) => { + let index = indexes.get(value); + if (index === undefined) { + index = strings.length; + strings.push(value); + indexes.set(value, index); + } + return index; + }; + const accounts: Record = {}; + for (const [credentialKey, cache] of caches) { + accounts[credentialKey] = { + accountKey: cache.accountKey, + sinceMs: cache.sinceMs, + untilMs: cache.untilMs, + fetchedAtMs: cache.fetchedAtMs, + records: cache.records.map((record) => [ + record.timestampMs, + intern(record.model), + intern(record.sessionId), + record.totals.uncachedInputTokens, + record.totals.cachedInputTokens, + record.totals.cacheCreationTokens, + record.totals.outputTokens, + record.reportedCostUsd, + record.dedupeKey, + ]), + }; + } + return JSON.stringify({ version: CURSOR_ACCOUNT_CACHE_VERSION, strings, accounts }); +} + +const isFiniteNumber = (value: unknown): value is number => + typeof value === "number" && Number.isFinite(value); + +/** + * Rebuilds the caches from a parsed document. A malformed account is dropped + * whole, costing one refetch; keeping its survivors would leave a gap the + * covered range claims is complete. + */ +export function decodeCursorAccountCaches(document: unknown): Map { + const caches = new Map(); + if (typeof document !== "object" || document === null) return caches; + const root = document as { + readonly version?: unknown; + readonly strings?: unknown; + readonly accounts?: unknown; + }; + const strings = root.strings; + if ( + root.version !== CURSOR_ACCOUNT_CACHE_VERSION || + !Array.isArray(strings) || + !strings.every((value): value is string => typeof value === "string") || + typeof root.accounts !== "object" || + root.accounts === null + ) { + return caches; + } + const stringAt = (index: unknown): string | undefined => + typeof index === "number" ? strings[index] : undefined; + + const decodeRecord = (row: unknown): UsageRecord | null => { + if (!Array.isArray(row)) return null; + const [ + timestampMs, + modelIndex, + sessionIndex, + uncached, + cached, + creation, + output, + cost, + key, + ]: readonly unknown[] = row; + const model = stringAt(modelIndex); + const sessionId = stringAt(sessionIndex); + if ( + !isFiniteNumber(timestampMs) || + model === undefined || + sessionId === undefined || + !isFiniteNumber(uncached) || + !isFiniteNumber(cached) || + !isFiniteNumber(creation) || + !isFiniteNumber(output) || + (cost !== null && !isFiniteNumber(cost)) || + (key !== null && typeof key !== "string") + ) { + return null; + } + return { + provider: "cursor", + timestampMs, + model, + rateModel: cursorRateModel(model), + sessionId, + totals: { + uncachedInputTokens: uncached, + cachedInputTokens: cached, + cacheCreationTokens: creation, + outputTokens: output, + reasoningTokens: 0, + }, + reportedCostUsd: cost, + speed: "standard", + dedupeKey: key, + }; + }; + + for (const [credentialKey, raw] of Object.entries(root.accounts)) { + if (typeof raw !== "object" || raw === null) continue; + const account = raw as Partial>; + if ( + typeof account.accountKey !== "string" || + !isFiniteNumber(account.sinceMs) || + !isFiniteNumber(account.untilMs) || + !isFiniteNumber(account.fetchedAtMs) || + !Array.isArray(account.records) + ) { + continue; + } + const records: UsageRecord[] = []; + for (const row of account.records) { + const record = decodeRecord(row); + if (record === null) break; + records.push(record); + } + if (records.length !== account.records.length) continue; + caches.set(credentialKey, { + accountKey: account.accountKey, + sinceMs: account.sinceMs, + untilMs: account.untilMs, + fetchedAtMs: account.fetchedAtMs, + records, + }); + } + return caches; +} diff --git a/apps/server/src/usage/cursorUsageReader.ts b/apps/server/src/usage/cursorUsageReader.ts index dd1860e7cbde..86ba226a0529 100644 --- a/apps/server/src/usage/cursorUsageReader.ts +++ b/apps/server/src/usage/cursorUsageReader.ts @@ -4,6 +4,10 @@ import * as NodeFSP from "node:fs/promises"; import * as NodeCrypto from "node:crypto"; import * as NodeTimersPromises from "node:timers/promises"; +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; + import type { UsageRecord } from "./usageTranscripts.ts"; import { CursorKeychainTimeoutError, @@ -289,3 +293,24 @@ export async function readCursorAccountUsage( cancel.abort(); } } + +/** Reads one range of Cursor account usage. A service so tests can stand in for Cursor's API. */ +export class CursorAccountReader extends Context.Service< + CursorAccountReader, + { + readonly read: ( + credentialSource: string | { readonly kind: "keychain" }, + sinceMs: number, + untilMs: number, + ) => Effect.Effect; + } +>()("t3/usage/cursorUsageReader/CursorAccountReader") {} + +/** Reads Cursor's dashboard API with the saved CLI or Keychain login. */ +export const layer = Layer.succeed( + CursorAccountReader, + CursorAccountReader.of({ + read: (credentialSource, sinceMs, untilMs) => + Effect.promise(() => readCursorAccountUsage(credentialSource, sinceMs, untilMs)), + }), +); diff --git a/apps/web/src/components/usage/UsagePage.refresh.test.tsx b/apps/web/src/components/usage/UsagePage.refresh.test.tsx index 61ad487ade70..bbe53e36936a 100644 --- a/apps/web/src/components/usage/UsagePage.refresh.test.tsx +++ b/apps/web/src/components/usage/UsagePage.refresh.test.tsx @@ -50,6 +50,7 @@ vi.mock("../../state/usage", () => ({ }, ], isPending: false, + shown: null, isPartial: false, refresh: async () => undefined, }), diff --git a/apps/web/src/components/usage/UsagePage.test.tsx b/apps/web/src/components/usage/UsagePage.test.tsx index 6118595c41bf..1ba173c3635b 100644 --- a/apps/web/src/components/usage/UsagePage.test.tsx +++ b/apps/web/src/components/usage/UsagePage.test.tsx @@ -75,6 +75,7 @@ beforeEach(() => { environments, selectedEnvironments: environments, isPending: false, + shown: null, isPartial: false, refresh: vi.fn(), }); diff --git a/apps/web/src/components/usage/UsagePage.tsx b/apps/web/src/components/usage/UsagePage.tsx index 2bda301c2dd3..93065e97994d 100644 --- a/apps/web/src/components/usage/UsagePage.tsx +++ b/apps/web/src/components/usage/UsagePage.tsx @@ -8,18 +8,17 @@ import { type EnvironmentId, type UsageProviderKind, } from "@t3tools/contracts"; -import { - CircleAlertIcon, - ChevronDownIcon, - CircleDashedIcon, - InfoIcon, - SlidersHorizontalIcon, -} from "lucide-react"; +import { CircleAlertIcon, ChevronDownIcon, InfoIcon, SlidersHorizontalIcon } from "lucide-react"; import { useEffect, useEffectEvent, useMemo, useRef, useState } from "react"; import { cursorKeychainAccessEnvironments, refreshUsageLimits, } from "@t3tools/client-runtime/state/usage"; +import { + updatingProvidersLabel, + usageEnvironmentProgress, + usageLoadingState, +} from "@t3tools/client-runtime/state/usage-progress"; import { isCompatibleUsageContractVersion, @@ -109,6 +108,8 @@ function isUsageWindowDays(value: number): value is UsagePagePreferences["window return WINDOW_OPTIONS.some((option) => option.days === value); } +const providerLabel = (provider: UsageProviderKind) => PROVIDER_PRESENTATION[provider].label; + export function UsagePage() { const [preferences, setPreferences] = useState(readUsagePagePreferences); useEscapeToGoBack(); @@ -145,11 +146,36 @@ export function UsagePage() { ); const { days: windowDays, window } = windowSelection; const isPast24Hours = windowDays === 1; - const { merged, environments, selectedEnvironments, isPending, isPartial, refresh } = useUsage( - window, - selectedEnvironmentIds, - hiddenProviders, - ); + const { + merged: answeredUsage, + environments, + selectedEnvironments, + shown, + isPartial, + refresh, + } = useUsage(window, selectedEnvironmentIds, hiddenProviders); + // Until a new window's first answer, the previous one stays on screen, muted. + const merged = shown?.merged ?? answeredUsage; + const shownWindow = shown?.window ?? window; + const shownHourly = shownWindow.resolution === "hour"; + const refreshingUsage = isRefreshing && !showingLimits; + // Usage kept from another window is all old, so every figure stays muted. + const showingKept = shown !== null && shown.window !== window; + const loading = useMemo(() => { + if (showingKept) { + return { partial: true, everyProvider: true, providers: new Set() }; + } + const state = usageLoadingState(selectedEnvironments, refreshingUsage); + // A hidden provider still refreshing must not add its row or mute the totals. + const providers = new Set([...state.providers].filter((p) => !hiddenProviders.has(p))); + return { + partial: state.everyProvider || providers.size > 0, + everyProvider: state.everyProvider, + providers, + }; + }, [showingKept, refreshingUsage, selectedEnvironments, hiddenProviders]); + const isProviderLoading = (provider: UsageProviderKind) => + loading.everyProvider || loading.providers.has(provider); const presentations = useAtomValue(environmentPresentations.presentationsAtom); const cursorAccessEnvironments = hiddenProviders.has("cursor") ? [] @@ -180,21 +206,21 @@ export function UsagePage() { ); const days = useMemo( - () => enumerateDays(window.sinceDay, window.untilDay), - [window.sinceDay, window.untilDay], + () => enumerateDays(shownWindow.sinceDay, shownWindow.untilDay), + [shownWindow.sinceDay, shownWindow.untilDay], ); const hours = useMemo( () => - window.sinceTime === undefined || window.untilTime === undefined + shownWindow.sinceTime === undefined || shownWindow.untilTime === undefined ? [] - : enumerateHourStarts(window.sinceTime, window.untilTime), - [window.sinceTime, window.untilTime], + : enumerateHourStarts(shownWindow.sinceTime, shownWindow.untilTime), + [shownWindow.sinceTime, shownWindow.untilTime], ); // Newest first: the window can run 90 periods, so the interesting end // belongs at the top of the table. const breakdownPeriods = useMemo( - () => (isPast24Hours ? merged.hourly : merged.daily).toReversed(), - [isPast24Hours, merged.daily, merged.hourly], + () => (shownHourly ? merged.hourly : merged.daily).toReversed(), + [shownHourly, merged.daily, merged.hourly], ); const breakdownModels = useMemo( () => @@ -203,7 +229,24 @@ export function UsagePage() { : merged.models, [breakdown, merged.models, metric], ); - const activeProviders = useMemo(() => providersWithUsage(merged.providers), [merged.providers]); + const providersWithData = useMemo(() => providersWithUsage(merged.providers), [merged.providers]); + // A provider still refreshing keeps its row and line before its usage lands. + const activeProviders = useMemo( + () => + PROVIDER_ORDER.filter( + (provider) => providersWithData.includes(provider) || loading.providers.has(provider), + ), + [loading.providers, providersWithData], + ); + const chartLoadingProviders = useMemo( + () => + new Set( + activeProviders.filter( + (provider) => loading.everyProvider || loading.providers.has(provider), + ), + ), + [activeProviders, loading.everyProvider, loading.providers], + ); const selectedModel = selectedModelKey === null ? undefined @@ -332,10 +375,11 @@ export function UsagePage() { if (showingLimits && connectedLimitsEnvironments) autoRefreshLimits(); }, [showingLimits, connectedLimitsEnvironments]); + // Names the period on screen, which is the previous one until the new one answers. const windowLabel = - isPast24Hours && window.sinceTime !== undefined && window.untilTime !== undefined - ? `${formatDateTimeShort(window.sinceTime, window.timeZone)} to ${formatDateTimeShort(window.untilTime, window.timeZone)}` - : `${formatDayShort(window.sinceDay)} to ${formatDayShort(window.untilDay)}`; + shownHourly && shownWindow.sinceTime !== undefined && shownWindow.untilTime !== undefined + ? `${formatDateTimeShort(shownWindow.sinceTime, shownWindow.timeZone)} to ${formatDateTimeShort(shownWindow.untilTime, shownWindow.timeZone)}` + : `${formatDayShort(shownWindow.sinceDay)} to ${formatDayShort(shownWindow.untilDay)}`; const topbarContent = (
@@ -350,6 +394,7 @@ export function UsagePage() { selectedEnvironmentIds={selectedEnvironmentIds} onSelectionChange={setSelectedEnvironmentIds} showUsageStatus={!showingLimits} + refreshing={refreshingUsage} isPartial={isPartial} duplicateSources={merged.duplicateSources} contractMismatches={merged.contractMismatches} @@ -510,7 +555,7 @@ export function UsagePage() { ) : null } /> - ) : isPending ? ( + ) : shown === null ? ( ) : !canReadDiagnostics ? (
@@ -522,7 +567,7 @@ export function UsagePage() { ))}
) : ( - <> +
{sourceMessages.map((message) => (

{message} @@ -531,13 +576,20 @@ export function UsagePage() {

- + {metric === "cost" ? formatUsd(merged.costUsd) : formatTokens(merged.totalTokens)} - {formatCount(merged.sessions)} sessions + + {formatCount(merged.sessions)} sessions + {metric === "cost" && ( <> {" · API estimate"} @@ -600,6 +652,8 @@ export function UsagePage() { const sessionLabel = `${formatCount(providerSessions)} ${ providerSessions === 1 ? "session" : "sessions" }`; + const providerLoading = isProviderLoading(provider); + const awaitingData = providerLoading && !providersWithData.includes(provider); return (
@@ -616,18 +670,38 @@ export function UsagePage() { {PROVIDER_PRESENTATION[provider].label} - + {sessionLabel} - - {metric === "cost" - ? formatUsd(totals?.costUsd ?? 0) - : formatTokens(totals?.totalTokens ?? 0)} + + {awaitingData + ? "—" + : metric === "cost" + ? formatUsd(totals?.costUsd ?? 0) + : formatTokens(totals?.totalTokens ?? 0)}
- + {/* Kept while awaiting data so the row does not grow when it lands. */} + {metric === "cost" ? `${formatPercent(share)} of cost · ${formatTokens(totals?.totalTokens ?? 0)} tokens` : `${formatPercent(share)} of tokens · ${formatUsd(totals?.costUsd ?? 0)}`} @@ -639,19 +713,20 @@ export function UsagePage() {

- {isPast24Hours ? "Hourly" : "Daily"}{" "} + {shownHourly ? "Hourly" : "Daily"}{" "} {metric === "tokens" ? "processed tokens" : "cost"}

@@ -659,14 +734,28 @@ export function UsagePage() {

Totals

- - + + - + @@ -674,7 +763,12 @@ export function UsagePage() {
{merged.totalTokens > 0 ? ( -
+
{metric === "tokens" ? ( ( @@ -752,6 +846,7 @@ export function UsagePage() { model, metric === "tokens" ? "tokens" : "cost", ); + const rowFigures = figureClass(isProviderLoading(model.provider)); return ( {model.model} -
+
- + {isModelCostUnknown(model) ? ( Unpriced ) : ( formatUsd(model.costUsd) )} - + {share === null ? "" : formatPercent(share)} - {formatTokens(model.totalTokens)} + + {formatTokens(model.totalTokens)} + ); }) @@ -813,7 +913,7 @@ export function UsagePage() { - {isPast24Hours ? "Hour" : "Day"} + {shownHourly ? "Hour" : "Day"} {activeProviders.map((provider) => ( {PROVIDER_PRESENTATION[provider].label} @@ -841,21 +941,34 @@ export function UsagePage() { > {"hourStart" in period - ? formatHourShort(period.hourStart, window.timeZone) + ? formatHourShort(period.hourStart, shownWindow.timeZone) : formatDayShort(period.day)} {activeProviders.map((provider) => ( {formatUsd(period.byProvider.get(provider)?.costUsd ?? 0)} ))} - + {formatUsd(period.costUsd)} - + {formatTokens(period.totalTokens)} @@ -865,7 +978,7 @@ export function UsagePage() { )}
- +
)} @@ -878,9 +991,9 @@ export function UsagePage() { chartWindow={{ days, hours, - resolution: isPast24Hours ? "hour" : "day", - timeZone: window.timeZone, - referenceTime: window.untilTime, + resolution: shownHourly ? "hour" : "day", + timeZone: shownWindow.timeZone, + referenceTime: shownWindow.untilTime, }} onSetPrice={() => { setSelectedModelKey(null); @@ -1052,11 +1165,28 @@ function ProviderMark({ ); } -function Metric({ label, value }: { readonly label: string; readonly value: string }) { +/** Mutes a figure that is still coming in. The delay keeps a quick answer from flashing. */ +function figureClass(loading: boolean) { + return cn("transition-opacity", loading && "opacity-40 delay-150"); +} + +function Metric({ + label, + value, + loading, +}: { + readonly label: string; + readonly value: string; + readonly loading: boolean; +}) { return (
{label} - {value} + + {value} +
); } @@ -1108,13 +1238,14 @@ function UsageCoverageNotice({ ); } -/** Environment selection and scan progress share a permanent header control. */ +/** Environment selection, with each environment's scan status in the menu. */ function UsageEnvironmentFilter({ environments, selectedEnvironments, selectedEnvironmentIds, onSelectionChange, showUsageStatus, + refreshing, isPartial, duplicateSources, contractMismatches, @@ -1125,6 +1256,7 @@ function UsageEnvironmentFilter({ readonly selectedEnvironmentIds: ReadonlySet | null; readonly onSelectionChange: (ids: ReadonlySet | null) => void; readonly showUsageStatus: boolean; + readonly refreshing: boolean; readonly isPartial: boolean; readonly duplicateSources: readonly string[]; readonly contractMismatches: MergedUsage["contractMismatches"]; @@ -1136,10 +1268,6 @@ function UsageEnvironmentFilter({ : selectedEnvironments.length === 1 ? selectedEnvironments[0]!.label : `${selectedEnvironments.length} environments`; - const pendingCount = selectedEnvironments.filter( - (environment) => - environment.error === null && (environment.isPending || environment.summary === null), - ).length; const hasIssue = selectedEnvironments.some((environment) => environment.error !== null) || contractMismatches.length > 0; @@ -1149,15 +1277,7 @@ function UsageEnvironmentFilter({ } className="group/usage-environment min-w-0 max-w-full"> {label} - {showUsageStatus && pendingCount > 0 ? ( - <> - - - {pendingCount} {pendingCount === 1 ? "environment" : "environments"} still scanning - {isPartial ? "; totals are partial" : ""} - - - ) : showUsageStatus && hasIssue ? ( + {showUsageStatus && hasIssue ? ( { + const column = (codex: number) => ({ + total: codex, + bands: [{ provider: "codex" as const, value: codex }], + }); + + it("holds unlabeled placeholder gridlines while loading providers have nothing to show", () => { + const scale = chartScale([column(0)], new Set(["codex" as const])); + + expect(scale.labeled).toBe(false); + expect(scale.ticks.length).toBeGreaterThan(1); + }); + + it("scales to what is on screen, loading or not", () => { + expect(chartScale([column(40)], new Set(["codex" as const]))).toMatchObject({ + max: 40, + labeled: true, + }); + }); +}); + describe("niceScale", () => { it("never puts the peak above the top of the scale", () => { // Regression: an earlier version stopped at the last step below the peak, diff --git a/apps/web/src/components/usage/UsageProviderChart.tsx b/apps/web/src/components/usage/UsageProviderChart.tsx index 62772a6478dc..d8571f0b36d8 100644 --- a/apps/web/src/components/usage/UsageProviderChart.tsx +++ b/apps/web/src/components/usage/UsageProviderChart.tsx @@ -10,17 +10,21 @@ import { formatTokens, formatUsd, } from "@t3tools/shared/usageFormat"; +import { cn } from "~/lib/utils"; import { PROVIDER_ORDER, PROVIDER_PRESENTATION } from "./usageProviders"; const VIEW_WIDTH = 960; const VIEW_HEIGHT = 260; const TICK_COUNT = 4; const PLOT_TOP = 8; +const NONE_LOADING: ReadonlySet = new Set(); export type UsageChartMetric = "tokens" | "cost"; interface UsageProviderChartProps { readonly providers: readonly UsageProviderKind[]; + /** Providers whose figures are still coming in: their lines are muted. */ + readonly loadingProviders?: ReadonlySet; readonly days: readonly string[]; readonly daily: readonly DailyTotals[]; readonly hours: readonly string[]; @@ -148,6 +152,10 @@ function curvePath(segments: readonly CurveSegment[]): string { return path; } +function areaPath(line: string) { + return line === "" ? "" : `${line} L${VIEW_WIDTH},${VIEW_HEIGHT} L0,${VIEW_HEIGHT} Z`; +} + /** * Builds a scale whose maximum is a readable 1/2/5 x 10^n step at or above the * peak. @@ -170,8 +178,71 @@ export function niceScale(peak: number, count: number): { max: number; ticks: re return { max, ticks }; } +const PLACEHOLDER_TICKS = Array.from({ length: TICK_COUNT + 1 }, (_, index) => index); + +/** + * Scales to the largest single provider-period. With nothing to show yet while + * providers load, unlabeled placeholder gridlines hold their usual spacing so + * nothing shifts on arrival. + */ +export function chartScale( + columns: readonly DayColumn[], + loadingProviders: ReadonlySet, +) { + // Not the sum: layered series each measure from zero, so a combined peak + // would leave the plot permanently half empty. + const peak = columns.reduce( + (max, column) => column.bands.reduce((inner, band) => Math.max(inner, band.value), max), + 0, + ); + return peak === 0 && loadingProviders.size > 0 + ? { max: TICK_COUNT, ticks: PLACEHOLDER_TICKS, labeled: false } + : { ...niceScale(peak, TICK_COUNT), labeled: true }; +} + +// Leave room above the top gridline so the constant-width stroke is not +// clipped when a series reaches the peak. +function valueToY(value: number, max: number) { + return max === 0 ? VIEW_HEIGHT : VIEW_HEIGHT - (value / max) * (VIEW_HEIGHT - PLOT_TOP); +} + +/** Per-provider paths in paint order, heaviest first. */ +function buildChart( + periods: readonly string[], + byPeriod: ReadonlyMap, + metric: UsageChartMetric, + providers: readonly UsageProviderKind[], + loadingProviders: ReadonlySet, +) { + const columns = buildPeriodColumns(periods, byPeriod, metric); + const scale = chartScale(columns, loadingProviders); + const stepX = periods.length < 2 ? 0 : VIEW_WIDTH / (periods.length - 1); + const paths = providers.map((provider) => { + const slot = PROVIDER_ORDER.indexOf(provider); + const line = curvePath( + smoothCurve( + columns.map((column, periodIndex) => ({ + x: periodIndex * stepX, + y: valueToY(column.bands[slot]?.value ?? 0, scale.max), + })), + ), + ); + return { + provider, + loading: loadingProviders.has(provider), + total: columns.reduce((sum, column) => sum + (column.bands[slot]?.value ?? 0), 0), + line, + area: areaPath(line), + }; + }); + + // Paint the heavier series first so the lighter one is not buried. + return { columns, scale, stepX, paths: paths.toSorted((a, b) => b.total - a.total) }; +} + export function UsageProviderChart({ providers, + loadingProviders = NONE_LOADING, days, daily, hours, @@ -194,59 +265,14 @@ export function UsageProviderChart({ const tooltipRef = useRef(null); const hoverPositionRef = useRef<{ x: number; y: number } | null>(null); - const { paths, ticks, stepX, toY, series } = useMemo(() => { - if (periods.length === 0) { - return { - paths: [], - series: [] as readonly DayColumn[], - stepX: 0, - ticks: [0] as readonly number[], - toY: () => VIEW_HEIGHT, - }; - } - - const columns = buildPeriodColumns(periods, byPeriod, metric); - // The scale tops out at the largest single provider-period, not the sum: - // layered series each measure from zero, so a combined peak would leave - // the plot permanently half empty. - const peak = columns.reduce( - (max, column) => column.bands.reduce((inner, band) => Math.max(inner, band.value), max), - 0, - ); - const { max, ticks: tickValues } = niceScale(peak, TICK_COUNT); - const step = periods.length === 1 ? 0 : VIEW_WIDTH / (periods.length - 1); - // Leave room above the top gridline so the constant-width stroke is not - // clipped when a series reaches the peak. - const toY = (value: number) => - max === 0 ? VIEW_HEIGHT : VIEW_HEIGHT - (value / max) * (VIEW_HEIGHT - PLOT_TOP); - - const built = providers.map((provider) => { - const providerIndex = PROVIDER_ORDER.indexOf(provider); - const line = curvePath( - smoothCurve( - columns.map((column, periodIndex) => ({ - x: periodIndex * step, - y: toY(column.bands[providerIndex]?.value ?? 0), - })), - ), - ); - return { - provider, - total: columns.reduce((sum, column) => sum + (column.bands[providerIndex]?.value ?? 0), 0), - area: line === "" ? "" : `${line} L${VIEW_WIDTH},${VIEW_HEIGHT} L0,${VIEW_HEIGHT} Z`, - line, - }; - }); - - // Paint the heavier series first so the lighter one is not buried. - return { - paths: built.toSorted((a, b) => b.total - a.total), - series: columns, - stepX: step, - ticks: tickValues, - toY, - }; - }, [byPeriod, metric, periods, providers]); + const { columns, scale, paths, stepX } = useMemo( + () => buildChart(periods, byPeriod, metric, providers, loadingProviders), + [byPeriod, loadingProviders, metric, periods, providers], + ); + const toY = (value: number) => valueToY(value, scale.max); + // The delay keeps a quick answer from flashing, as with the page's figures. + const seriesClassName = (loading: boolean) => + cn("transition-opacity", loading && "opacity-40 delay-150"); const format = metric === "tokens" ? formatTokens : formatUsd; @@ -307,7 +333,8 @@ export function UsageProviderChart({ ); const hoveredPeriod = hoverIndex === null ? undefined : periods[hoverIndex]; - const hoveredColumn = hoverIndex === null ? undefined : series[hoverIndex]; + const hoveredColumn = hoverIndex === null ? undefined : columns[hoverIndex]; + const partial = providers.some((provider) => loadingProviders.has(provider)); const formatPeriod = (period: string) => resolution === "hour" ? formatHourShort(period, timeZone) : formatDayShort(period); const formatTooltipPeriod = (period: string) => @@ -320,15 +347,17 @@ export function UsageProviderChart({
{/* Axis labels sit outside the plot so they stay aligned to gridlines. */}
- {ticks.map((tick) => ( - - {tick === 0 ? "0" : format(tick)} - - ))} + {scale.labeled + ? scale.ticks.map((tick) => ( + + {tick === 0 ? "0" : format(tick)} + + )) + : null}
- {ticks.map((tick) => { + {scale.ticks.map((tick) => { const y = toY(tick); return ( ( + {paths.map(({ provider, loading, area }) => ( ))} - {paths.map(({ provider, line }) => ( + {paths.map(({ provider, loading, line }) => ( {label} - + {format( hoveredColumn?.bands.find((band) => band.provider === provider)?.value ?? 0, )} @@ -430,7 +468,12 @@ export function UsageProviderChart({ })}
Total - + {format(hoveredColumn?.total ?? 0)}
diff --git a/apps/web/src/state/usage.test.tsx b/apps/web/src/state/usage.test.tsx index c084d0a5e9ab..7b2efd76b16d 100644 --- a/apps/web/src/state/usage.test.tsx +++ b/apps/web/src/state/usage.test.tsx @@ -33,6 +33,7 @@ function environment( label: id, isPending: cost === null, canReadDiagnostics: true, + isConnected: true, error: null, needsCursorKeychainAccess: false, summary: @@ -90,11 +91,13 @@ let latest: UsageView; function Probe({ selected, hidden, + window = input, }: { selected: ReadonlySet | null; hidden?: ReadonlySet; + window?: typeof input; }) { - const usage = useUsage(input, selected, hidden); + const usage = useUsage(window, selected, hidden); useLayoutEffect(() => { latest = usage; }, [usage]); @@ -181,6 +184,45 @@ describe("usage environment selection", () => { expect(latest.isPending).toBe(false); expect(latest.isPartial).toBe(false); }); + + it("shows the last answered usage until the next window answers, for the same selection", async () => { + const selected = new Set([EnvironmentId.make("a")]); + await act(() => renderer?.update()); + expect(latest.shown?.merged.costUsd).toBe(10); + + // A new window that nothing has answered yet. + const nextWindow = { ...input, sinceDay: UsageDay.make("2026-08-28") }; + testState.environments = [environment("a", null)]; + await act(() => renderer?.update()); + expect(latest.isPending).toBe(true); + expect(latest.shown?.window).toBe(input); + expect(latest.shown?.merged.costUsd).toBe(10); + + // A window that fails everywhere keeps it too. + testState.environments = [{ ...environment("a", null), isPending: false, error: "Offline" }]; + await act(() => renderer?.update()); + expect(latest.isPending).toBe(false); + expect(latest.shown?.window).toBe(input); + expect(latest.shown?.merged.costUsd).toBe(10); + + // Once the new window answers, it replaces the kept one. + testState.environments = [environment("a", 30)]; + await act(() => renderer?.update()); + expect(latest.shown?.window).toBe(nextWindow); + expect(latest.shown?.merged.costUsd).toBe(30); + + // A different provider filter does not reuse usage merged with the old one. + testState.environments = [environment("a", null)]; + await act(() => + renderer?.update(