Skip to content
Merged
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
5 changes: 3 additions & 2 deletions apps/server/src/provider/builtInDrivers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,8 @@ export const BUILT_IN_DRIVERS: ReadonlyArray<AnyProviderDriver<BuiltInDriversEnv
export type BuiltInUsageReadersEnv =
| ProviderUsageReaderEnv<typeof ClaudeDriver>
| ProviderUsageReaderEnv<typeof CodexDriver>
| ProviderUsageReaderEnv<typeof GrokDriver>;
| ProviderUsageReaderEnv<typeof GrokDriver>
| ProviderUsageReaderEnv<typeof OpenCodeDriver>;

/**
* The drivers that keep usage history, in the order the usage page reads
Expand All @@ -83,4 +84,4 @@ export type BuiltInUsageReadersEnv =
*/
export const BUILT_IN_USAGE_DRIVERS: ReadonlyArray<
AnyProviderDriver<BuiltInDriversEnv, BuiltInUsageReadersEnv>
> = [ClaudeDriver, CodexDriver, GrokDriver];
> = [ClaudeDriver, CodexDriver, GrokDriver, OpenCodeDriver];
66 changes: 37 additions & 29 deletions apps/server/src/usage/UsageService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,6 @@ import type {
TranscriptUsageFormat,
UsageRecord,
} from "@t3tools/provider-core/server/usage";
import { readOpenCodeUsage } from "./opencodeUsageReader.ts";
import { makeAntigravityUsageCache, readAntigravityUsage } from "./antigravityUsageReader.ts";
import {
CURSOR_ACCOUNT_CACHE_FILE_NAME,
Expand Down Expand Up @@ -127,6 +126,11 @@ const transcriptReaders = BUILT_IN_USAGE_DRIVERS.flatMap((driver) =>
driver.usage?.kind === "transcripts" ? [{ driver, reader: driver.usage }] : [],
);

/** The scan readers, in driver order. */
const scanReaders = BUILT_IN_USAGE_DRIVERS.flatMap((driver) =>
driver.usage?.kind === "scan" ? [{ driver, reader: driver.usage }] : [],
);

/** Transcript formats by provider, for decoding the persisted scan cache. */
const transcriptFormats = new Map(
transcriptReaders.map(({ reader }) => [reader.provider, reader.format] as const),
Expand Down Expand Up @@ -851,32 +855,36 @@ export const make = Effect.gen(function* () {
}
return [...canonical];
});
const dataHome = hostEnvironment["XDG_DATA_HOME"]?.trim();

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" },
);
});
const scans = Effect.forEach(
scanReaders,
({ driver, reader }) =>
reader
.scan({
instances: usageInstances(driver, settings),
windowStartMs,
retentionCutoffMs,
awaitRefresh,
})
.pipe(
Effect.flatMap((sources) =>
Effect.forEach(sources, ({ volumeId, ...source }) =>
Effect.map(
volumeId === undefined
? Effect.promise(() => readDirectoryVolumeId(source.dir))
: Effect.succeed(volumeId),
(resolved): ScannedDir => ({
...source,
provider: reader.provider,
volumeId: resolved,
}),
),
),
),
Effect.provideContext(readerContext),
),
{ concurrency: "unbounded" },
).pipe(Effect.map((sources) => sources.flat()));

const antigravity = Effect.gen(function* () {
const antigravityRoots = yield* envRoots("ANTIGRAVITY_DATA_DIR", [
Expand Down Expand Up @@ -999,18 +1007,18 @@ export const make = Effect.gen(function* () {
// 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(
const [transcripts, scanDirs, antigravityDirs, cursorDirs] = yield* Effect.all(
[
Effect.forEach(dirs, (dir) => scanTranscriptDir(dir, windowStartMs)),
openCode,
scans,
antigravity,
cursor,
],
{ concurrency: "unbounded" },
);
const scanned: readonly ScannedDir[] = [
...transcripts,
...openCodeDirs,
...scanDirs,
...antigravityDirs,
...cursorDirs,
];
Expand Down
50 changes: 0 additions & 50 deletions apps/server/src/usage/usageTranscriptReader.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ import { afterEach, assert, beforeEach, describe, it } from "@effect/vitest";

import { TEST_FORMATS } from "./usageTestFormats.ts";
import { readTranscriptRecords } from "./usageTranscriptReader.ts";
import { readOpenCodeUsage } from "./opencodeUsageReader.ts";
import { readCursorAccountUsage } from "./cursorUsageReader.ts";
import { makeAntigravityUsageCache, readAntigravityUsage } from "./antigravityUsageReader.ts";

Expand Down Expand Up @@ -482,55 +481,6 @@ describe("SQLite usage readers", () => {
assert.isFalse(requested);
});

it("counts migrated OpenCode messages once and sees subsequent WAL writes", async () => {
const db = new NodeSqlite.DatabaseSync(NodePath.join(dir, "opencode.db"));
try {
db.exec(
"PRAGMA journal_mode = WAL; PRAGMA wal_autocheckpoint = 0; CREATE TABLE message (id TEXT, session_id TEXT, data TEXT)",
);
const message = {
id: "msg-1",
sessionID: "session-1",
role: "assistant",
modelID: "claude-sonnet-4-5",
time: { created: 1780000000000 },
cost: 0.25,
tokens: { input: 100, output: 20, reasoning: 5, cache: { read: 30, write: 10 } },
};
const insert = db.prepare("INSERT INTO message VALUES (?, ?, ?)");
insert.run(message.id, message.sessionID, JSON.stringify(message));
const legacy = NodePath.join(dir, "storage", "message", message.sessionID);
await NodeFSP.mkdir(legacy, { recursive: true });
await NodeFSP.writeFile(NodePath.join(legacy, "msg-1.json"), JSON.stringify(message));
const first = await readOpenCodeUsage(dir, 0);
assert.isFalse(first.error);
const records = first.files.flatMap((file) => file.records);
assert.strictEqual(records.length, 1);
assert.deepStrictEqual(records[0]?.totals, {
uncachedInputTokens: 100,
cachedInputTokens: 30,
cacheCreationTokens: 10,
outputTokens: 25,
reasoningTokens: 5,
});
assert.strictEqual(records[0]?.reportedCostUsd, 0.25);
insert.run(
"msg-2",
message.sessionID,
JSON.stringify({ ...message, id: "msg-2", time: { created: 1780000001000 } }),
);
const next = await readOpenCodeUsage(dir, 1780000001000);
assert.isFalse(next.error);
assert.deepStrictEqual(
next.files.flatMap((file) => file.records).map((record) => record.dedupeKey),
["opencode:msg-2"],
);
assert.isAbove((await NodeFSP.stat(NodePath.join(dir, "opencode.db-wal"))).size, 0);
} finally {
db.close();
}
});

it("deduplicates Antigravity generation and step usage while preserving retry model and token buckets", async () => {
const db = new NodeSqlite.DatabaseSync(NodePath.join(dir, "session-1.db"));
const stamp = protoNumber(1, 1780000000);
Expand Down
3 changes: 2 additions & 1 deletion packages/provider-core/src/server/usage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,8 @@ export interface ProviderUsageInstance<Config> {
/** One source a `scan` reader read. */
export interface ProviderUsageScan {
readonly dir: string;
readonly volumeId: string;
/** Identity of the source across hosts. Defaults to the filesystem identity of `dir`. */
readonly volumeId?: string;
/** Overrides the server's host name in the source fingerprint, for account-wide sources. */
readonly hostId?: string;
readonly status?: UsageSource["status"];
Expand Down
8 changes: 7 additions & 1 deletion packages/provider-opencode/src/server/driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import * as OpenCode2AdapterV2 from "./v2/adapter.ts";
import type * as ProviderAdapter from "@t3tools/provider-core/server/ProviderAdapter";
import type { ProviderTextGeneration } from "@t3tools/provider-core/server/textGeneration";
import { ProviderDriverError } from "@t3tools/provider-core/server/errors";
import { openCodeUsageReader, type OpenCodeUsageReaderEnv } from "./usage.ts";
import { readOpenCodeGoUsageLimits } from "./usageLimits.ts";
import {
checkOpenCodeProviderStatus,
Expand Down Expand Up @@ -182,14 +183,19 @@ export type OpenCodeDriverEnv =
| OpenCodeRuntime.OpenCodeRuntime
| Path.Path;

export const OpenCodeDriver: ProviderDriver<OpenCodeSettings, OpenCodeDriverEnv> = {
export const OpenCodeDriver: ProviderDriver<
OpenCodeSettings,
OpenCodeDriverEnv,
OpenCodeUsageReaderEnv
> = {
driverKind: DRIVER_KIND,
metadata: {
displayName: "OpenCode",
supportsMultipleInstances: true,
},
configSchema: OpenCodeSettings,
defaultConfig: (): OpenCodeSettings => decodeOpenCodeSettings({}),
usage: openCodeUsageReader,
create: ({ instanceId, displayName, accentColor, environment, enabled, config }) =>
Effect.gen(function* () {
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
Expand Down
71 changes: 71 additions & 0 deletions packages/provider-opencode/src/server/usage.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
// @effect-diagnostics nodeBuiltinImport:off - the reader under test opens real
// SQLite databases and walks a real legacy JSON store.
import * as NodeFSP from "node:fs/promises";
import * as NodeOS from "node:os";
import * as NodePath from "node:path";
import * as NodeSqlite from "node:sqlite";

import { afterEach, assert, beforeEach, describe, it } from "@effect/vitest";

import { readOpenCodeUsage } from "./usage.ts";

let dir: string;

beforeEach(async () => {
dir = await NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "usage-reader-test-"));
});

afterEach(async () => {
await NodeFSP.rm(dir, { recursive: true, force: true });
});

describe("readOpenCodeUsage", () => {
it("counts migrated OpenCode messages once and sees subsequent WAL writes", async () => {
const db = new NodeSqlite.DatabaseSync(NodePath.join(dir, "opencode.db"));
try {
db.exec(
"PRAGMA journal_mode = WAL; PRAGMA wal_autocheckpoint = 0; CREATE TABLE message (id TEXT, session_id TEXT, data TEXT)",
);
const message = {
id: "msg-1",
sessionID: "session-1",
role: "assistant",
modelID: "claude-sonnet-4-5",
time: { created: 1780000000000 },
cost: 0.25,
tokens: { input: 100, output: 20, reasoning: 5, cache: { read: 30, write: 10 } },
};
const insert = db.prepare("INSERT INTO message VALUES (?, ?, ?)");
insert.run(message.id, message.sessionID, JSON.stringify(message));
const legacy = NodePath.join(dir, "storage", "message", message.sessionID);
await NodeFSP.mkdir(legacy, { recursive: true });
await NodeFSP.writeFile(NodePath.join(legacy, "msg-1.json"), JSON.stringify(message));
const first = await readOpenCodeUsage(dir, 0);
assert.isFalse(first.error);
const records = first.files.flatMap((file) => file.records);
assert.strictEqual(records.length, 1);
assert.deepStrictEqual(records[0]?.totals, {
uncachedInputTokens: 100,
cachedInputTokens: 30,
cacheCreationTokens: 10,
outputTokens: 25,
reasoningTokens: 5,
});
assert.strictEqual(records[0]?.reportedCostUsd, 0.25);
insert.run(
"msg-2",
message.sessionID,
JSON.stringify({ ...message, id: "msg-2", time: { created: 1780000001000 } }),
);
const next = await readOpenCodeUsage(dir, 1780000001000);
assert.isFalse(next.error);
assert.deepStrictEqual(
next.files.flatMap((file) => file.records).map((record) => record.dedupeKey),
["opencode:msg-2"],
);
assert.isAbove((await NodeFSP.stat(NodePath.join(dir, "opencode.db-wal"))).size, 0);
} finally {
db.close();
}
});
});
Original file line number Diff line number Diff line change
@@ -1,11 +1,29 @@
// node:sqlite reads live OpenCode databases; Node fs walks legacy JSON history.
// @effect-diagnostics nodeBuiltinImport:off
/**
* Usage history for OpenCode, read from its SQLite databases and legacy JSON
* message store under each data directory.
*
* @module provider-opencode/server/usage
*/
import * as NodeFSP from "node:fs/promises";
import * as NodeOS from "node:os";
import * as NodePath from "node:path";
import * as NodeSqlite from "node:sqlite";
import * as NodeTimersPromises from "node:timers/promises";

import { totalTokens, type UsageRecord } from "@t3tools/provider-core/server/usage";
import { expandHomePath } from "@t3tools/provider-core/server/pathExpansion";
import {
totalTokens,
type ProviderUsageReader,
type UsageRecord,
} from "@t3tools/provider-core/server/usage";
import { HostProcessEnvironment } from "@t3tools/shared/hostProcess";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Path from "effect/Path";

import type { OpenCodeSettings } from "../settings.ts";

function object(value: unknown): Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value)
Expand Down Expand Up @@ -186,3 +204,55 @@ export async function readOpenCodeUsage(
}
return { files, missing: !found && !error, error };
}

/**
* The data directories to read: `OPENCODE_DATA_DIR` (comma-separated) or the
* XDG default, canonicalized so aliases count once.
*/
const resolveOpenCodeDataDirs = Effect.fn("resolveOpenCodeDataDirs")(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const environment = yield* HostProcessEnvironment;
const roots = environment["OPENCODE_DATA_DIR"]
?.split(",")
.map((value) => value.trim())
.filter(Boolean);
const dataHome = environment["XDG_DATA_HOME"]?.trim();
const defaults = [
path.join(
dataHome && path.isAbsolute(dataHome)
? dataHome
: path.join(NodeOS.homedir(), ".local", "share"),
"opencode",
),
];
const canonical = new Set<string>();
for (const root of roots?.length ? roots : defaults) {
const resolved = path.resolve(expandHomePath(root));
canonical.add(yield* fileSystem.realPath(resolved).pipe(Effect.orElseSucceed(() => resolved)));
}
return [...canonical];
});

export type OpenCodeUsageReaderEnv = FileSystem.FileSystem | Path.Path;

export const openCodeUsageReader: ProviderUsageReader<OpenCodeSettings, OpenCodeUsageReaderEnv> = {
kind: "scan",
provider: "opencode",
scan: Effect.fn("openCodeUsageReader.scan")(function* ({ windowStartMs }) {
const roots = yield* resolveOpenCodeDataDirs();
return yield* Effect.forEach(
roots,
(dir) =>
Effect.promise(() => readOpenCodeUsage(dir, windowStartMs)).pipe(
Effect.map((result) => ({
dir,
files: result.missing && !result.error ? null : result.files,
status: result.error ? ("partial" as const) : ("ok" as const),
...(result.error ? { message: "Some OpenCode history could not be read." } : {}),
})),
),
{ concurrency: "unbounded" },
);
}),
};
Loading