Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 72 additions & 1 deletion apps/server/src/usage/UsageService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import * as NodePath from "node:path";

import { assert, describe, it } from "@effect/vitest";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { HostProcessEnvironment } from "@t3tools/shared/hostProcess";
import { HostProcessEnvironment, HostProcessPlatform } from "@t3tools/shared/hostProcess";
import { UsageDay, type UsageSummaryInput } from "@t3tools/contracts";
import * as Duration from "effect/Duration";
import * as Deferred from "effect/Deferred";
Expand Down Expand Up @@ -218,6 +218,77 @@ describe("UsageService", () => {
}).pipe(Effect.scoped),
);

it.live("canonicalizes aliased homes and keeps missing homes visible", () =>
Effect.gen(function* () {
const platform = yield* HostProcessPlatform;
const home = yield* Effect.promise(() =>
NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "usage-service-homes-test-")),
);
yield* Effect.addFinalizer(() =>
Effect.promise(() => NodeFSP.rm(home, { recursive: true, force: true })),
);

const claudeHome = NodePath.join(home, "claude");
const aliasHome = NodePath.join(home, "claude-alias");
const missingHome = NodePath.join(home, "claude-missing");
const transcriptDir = NodePath.join(claudeHome, "projects", "proj");
yield* Effect.promise(() => NodeFSP.mkdir(transcriptDir, { recursive: true }));
yield* Effect.promise(() =>
NodeFSP.writeFile(NodePath.join(transcriptDir, "session.jsonl"), claudeLine(1, 5)),
);
yield* Effect.promise(() =>
NodeFSP.symlink(claudeHome, aliasHome, platform === "win32" ? "junction" : "dir"),
);

const settings = {
providerInstances: {
claudeAgent: {
driver: "claudeAgent" as const,
config: { homePath: claudeHome },
},
claude_alias: {
driver: "claudeAgent" as const,
config: { homePath: aliasHome },
},
claude_missing: {
driver: "claudeAgent" as const,
config: { homePath: missingHome },
},
codex: {
driver: "codex" as const,
config: { homePath: NodePath.join(home, "codex") },
},
},
};
const service = yield* UsageService.make.pipe(
Effect.provide(serviceLayers({ prefix: "usage-service-homes-test", home, settings })),
);

const summary = yield* service.readSummary(WINDOW);
const canonicalDir = yield* Effect.promise(() =>
NodeFSP.realpath(NodePath.join(claudeHome, "projects")),
);
const claudeSources = summary.sources.filter(
(source) => source.fingerprint.provider === "claude",
);
const canonicalSourceIndex = summary.sources.findIndex(
(source) => source.fingerprint.resolvedHomePath === canonicalDir,
);

assert.strictEqual(claudeSources.length, 2);
assert.strictEqual(claudeSources[0]?.fingerprint.resolvedHomePath, canonicalDir);
assert.strictEqual(claudeSources[0]?.status, "ok");
assert.strictEqual(
claudeSources[1]?.fingerprint.resolvedHomePath,
NodePath.join(missingHome, "projects"),
);
assert.strictEqual(claudeSources[1]?.status, "missing");
assert.strictEqual(summary.buckets.length, 1);
assert.strictEqual(summary.buckets[0]?.sourceIndex, canonicalSourceIndex);
assert.strictEqual(totalOutputTokens(summary), 5);
}).pipe(Effect.scoped),
);

it.live("shares one scan between concurrent identical requests", () =>
Effect.gen(function* () {
const { transcript, settings, home } = yield* setup;
Expand Down
70 changes: 45 additions & 25 deletions apps/server/src/usage/UsageService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,12 +40,10 @@ import * as Semaphore from "effect/Semaphore";
import { HttpClient, HttpClientResponse } from "effect/unstable/http";

import { ServerConfig } from "../config.ts";
import { expandHomePath } from "../pathExpansion.ts";
import * as ServerSettings from "../serverSettings.ts";
import { resolveClaudeHomePath } from "../provider/Drivers/ClaudeHome.ts";
import { resolveCodexHomeLayout } from "../provider/Drivers/CodexHomeLayout.ts";
import { UsageAggregator } from "./usageAggregation.ts";
import { createOverrideRateTable, parseRateTable, type RateTable } from "./usagePricing.ts";
import { resolveUsageProviderHomes } from "./usageProviderHomes.ts";
import {
listTranscriptFiles,
readDirectoryVolumeId,
Expand Down Expand Up @@ -245,30 +243,51 @@ export const make = Effect.gen(function* () {
),
);

/** Resolves the transcript directory for each provider. */
/**
* Resolves the transcript directories for each provider. Claude and Codex
* can be configured multiple times via provider instances, so both may
* contribute several directories; distinct instances sharing a home
* collapse to one entry so their records are not double counted.
*/
const resolveTranscriptDirs = Effect.fn("UsageService.resolveTranscriptDirs")(function* (
settings: ServerSettingsValue,
) {
const claudeHome = yield* resolveClaudeHomePath(settings.providers.claudeAgent);
const claudeDir = yield* resolveClaudeTranscriptDir(claudeHome);
const codexLayout = yield* resolveCodexHomeLayout(settings.providers.codex);
// Grok Settings only expose the binary path; home is `$GROK_HOME` or `~/.grok`.
// Empty/whitespace GROK_HOME must fall back: coalescing alone would scan cwd.
const grokHomeEnv = hostEnvironment["GROK_HOME"]?.trim() ?? "";
const grokHome =
grokHomeEnv.length > 0
? path.resolve(expandHomePath(grokHomeEnv))
: path.join(NodeOS.homedir(), ".grok");

return [
{ provider: "claude" as const, dir: claudeDir },
{ provider: "codex" as const, dir: path.join(codexLayout.sharedHomePath, "sessions") },
{
provider: "grok" as const,
dir: path.join(grokHome, "sessions"),
fileName: "updates.jsonl",
},
];
const homes = yield* resolveUsageProviderHomes(settings, hostEnvironment);

const dirs: Array<{
provider: UsageProviderKind;
dir: string;
fileName?: string;
}> = [];
const seen = new Set<string>();
// Two configured homes can name one physical transcript directory through
// symlinks (including Codex shadow overlays). Canonicalize the final
// transcript directory before de-duplicating it. If the directory does not
// exist, keep its configured path so the scan reports a missing source.
const pushDir = Effect.fn("UsageService.pushTranscriptDir")(function* (
provider: UsageProviderKind,
dir: string,
fileName?: string,
) {
const canonical = yield* fileSystem
.realPath(dir)
.pipe(Effect.catchCause(() => Effect.succeed(dir)));
const key = `${provider}\u0000${canonical}`;
if (seen.has(key)) return;
seen.add(key);
dirs.push({ provider, dir: canonical, ...(fileName === undefined ? {} : { fileName }) });
});

for (const home of homes.claudeHomePaths) {
// Distinct homes can probe to the same transcript dir (e.g. `~/x` with
// a nested `.claude` next to `~/x/.claude` itself), so dedupe post-probe.
yield* pushDir("claude", yield* resolveClaudeTranscriptDir(home));
}
for (const dir of homes.codexSessionDirs) {
yield* pushDir("codex", dir);
}
yield* pushDir("grok", homes.grokSessionsDir, "updates.jsonl");
return dirs;
Comment thread
cursor[bot] marked this conversation as resolved.
});

/**
Expand Down Expand Up @@ -482,6 +501,7 @@ export const make = Effect.gen(function* () {
const walkedRoots: string[] = [];

for (const { provider, dir, volumeId, files } of scannedDirs) {
const sourceIndex = sources.length;
if (files === null) {
sources.push({
fingerprint: { hostId, provider, resolvedHomePath: dir, volumeId },
Expand Down Expand Up @@ -512,7 +532,7 @@ export const make = Effect.gen(function* () {
for (const record of file.records) {
// Only sessions that contributed in-window count: the mtime slack
// admits boundary files whose records fall outside the range.
if (aggregator.add(record) && record.sessionId.length > 0) {
if (aggregator.add(record, sourceIndex) && record.sessionId.length > 0) {
sessionIds.add(record.sessionId);
}
}
Expand Down
25 changes: 21 additions & 4 deletions apps/server/src/usage/usageAggregation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ function aggregate(
...hourlyBounds,
rates,
});
for (const item of records) aggregator.add(item);
for (const item of records) aggregator.add(item, 0);
return aggregator.finish();
}

Expand All @@ -83,6 +83,7 @@ describe("UsageAggregator", () => {

expect(result.duplicatesDropped).toBe(2);
expect(result.buckets).toHaveLength(1);
expect(result.buckets[0]?.sourceIndex).toBe(0);
expect(result.buckets[0]?.records).toBe(1);
expect(result.buckets[0]?.totals.outputTokens).toBe(50);
});
Expand Down Expand Up @@ -187,9 +188,25 @@ describe("UsageAggregator", () => {
rates,
});

expect(aggregator.add(record({ dedupeKey: "msg_1:" }))).toBe(true);
expect(aggregator.add(record({ dedupeKey: "msg_1:" }))).toBe(false);
expect(aggregator.add(record({ timestampMs: Date.parse("2026-07-01T12:00:00Z") }))).toBe(false);
expect(aggregator.add(record({ dedupeKey: "msg_1:" }), 0)).toBe(true);
expect(aggregator.add(record({ dedupeKey: "msg_1:" }), 0)).toBe(false);
expect(aggregator.add(record({ timestampMs: Date.parse("2026-07-01T12:00:00Z") }), 0)).toBe(
false,
);
});

it("keeps otherwise identical buckets separate by transcript source", () => {
const aggregator = new UsageAggregator({
timeZone: "UTC",
sinceDay: "2026-08-01",
untilDay: "2026-08-31",
rates,
});

aggregator.add(record({ sessionId: "session-a" }), 0);
aggregator.add(record({ sessionId: "session-b" }), 1);

expect(aggregator.finish().buckets.map((bucket) => bucket.sourceIndex)).toEqual([0, 1]);
});

it("separates providers and models into their own buckets", () => {
Expand Down
15 changes: 9 additions & 6 deletions apps/server/src/usage/usageAggregation.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
// @effect-diagnostics globalDate:off
/**
* Folds parsed transcript records into `(day, hourStart?, provider, model)`
* buckets.
* Folds parsed transcript records into
* `(sourceIndex, day, hourStart?, provider, model)` buckets.
*
* `Intl.DateTimeFormat` is the only reliable way to resolve a wall-clock day in
* an arbitrary IANA zone, and it takes a `Date`. That is why the raw `Date`
Expand Down Expand Up @@ -112,7 +112,7 @@ export class UsageAggregator {
* can derive per-window facts (distinct sessions, for one) from the records
* that landed rather than everything the mtime prefilter happened to admit.
*/
add(record: UsageRecord): boolean {
add(record: UsageRecord, sourceIndex: number): boolean {
if (record.dedupeKey !== null) {
if (this.#seen.has(record.dedupeKey)) {
this.#duplicatesDropped += 1;
Expand Down Expand Up @@ -146,7 +146,7 @@ export class UsageAggregator {
this.#hourlyWindow.sinceTimeMs +
Math.floor((record.timestampMs - this.#hourlyWindow.sinceTimeMs) / HOUR_MS) * HOUR_MS,
).toISOString();
const key = `${day}\u0000${hourStart}\u0000${record.provider}\u0000${record.model}`;
const key = `${day}\u0000${hourStart}\u0000${record.provider}\u0000${record.model}\u0000${sourceIndex}`;
let bucket = this.#buckets.get(key);
if (bucket === undefined) {
bucket = {
Expand Down Expand Up @@ -187,8 +187,10 @@ export class UsageAggregator {
finish(): AggregateResult {
const buckets: UsageBucket[] = [];
for (const [key, bucket] of this.#buckets) {
const [day = "", hourStart = "", provider = "", model = ""] = key.split("\u0000");
const [day = "", hourStart = "", provider = "", model = "", sourceIndex = ""] =
key.split("\u0000");
buckets.push({
sourceIndex: Number(sourceIndex),
day: day as UsageDay,
...(hourStart === "" ? {} : { hourStart }),
provider: provider as UsageBucket["provider"],
Expand All @@ -208,7 +210,8 @@ export class UsageAggregator {
a.day.localeCompare(b.day) ||
(a.hourStart ?? "").localeCompare(b.hourStart ?? "") ||
a.provider.localeCompare(b.provider) ||
a.model.localeCompare(b.model),
a.model.localeCompare(b.model) ||
a.sourceIndex - b.sourceIndex,
);

return {
Expand Down
Loading
Loading