Skip to content
Open
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