From b9e723cd899e494c06cce371a6829e57998231fd Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Mon, 31 Aug 2026 20:10:29 +0300 Subject: [PATCH 01/10] fix(usage): show opencode tokens and models on usage page Usage page only scanned claude/codex/grok transcripts, so opencode sessions showed zero tokens and no model breakdown even though the SDK exposes tokens/cost on every assistant message. The live adapter also never emitted thread.token-usage.updated, so the context window bar stayed empty. - Add opencode to UsageProviderKind and bump USAGE_CONTRACT_VERSION to 6 (compatible since 4) - Add OpenCode presentation for web and mobile - Emit live token usage from OpenCode adapter (message.updated, step-finish parts, session.updated, generic step events) and include usage in turn.completed - Scan opencode's SQLite DB (~/.local/share/opencode/opencode.db) for historical buckets with caching and dedup - Update docs, scan cache guard, and chart tests Verified with real opencode.db (4085 messages) and vp test run for usage, OpenCodeAdapter, and chart suites. Model: muse-spark-1.2-contributor-free Harness: opencode --- .../src/features/usage/usageProviders.ts | 4 +- .../src/provider/Layers/OpenCodeAdapter.ts | 233 +++++++++++++++++- apps/server/src/usage/UsageService.ts | 111 +++++++++ apps/server/src/usage/usageOpenCodeDb.ts | 140 +++++++++++ apps/server/src/usage/usageScanCache.ts | 3 +- apps/server/src/usage/usageTranscripts.ts | 76 ++++++ .../usage/UsageProviderChart.test.ts | 1 + .../src/components/usage/usageProviders.ts | 7 +- docs/user/usage.md | 11 +- packages/contracts/src/usage.ts | 17 +- packages/shared/src/usageMerge.test.ts | 2 +- 11 files changed, 585 insertions(+), 20 deletions(-) create mode 100644 apps/server/src/usage/usageOpenCodeDb.ts diff --git a/apps/mobile/src/features/usage/usageProviders.ts b/apps/mobile/src/features/usage/usageProviders.ts index 2576ac21fb07..fc18ea050c38 100644 --- a/apps/mobile/src/features/usage/usageProviders.ts +++ b/apps/mobile/src/features/usage/usageProviders.ts @@ -5,12 +5,13 @@ import { useAppearancePreferences } from "../settings/appearance/AppearancePrefe * Series and table order. The chart stacks providers from the bottom in this * order, so it also fixes which band sits on top of the bars. */ -export const PROVIDER_ORDER: readonly UsageProviderKind[] = ["codex", "claude", "grok"]; +export const PROVIDER_ORDER: readonly UsageProviderKind[] = ["codex", "claude", "grok", "opencode"]; export const PROVIDER_LABEL: Record = { claude: "Claude Code", codex: "Codex", grok: "Grok Build", + opencode: "OpenCode", }; /** @@ -23,5 +24,6 @@ export function useProviderColors(): Record { claude: "#d97757", codex: scheme === "dark" ? "#e6e6e6" : "#3c3c43", grok: scheme === "dark" ? "#a1a1aa" : "#52525b", + opencode: scheme === "dark" ? "#8a8a9a" : "#5a5a6a", }; } diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index dce8d3b6568c..e539144b7703 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -237,6 +237,9 @@ interface OpenCodeSessionContext { activeTurnId: TurnId | undefined; activeAgent: string | undefined; activeVariant: string | undefined; + lastTokens: OpenCodeTokens | null; + lastCost: number | null; + lastModel: string | null; /** * One-shot guard flipped by `stopOpenCodeContext` / `emitUnexpectedExit`. * The session lifecycle is owned by `sessionScope`; this Ref exists only @@ -528,6 +531,65 @@ function sessionErrorMessage(error: unknown): string { : "OpenCode session failed."; } +interface OpenCodeTokens { + readonly input: number; + readonly output: number; + readonly reasoning: number; + readonly cache: { readonly read: number; readonly write: number }; + readonly total?: number; +} + +function readOpenCodeTokens(value: unknown): OpenCodeTokens | null { + if (typeof value !== "object" || value === null) return null; + const record = value as Record; + const input = typeof record.input === "number" ? record.input : null; + const output = typeof record.output === "number" ? record.output : null; + if (input === null || output === null) return null; + const reasoning = typeof record.reasoning === "number" ? record.reasoning : 0; + const cacheRaw = record.cache; + let cacheRead = 0; + let cacheWrite = 0; + if (typeof cacheRaw === "object" && cacheRaw !== null) { + const cache = cacheRaw as Record; + cacheRead = typeof cache.read === "number" ? cache.read : 0; + cacheWrite = typeof cache.write === "number" ? cache.write : 0; + } + const total = typeof record.total === "number" ? record.total : undefined; + return { + input, + output, + reasoning, + cache: { read: cacheRead, write: cacheWrite }, + ...(total !== undefined ? { total } : {}), + }; +} + +function openCodeTokensToSnapshot(tokens: OpenCodeTokens): { + readonly usedTokens: number; + readonly inputTokens: number; + readonly cachedInputTokens: number; + readonly outputTokens: number; + readonly reasoningOutputTokens: number; +} { + const usedTokens = tokens.total ?? tokens.input + tokens.output; + const cachedInputTokens = tokens.cache.read; + const inputTokens = Math.max(0, tokens.input - cachedInputTokens - tokens.cache.write); + return { + usedTokens: Math.max(0, usedTokens), + inputTokens, + cachedInputTokens, + outputTokens: tokens.output, + reasoningOutputTokens: Math.min(tokens.output, tokens.reasoning), + }; +} + +function tryExtractOpenCodeTokensFromRecord( + record: Record, +): OpenCodeTokens | null { + if ("tokens" in record) return readOpenCodeTokens(record.tokens); + return null; +} + function updateProviderSession( context: OpenCodeSessionContext, patch: Partial, @@ -683,6 +745,66 @@ export function makeOpenCodeAdapter( }, ) => writeNativeEvent(threadId, event).pipe(Effect.catchCause(() => Effect.void)); + const emitOpenCodeTokenUsage = Effect.fn("emitOpenCodeTokenUsage")(function* ( + context: OpenCodeSessionContext, + tokens: OpenCodeTokens, + turnId: TurnId | undefined, + raw: unknown, + ) { + // Remember for turn.completed + context.lastTokens = tokens; + if ( + tokens.input + tokens.output + tokens.cache.read + tokens.cache.write === 0 && + tokens.total === undefined + ) { + return; + } + const snapshot = openCodeTokensToSnapshot(tokens); + if (snapshot.usedTokens <= 0) return; + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + raw, + })), + type: "thread.token-usage.updated", + payload: { + usage: snapshot, + }, + }); + }); + + const rememberOpenCodeCostAndModel = ( + context: OpenCodeSessionContext, + cost: unknown, + model: unknown, + providerId: unknown, + ) => { + if (typeof cost === "number" && Number.isFinite(cost)) context.lastCost = cost; + if (typeof model === "string" && model.trim().length > 0) { + const provider = + typeof providerId === "string" && providerId.trim().length > 0 + ? providerId.trim() + : "opencode"; + context.lastModel = `${provider}/${model.trim()}`; + } else if (typeof model === "object" && model !== null) { + const record = model as Record; + const id = + typeof record.id === "string" + ? record.id + : typeof record.modelID === "string" + ? record.modelID + : null; + const prov = + typeof record.providerID === "string" + ? record.providerID + : typeof record.provider === "string" + ? record.provider + : "opencode"; + if (id) context.lastModel = `${prov}/${id}`; + } + }; + const emitUnexpectedExit = Effect.fn("emitUnexpectedExit")(function* ( context: OpenCodeSessionContext, message: string, @@ -839,6 +961,26 @@ export function makeOpenCodeAdapter( }, }); } + // Session carries cumulative tokens/cost - emit live usage + { + const info = (event.properties as Record).info as + | Record + | undefined; + if (info) { + const tokens = readOpenCodeTokens(info.tokens); + if (tokens) { + rememberOpenCodeCostAndModel( + context, + info.cost, + info.model ?? info.modelID, + (info as Record).providerID, + ); + if (info.cost !== undefined) + context.lastCost = typeof info.cost === "number" ? info.cost : context.lastCost; + yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); + } + } + } break; } @@ -851,6 +993,17 @@ export function makeOpenCodeAdapter( } yield* emitAssistantTextDelta(context, part, turnId, event); } + const info = event.properties.info as unknown as Record; + const tokens = readOpenCodeTokens(info.tokens); + if (tokens) { + rememberOpenCodeCostAndModel( + context, + info.cost, + info.modelID ?? info.model, + info.providerID, + ); + yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); + } } break; } @@ -914,6 +1067,21 @@ export function makeOpenCodeAdapter( yield* emitAssistantTextDelta(context, part, turnId, event); } + // Step-finish and other token-carrying parts + { + const partRecord = part as unknown as Record; + const tokens = readOpenCodeTokens(partRecord.tokens); + if (tokens) { + rememberOpenCodeCostAndModel( + context, + partRecord.cost, + (partRecord.modelID as string | undefined) ?? (partRecord.model as unknown), + partRecord.providerID as string | undefined, + ); + yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); + } + } + if (part.type === "tool") { const itemType = toToolLifecycleItemType(part.tool); const title = @@ -1076,6 +1244,17 @@ export function makeOpenCodeAdapter( if (event.properties.status.type === "idle" && turnId) { context.activeTurnId = undefined; yield* updateProviderSession(context, { status: "ready" }, { clearActiveTurnId: true }); + const completedPayload: Record = { state: "completed" }; + if (context.lastTokens) { + completedPayload.usage = context.lastTokens; + completedPayload.modelUsage = context.lastModel + ? { [context.lastModel]: context.lastTokens } + : undefined; + } + if (context.lastCost !== null) completedPayload.totalCostUsd = context.lastCost; + // Reset for next turn - keep model but clear tokens to avoid leaking to next turn's idle if no tokens + context.lastTokens = null; + context.lastCost = null; yield* emit({ ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -1083,9 +1262,7 @@ export function makeOpenCodeAdapter( raw: event, })), type: "turn.completed", - payload: { - state: "completed", - }, + payload: completedPayload as never, }); } break; @@ -1132,8 +1309,53 @@ export function makeOpenCodeAdapter( break; } - default: + default: { + // Generic token extraction for events like `session.next.step.ended`, `session.diff`, etc. + // These carry `tokens`/`cost`/`model` at top-level properties or nested. + const props = (event as Record).properties as + | Record + | undefined; + if (props) { + const tokens = + readOpenCodeTokens(props.tokens) ?? + readOpenCodeTokens((props.part as Record | undefined)?.tokens) ?? + readOpenCodeTokens((props.info as Record | undefined)?.tokens) ?? + null; + if (tokens) { + const cost = + (props.cost as number | undefined) ?? + (props.info as Record | undefined)?.cost ?? + (props.part as Record | undefined)?.cost; + const model = + (props.model as unknown) ?? + (props.info as Record | undefined)?.model ?? + (props.info as Record | undefined)?.modelID ?? + (props.part as Record | undefined)?.modelID; + const providerId = + (props.providerID as string | undefined) ?? + (props.info as Record | undefined)?.providerID ?? + (props.part as Record | undefined)?.providerID; + rememberOpenCodeCostAndModel(context, cost, model, providerId); + yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); + } else { + // Also try nested usage like `event.properties.tokens` already covered, try `event.properties.data.tokens` + const data = props.data as Record | undefined; + if (data) { + const nestedTokens = readOpenCodeTokens(data.tokens); + if (nestedTokens) { + rememberOpenCodeCostAndModel( + context, + data.cost, + data.modelID ?? data.model, + data.providerID, + ); + yield* emitOpenCodeTokenUsage(context, nestedTokens, turnId, event); + } + } + } + } break; + } } }); @@ -1402,6 +1624,9 @@ export function makeOpenCodeAdapter( activeTurnId: undefined, activeAgent: undefined, activeVariant: undefined, + lastTokens: null, + lastCost: null, + lastModel: null, stopped: yield* Ref.make(false), sessionScope: started.sessionScope, }; diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index 224662e9dca7..5f3101010c85 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -1,3 +1,4 @@ +// @effect-diagnostics nodeBuiltinImport:off /** * UsageService - scans provider transcripts and returns priced usage buckets. * @@ -12,6 +13,7 @@ * @module UsageService */ import * as NodeOS from "node:os"; +import * as NodePath from "node:path"; import { USAGE_CONTRACT_VERSION, @@ -54,6 +56,12 @@ import { type ScanCache, } from "./usageScanCache.ts"; import type { UsageRecord } from "./usageTranscripts.ts"; +import { + readOpenCodeDbRecords, + readOpenCodeDbVolumeId, + resolveOpenCodeDbPath, + statOpenCodeDb, +} from "./usageOpenCodeDb.ts"; const LITELLM_RATES_URL = "https://raw.githubusercontent.com/BerriAI/litellm/main/model_prices_and_context_window.json"; @@ -305,6 +313,30 @@ export const make = Effect.gen(function* () { return records; }); + /** Parses the OpenCode SQLite DB, reusing cached result when unchanged. */ + const readOpenCodeDbRecordsCached = ( + dbPath: string, + size: number, + mtimeMs: number, + ): Effect.Effect => + Effect.gen(function* () { + const cached = fileCache.get(dbPath); + if ( + cached && + cached.size === size && + cached.mtimeMs === mtimeMs && + cached.provider === "opencode" + ) { + return cached.records; + } + const parsed = yield* Effect.promise(() => readOpenCodeDbRecords(dbPath)); + if (parsed === null) return []; + const records = dedupeWithinFile(parsed); + fileCache.set(dbPath, { size, mtimeMs, provider: "opencode" as const, records }); + cacheDirty = true; + return records; + }); + const readSummary = Effect.fn("UsageService.readSummary")(function* (input: UsageSummaryInput) { if (input.sinceDay > input.untilDay) { return yield* new UsageReadError({ @@ -425,6 +457,85 @@ export const make = Effect.gen(function* () { }); } + // OpenCode: SQLite DB instead of transcript files + { + const dbPath = yield* Effect.promise(() => resolveOpenCodeDbPath()); + if (dbPath === null) { + const fallbackPath = NodePath.join( + NodeOS.homedir(), + ".local", + "share", + "opencode", + "opencode.db", + ); + const volumeId = yield* Effect.promise(() => readOpenCodeDbVolumeId(fallbackPath)); + sources.push({ + fingerprint: { + hostId, + provider: "opencode" as const, + resolvedHomePath: fallbackPath, + volumeId, + }, + status: "missing", + scannedFiles: 0, + skippedFiles: 0, + malformedRecords: 0, + distinctSessions: 0, + message: "No OpenCode database on this environment.", + }); + } else { + const volumeId = yield* Effect.promise(() => readOpenCodeDbVolumeId(dbPath)); + const stat = yield* Effect.promise(() => statOpenCodeDb(dbPath)); + if (stat === null) { + sources.push({ + fingerprint: { + hostId, + provider: "opencode" as const, + resolvedHomePath: dbPath, + volumeId, + }, + status: "failed", + scannedFiles: 0, + skippedFiles: 0, + malformedRecords: 0, + distinctSessions: 0, + message: "OpenCode database could not be read.", + }); + } else { + walkedRoots.push(NodePath.dirname(dbPath)); + livePaths.add(dbPath); + const records = yield* readOpenCodeDbRecordsCached(dbPath, stat.size, stat.mtimeMs); + const sessionIds = new Set(); + let scannedFiles = 0; + let skippedFiles = 0; + if (records.length === 0) { + skippedFiles = 1; + } else { + scannedFiles = 1; + } + for (const record of records) { + if (aggregator.add(record) && record.sessionId.length > 0) { + sessionIds.add(record.sessionId); + } + } + sources.push({ + fingerprint: { + hostId, + provider: "opencode" as const, + resolvedHomePath: dbPath, + volumeId, + }, + status: "ok", + scannedFiles, + skippedFiles, + malformedRecords: 0, + distinctSessions: sessionIds.size, + message: null, + }); + } + } + } + const pruned = pruneScanCache(fileCache, { livePaths, walkedRoots, diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts new file mode 100644 index 000000000000..2d0b552e17d7 --- /dev/null +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -0,0 +1,140 @@ +// @effect-diagnostics nodeBuiltinImport:off +/** + * SQLite reader for OpenCode's `opencode.db`. + * + * OpenCode stores all session history in `~/.local/share/opencode/opencode.db` + * (Linux) / `~/Library/Application Support/opencode/opencode.db` (macOS). The + * `message` table holds one row per user/assistant message with JSON `data` + * containing `tokens`, `cost`, `modelID`, etc. We scan that table directly + * rather than the file-based transcript approach used for Claude/Codex/Grok. + * + * @module usageOpenCodeDb + */ +import * as NodeFS from "node:fs"; +import * as NodeOS from "node:os"; +import * as NodePath from "node:path"; + +import type { UsageRecord } from "./usageTranscripts.ts"; +import { parseOpenCodeMessage } from "./usageTranscripts.ts"; + +/** + * Resolves the OpenCode database path. + * + * Checks in order: + * 1. `$OPENCODE_DATA_DIR/opencode.db` if set + * 2. `$XDG_DATA_HOME/opencode/opencode.db` if `XDG_DATA_HOME` set + * 3. `~/.local/share/opencode/opencode.db` (Linux) + * 4. `~/Library/Application Support/opencode/opencode.db` (macOS fallback) + * + * Returns the first path that exists, or `null` if none exist. + */ +export async function resolveOpenCodeDbPath(): Promise { + const candidates: string[] = []; + + const dataDirEnv = process.env.OPENCODE_DATA_DIR?.trim(); + if (dataDirEnv && dataDirEnv.length > 0) { + candidates.push(NodePath.join(dataDirEnv, "opencode.db")); + } + + const xdgDataHome = process.env.XDG_DATA_HOME?.trim(); + if (xdgDataHome && xdgDataHome.length > 0) { + candidates.push(NodePath.join(xdgDataHome, "opencode", "opencode.db")); + } + + candidates.push(NodePath.join(NodeOS.homedir(), ".local", "share", "opencode", "opencode.db")); + candidates.push( + NodePath.join(NodeOS.homedir(), "Library", "Application Support", "opencode", "opencode.db"), + ); + // Also check XDG on macOS via ~/.local/share + // And legacy ~/.opencode + candidates.push(NodePath.join(NodeOS.homedir(), ".opencode", "opencode.db")); + + for (const candidate of candidates) { + try { + await NodeFS.promises.access(candidate, NodeFS.constants.R_OK); + return candidate; + } catch { + // not accessible, try next + } + } + return null; +} + +/** + * Reads OpenCode message records from the SQLite database. + * + * Opens the database read-only so the live OpenCode server can continue + * writing with WAL. Returns `null` if the file cannot be opened or the + * schema is unexpected; an empty array if the file exists but has no + * relevant rows. Caller is responsible for caching by `(size, mtime)`. + */ +export async function readOpenCodeDbRecords( + dbPath: string, +): Promise { + let DatabaseSync: typeof import("node:sqlite").DatabaseSync; + try { + const sqlite = await import("node:sqlite"); + DatabaseSync = sqlite.DatabaseSync; + } catch { + return null; + } + + let db: InstanceType | null = null; + try { + db = new DatabaseSync(dbPath, { readOnly: true, enableForeignKeyConstraints: false }); + } catch { + return null; + } + + try { + // Verify table exists + const tableCheck = db + .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name='message'") + .get() as { name: string } | undefined; + if (!tableCheck) return []; + + const stmt = db.prepare("SELECT id, session_id, data FROM message"); + const rows = stmt.all() as Array<{ id: string; session_id: string; data: string }>; + const records: UsageRecord[] = []; + for (const row of rows) { + const record = parseOpenCodeMessage(row.data, row.id, row.session_id); + if (record !== null) records.push(record); + } + return records; + } catch { + return null; + } finally { + try { + db?.close(); + } catch {} + } +} + +/** + * Gets file stats for caching (size, mtime) for the OpenCode DB. + */ +export async function statOpenCodeDb( + dbPath: string, +): Promise<{ size: number; mtimeMs: number } | null> { + try { + const stats = await NodeFS.promises.stat(dbPath); + return { size: stats.size, mtimeMs: stats.mtimeMs }; + } catch { + return null; + } +} + +/** + * Reads the directory `device:inode` for fingerprinting the OpenCode DB source. + */ +export async function readOpenCodeDbVolumeId(dbPath: string): Promise { + try { + const dir = NodePath.dirname(dbPath); + const stats = await NodeFS.promises.stat(dir); + const dev = (stats as unknown as { dev: number }).dev ?? 0; + const ino = (stats as unknown as { ino: number }).ino ?? 0; + return `${dev}:${ino}`; + } catch { + return ""; + } +} diff --git a/apps/server/src/usage/usageScanCache.ts b/apps/server/src/usage/usageScanCache.ts index 02daf5ebbd70..dcae53d8bef0 100644 --- a/apps/server/src/usage/usageScanCache.ts +++ b/apps/server/src/usage/usageScanCache.ts @@ -134,7 +134,8 @@ export function decodeScanCache(document: unknown): ScanCache { if (typeof raw !== "object" || raw === null) continue; const entry = raw as Partial; if (typeof entry.s !== "number" || typeof entry.m !== "number") continue; - if (entry.p !== "claude" && entry.p !== "codex" && entry.p !== "grok") continue; + if (entry.p !== "claude" && entry.p !== "codex" && entry.p !== "grok" && entry.p !== "opencode") + continue; if (!isRecordArray(entry.r)) continue; const provider: UsageProviderKind = entry.p; diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index 2aea60709666..d4ebb2d34285 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -70,6 +70,7 @@ export function totalTokens(totals: UsageTokenTotals): number { export function mightCarryUsage(line: string, provider: UsageProviderKind): boolean { if (provider === "claude") return line.includes('"usage"'); if (provider === "grok") return line.includes('"turn_completed"'); + if (provider === "opencode") return line.includes('"tokens"'); return line.includes('"token_count"'); } @@ -485,4 +486,79 @@ export function parseGrokLine(line: string): readonly UsageRecord[] { return results; } +/* -------------------------------------------------------------------------- */ +/* OpenCode */ +/* -------------------------------------------------------------------------- */ + +/** + * Parses one message record from the OpenCode SQLite `message` table. + * + * Expects the JSON stored in `message.data` and the `message.id` / + * `message.session_id` columns. Only assistant messages with token counts are + * emitted; user messages and assistants without usage are ignored. + */ +export function parseOpenCodeMessage( + data: string, + messageId: string, + sessionId: string, +): UsageRecord | null { + let parsed: unknown; + try { + parsed = JSON.parse(data); + } catch { + return null; + } + if (typeof parsed !== "object" || parsed === null) return null; + const record = parsed as Record; + if (record.role !== "assistant") return null; + + const tokensRaw = record.tokens; + if (typeof tokensRaw !== "object" || tokensRaw === null) return null; + const tokens = tokensRaw as Record; + const cacheRaw = tokens.cache as Record | undefined; + + const inputTokens = int(tokens.input); + const outputTokens = int(tokens.output); + const reasoningTokens = int(tokens.reasoning); + const cachedRead = int(cacheRaw?.read); + const cacheWrite = int(cacheRaw?.write); + + const totals: UsageTokenTotals = { + uncachedInputTokens: Math.max(0, inputTokens - cachedRead - cacheWrite), + cachedInputTokens: cachedRead, + cacheCreationTokens: cacheWrite, + outputTokens, + reasoningTokens: Math.min(outputTokens, reasoningTokens), + }; + + if (totalTokens(totals) === 0) return null; + + // Prefer completed timestamp for bucketing, fall back to created. + let timestampMs: number | null = null; + const timeRaw = record.time as Record | undefined; + if (timeRaw) { + if (typeof timeRaw.completed === "number" && Number.isFinite(timeRaw.completed)) { + timestampMs = timeRaw.completed; + } else if (typeof timeRaw.created === "number" && Number.isFinite(timeRaw.created)) { + timestampMs = timeRaw.created; + } + } + if (timestampMs === null) return null; + + const modelId = typeof record.modelID === "string" ? record.modelID : ""; + const providerId = typeof record.providerID === "string" ? record.providerID : "opencode"; + const model = modelId.length > 0 ? `${providerId}/${modelId}` : providerId; + + const cost = record.cost; + return { + provider: "opencode", + timestampMs, + model, + sessionId, + totals, + reportedCostUsd: typeof cost === "number" && Number.isFinite(cost) ? cost : null, + dedupeKey: messageId.length > 0 ? messageId : null, + }; +} + export { EMPTY_TOTALS }; diff --git a/apps/web/src/components/usage/UsageProviderChart.test.ts b/apps/web/src/components/usage/UsageProviderChart.test.ts index a4114cfdfb57..9b61e918a523 100644 --- a/apps/web/src/components/usage/UsageProviderChart.test.ts +++ b/apps/web/src/components/usage/UsageProviderChart.test.ts @@ -87,6 +87,7 @@ describe("buildDayColumns", () => { { provider: "codex", value: 10 }, { provider: "claude", value: 20 }, { provider: "grok", value: 0 }, + { provider: "opencode", value: 0 }, ]); }); diff --git a/apps/web/src/components/usage/usageProviders.ts b/apps/web/src/components/usage/usageProviders.ts index efad95e531ad..20b4cfff304a 100644 --- a/apps/web/src/components/usage/usageProviders.ts +++ b/apps/web/src/components/usage/usageProviders.ts @@ -1,6 +1,6 @@ import type { UsageProviderKind } from "@t3tools/contracts"; -import { ClaudeAI, GrokIcon, type Icon, OpenAI } from "../Icons"; +import { ClaudeAI, GrokIcon, type Icon, OpenAI, OpenCodeIcon } from "../Icons"; type UsageProviderPresentation = { readonly label: string; @@ -30,6 +30,11 @@ export const PROVIDER_PRESENTATION = { color: "color-mix(in oklab, var(--contrast-foreground) 72%, var(--background))", mark: GrokIcon, }, + opencode: { + label: "OpenCode", + color: "color-mix(in oklab, var(--contrast-foreground) 58%, var(--background))", + mark: OpenCodeIcon, + }, } satisfies Record; /** Stable provider reading order across charts, summaries, tables, and hover rows. */ diff --git a/docs/user/usage.md b/docs/user/usage.md index ff38c730c1cd..00dbeb427452 100644 --- a/docs/user/usage.md +++ b/docs/user/usage.md @@ -1,13 +1,16 @@ # Review usage -The Usage page combines Codex, Claude Code, and Grok Build activity from your connected -environments. It reads the providers' local session history and shows API-equivalent token cost, -processed tokens, cache savings, provider shares, and model breakdowns. Subscription billing is -separate from the raw token cost shown here. +The Usage page combines Codex, Claude Code, Grok Build, and OpenCode activity from your +connected environments. It reads the providers' local session history and shows API-equivalent +token cost, processed tokens, cache savings, provider shares, and model breakdowns. Subscription +billing is separate from the raw token cost shown here. Grok Build totals come from persisted session updates. Interactive turns that never wrote a completed-turn record will not appear. +OpenCode totals come from its local SQLite database (`~/.local/share/opencode/opencode.db` on +Linux). Like the other providers, every turn that wrote a message with token counts is included. + Use **Past 24h** for an hourly chart covering the exact rolling 24-hour period. The **7 days**, **30 days**, and **90 days** ranges use daily resolution. Cost and token toggles update both the headline and chart, and refreshing rescans every connected environment. diff --git a/packages/contracts/src/usage.ts b/packages/contracts/src/usage.ts index 8c099ddb33aa..e8426fbc7a34 100644 --- a/packages/contracts/src/usage.ts +++ b/packages/contracts/src/usage.ts @@ -3,9 +3,10 @@ * * Each environment scans the provider CLIs' own on-disk session transcripts * (`~/.claude/projects/**\/*.jsonl`, `~/.codex/sessions/**\/*.jsonl`, - * `~/.grok/sessions/**\/updates.jsonl`) rather than relying on T3 Code's own - * orchestration projections, so usage stays complete even for turns that were - * never driven through T3 Code. This mirrors the approach `ccusage` takes. + * `~/.grok/sessions/**\/updates.jsonl`, `~/.local/share/opencode/opencode.db`) + * rather than relying on T3 Code's own orchestration projections, so usage + * stays complete even for turns that were never driven through T3 Code. This + * mirrors the approach `ccusage` takes. * * Environments return pre-aggregated `(day, hourStart?, provider, model)` * buckets. Raw transcript records never cross the wire. @@ -21,18 +22,18 @@ import { NonNegativeInt, TrimmedNonEmptyString } from "./baseSchemas.ts"; * client renders partial coverage when an environment reports an older version * rather than failing the whole page. */ -export const USAGE_CONTRACT_VERSION = 5 as const; +export const USAGE_CONTRACT_VERSION = 6 as const; /** * Oldest {@link UsageSummary} version a current client will still merge. * - * v5 only adds `grok` to {@link UsageProviderKind}; v4 Claude/Codex buckets - * remain valid, so mixed-version environments keep those totals instead of - * treating every older server as stale. + * v6 only adds `opencode` to {@link UsageProviderKind}; v4 Claude/Codex and + * v5 Grok buckets remain valid, so mixed-version environments keep those totals + * instead of treating every older server as stale. */ export const USAGE_MERGE_COMPATIBLE_SINCE = 4 as const; -export const UsageProviderKind = Schema.Literals(["claude", "codex", "grok"]); +export const UsageProviderKind = Schema.Literals(["claude", "codex", "grok", "opencode"]); export type UsageProviderKind = typeof UsageProviderKind.Type; /** diff --git a/packages/shared/src/usageMerge.test.ts b/packages/shared/src/usageMerge.test.ts index 6c706395c6ff..19e750356685 100644 --- a/packages/shared/src/usageMerge.test.ts +++ b/packages/shared/src/usageMerge.test.ts @@ -158,7 +158,7 @@ describe("mergeUsage", () => { summary( [bucket()], [{ provider: "claude", hostId: "linux", homePath: "/b" }], - USAGE_CONTRACT_VERSION - 2, + USAGE_CONTRACT_VERSION - 3, ), ), ], From 8aeb626759d95c77897819f3d166fd60278a98bf Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Mon, 31 Aug 2026 20:43:03 +0300 Subject: [PATCH 02/10] fix(usage): address review findings for opencode integration - Fix cache double subtraction: opencode input is exclusive of cache, so uncached should be input directly (not input - cache) - Fix session.updated overwriting: session tokens are cumulative, do not emit as thread.token-usage or store for turn.completed - Fix WAL cache staleness: stat now includes opencode.db-wal size/mtime, busting cache on WAL writes - Add V2 session_message table scan alongside legacy message table - Fix null vs empty: preserve readOpenCodeDbRecords null failure, reporting source as failed instead of ok - Thread HostProcessEnvironment into resolveOpenCodeDbPath instead of module-global process.env - Use Path service (path.join/dirname) instead of node:path import, removing suppression comment - Fix as never cast on turn.completed payload, spreading only defined keys with proper typing --- .../src/provider/Layers/OpenCodeAdapter.ts | 55 ++++++------- apps/server/src/usage/UsageService.ts | 77 +++++++++++-------- apps/server/src/usage/usageOpenCodeDb.ts | 68 ++++++++++++---- apps/server/src/usage/usageTranscripts.ts | 3 +- 4 files changed, 123 insertions(+), 80 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index e539144b7703..0db02728e314 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -571,9 +571,11 @@ function openCodeTokensToSnapshot(tokens: OpenCodeTokens): { readonly outputTokens: number; readonly reasoningOutputTokens: number; } { - const usedTokens = tokens.total ?? tokens.input + tokens.output; + // OpenCode reports `input` exclusive of cache (unlike Codex/Grok). Don't subtract. const cachedInputTokens = tokens.cache.read; - const inputTokens = Math.max(0, tokens.input - cachedInputTokens - tokens.cache.write); + const inputTokens = tokens.input; + const usedTokens = + tokens.total ?? tokens.input + tokens.output + tokens.cache.read + tokens.cache.write; return { usedTokens: Math.max(0, usedTokens), inputTokens, @@ -961,26 +963,11 @@ export function makeOpenCodeAdapter( }, }); } - // Session carries cumulative tokens/cost - emit live usage - { - const info = (event.properties as Record).info as - | Record - | undefined; - if (info) { - const tokens = readOpenCodeTokens(info.tokens); - if (tokens) { - rememberOpenCodeCostAndModel( - context, - info.cost, - info.model ?? info.modelID, - (info as Record).providerID, - ); - if (info.cost !== undefined) - context.lastCost = typeof info.cost === "number" ? info.cost : context.lastCost; - yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); - } - } - } + // Note: `info.tokens`/`cost` on session.updated are cumulative for the whole + // session, not per-turn. Emitting them as `thread.token-usage.updated` would + // make the context bar track lifetime usage and cause turn.completed to report + // the entire session as one turn, so they are intentionally ignored here. + // Per-turn usage comes from message.updated / step-finish parts. break; } @@ -1244,17 +1231,12 @@ export function makeOpenCodeAdapter( if (event.properties.status.type === "idle" && turnId) { context.activeTurnId = undefined; yield* updateProviderSession(context, { status: "ready" }, { clearActiveTurnId: true }); - const completedPayload: Record = { state: "completed" }; - if (context.lastTokens) { - completedPayload.usage = context.lastTokens; - completedPayload.modelUsage = context.lastModel - ? { [context.lastModel]: context.lastTokens } - : undefined; - } - if (context.lastCost !== null) completedPayload.totalCostUsd = context.lastCost; - // Reset for next turn - keep model but clear tokens to avoid leaking to next turn's idle if no tokens + const usage = context.lastTokens; + const cost = context.lastCost; + const model = context.lastModel; context.lastTokens = null; context.lastCost = null; + context.lastModel = null; yield* emit({ ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -1262,7 +1244,16 @@ export function makeOpenCodeAdapter( raw: event, })), type: "turn.completed", - payload: completedPayload as never, + payload: { + state: "completed", + ...(usage + ? { + usage, + ...(model ? { modelUsage: { [model]: usage } } : {}), + } + : {}), + ...(cost !== null ? { totalCostUsd: cost } : {}), + }, }); } break; diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index 5f3101010c85..c454c42ddc24 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -1,4 +1,3 @@ -// @effect-diagnostics nodeBuiltinImport:off /** * UsageService - scans provider transcripts and returns priced usage buckets. * @@ -13,7 +12,6 @@ * @module UsageService */ import * as NodeOS from "node:os"; -import * as NodePath from "node:path"; import { USAGE_CONTRACT_VERSION, @@ -318,7 +316,7 @@ export const make = Effect.gen(function* () { dbPath: string, size: number, mtimeMs: number, - ): Effect.Effect => + ): Effect.Effect => Effect.gen(function* () { const cached = fileCache.get(dbPath); if ( @@ -330,7 +328,7 @@ export const make = Effect.gen(function* () { return cached.records; } const parsed = yield* Effect.promise(() => readOpenCodeDbRecords(dbPath)); - if (parsed === null) return []; + if (parsed === null) return null; const records = dedupeWithinFile(parsed); fileCache.set(dbPath, { size, mtimeMs, provider: "opencode" as const, records }); cacheDirty = true; @@ -459,9 +457,9 @@ export const make = Effect.gen(function* () { // OpenCode: SQLite DB instead of transcript files { - const dbPath = yield* Effect.promise(() => resolveOpenCodeDbPath()); + const dbPath = yield* Effect.promise(() => resolveOpenCodeDbPath(hostEnvironment)); if (dbPath === null) { - const fallbackPath = NodePath.join( + const fallbackPath = path.join( NodeOS.homedir(), ".local", "share", @@ -502,36 +500,53 @@ export const make = Effect.gen(function* () { message: "OpenCode database could not be read.", }); } else { - walkedRoots.push(NodePath.dirname(dbPath)); + walkedRoots.push(path.dirname(dbPath)); livePaths.add(dbPath); const records = yield* readOpenCodeDbRecordsCached(dbPath, stat.size, stat.mtimeMs); - const sessionIds = new Set(); - let scannedFiles = 0; - let skippedFiles = 0; - if (records.length === 0) { - skippedFiles = 1; + if (records === null) { + sources.push({ + fingerprint: { + hostId, + provider: "opencode" as const, + resolvedHomePath: dbPath, + volumeId, + }, + status: "failed", + scannedFiles: 0, + skippedFiles: 0, + malformedRecords: 0, + distinctSessions: 0, + message: "OpenCode database could not be read.", + }); } else { - scannedFiles = 1; - } - for (const record of records) { - if (aggregator.add(record) && record.sessionId.length > 0) { - sessionIds.add(record.sessionId); + const sessionIds = new Set(); + let scannedFiles = 0; + let skippedFiles = 0; + if (records.length === 0) { + skippedFiles = 1; + } else { + scannedFiles = 1; + } + for (const record of records) { + if (aggregator.add(record) && record.sessionId.length > 0) { + sessionIds.add(record.sessionId); + } } + sources.push({ + fingerprint: { + hostId, + provider: "opencode" as const, + resolvedHomePath: dbPath, + volumeId, + }, + status: "ok", + scannedFiles, + skippedFiles, + malformedRecords: 0, + distinctSessions: sessionIds.size, + message: null, + }); } - sources.push({ - fingerprint: { - hostId, - provider: "opencode" as const, - resolvedHomePath: dbPath, - volumeId, - }, - status: "ok", - scannedFiles, - skippedFiles, - malformedRecords: 0, - distinctSessions: sessionIds.size, - message: null, - }); } } } diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index 2d0b552e17d7..ee079040a99c 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -28,15 +28,17 @@ import { parseOpenCodeMessage } from "./usageTranscripts.ts"; * * Returns the first path that exists, or `null` if none exist. */ -export async function resolveOpenCodeDbPath(): Promise { +export async function resolveOpenCodeDbPath( + env: NodeJS.ProcessEnv = process.env, +): Promise { const candidates: string[] = []; - const dataDirEnv = process.env.OPENCODE_DATA_DIR?.trim(); + const dataDirEnv = env.OPENCODE_DATA_DIR?.trim(); if (dataDirEnv && dataDirEnv.length > 0) { candidates.push(NodePath.join(dataDirEnv, "opencode.db")); } - const xdgDataHome = process.env.XDG_DATA_HOME?.trim(); + const xdgDataHome = env.XDG_DATA_HOME?.trim(); if (xdgDataHome && xdgDataHome.length > 0) { candidates.push(NodePath.join(xdgDataHome, "opencode", "opencode.db")); } @@ -67,6 +69,9 @@ export async function resolveOpenCodeDbPath(): Promise { * writing with WAL. Returns `null` if the file cannot be opened or the * schema is unexpected; an empty array if the file exists but has no * relevant rows. Caller is responsible for caching by `(size, mtime)`. + * + * Reads both the legacy `message` table and the current `session_message` + * table (OpenCode v2), unioning the results. */ export async function readOpenCodeDbRecords( dbPath: string, @@ -87,18 +92,33 @@ export async function readOpenCodeDbRecords( } try { - // Verify table exists - const tableCheck = db - .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name='message'") - .get() as { name: string } | undefined; - if (!tableCheck) return []; - - const stmt = db.prepare("SELECT id, session_id, data FROM message"); - const rows = stmt.all() as Array<{ id: string; session_id: string; data: string }>; const records: UsageRecord[] = []; - for (const row of rows) { - const record = parseOpenCodeMessage(row.data, row.id, row.session_id); - if (record !== null) records.push(record); + const tables: string[] = []; + try { + const rows = db + .prepare( + "SELECT name FROM sqlite_master WHERE type='table' AND name IN ('message','session_message')", + ) + .all() as Array<{ name: string }>; + for (const row of rows) tables.push(row.name); + } catch { + return null; + } + if (tables.length === 0) return []; + + for (const table of tables) { + try { + // `message` has `session_id`, `session_message` has `session_id` as well but different time columns; data is same shape. + const stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); + const rows = stmt.all() as Array<{ id: string; session_id: string; data: string }>; + for (const row of rows) { + const record = parseOpenCodeMessage(row.data, row.id, row.session_id); + if (record !== null) records.push(record); + } + } catch { + // One table failing shouldn't hide the other; treat as empty for that table. + continue; + } } return records; } catch { @@ -111,14 +131,30 @@ export async function readOpenCodeDbRecords( } /** - * Gets file stats for caching (size, mtime) for the OpenCode DB. + * Gets file stats for caching (size, mtime) for the OpenCode DB, including WAL. + * + * When OpenCode runs in WAL mode, committed messages live in `opencode.db-wal` + * while the main file's mtime/size can stay stale until a checkpoint. Including + * the WAL's fingerprint ensures a refresh sees new rows immediately. */ export async function statOpenCodeDb( dbPath: string, ): Promise<{ size: number; mtimeMs: number } | null> { try { const stats = await NodeFS.promises.stat(dbPath); - return { size: stats.size, mtimeMs: stats.mtimeMs }; + let size = stats.size; + let mtimeMs = stats.mtimeMs; + // Include WAL file if present - its mtime/size changes on every committed write. + const walPath = `${dbPath}-wal`; + try { + const walStats = await NodeFS.promises.stat(walPath); + size += walStats.size; + mtimeMs = Math.max(mtimeMs, walStats.mtimeMs); + } catch { + // No WAL file, that's fine. + } + // Also include shm for completeness, though its mtime is less meaningful. + return { size, mtimeMs }; } catch { return null; } diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index d4ebb2d34285..b9220b987448 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -524,7 +524,8 @@ export function parseOpenCodeMessage( const cacheWrite = int(cacheRaw?.write); const totals: UsageTokenTotals = { - uncachedInputTokens: Math.max(0, inputTokens - cachedRead - cacheWrite), + // OpenCode reports `input` exclusive of cache (unlike Codex/Grok which are inclusive). + uncachedInputTokens: inputTokens, cachedInputTokens: cachedRead, cacheCreationTokens: cacheWrite, outputTokens, From 52f99f694710a24fcbbc9899f2c7e9aee6b28b4f Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Mon, 31 Aug 2026 20:57:51 +0300 Subject: [PATCH 03/10] fix(usage): address second review round for opencode - Fix V2 session_message parsing: handle role in type column, nested model.id/providerID, and column times - Fix blocking DB read: use iterate() with setImmediate yields every 100 rows instead of all() - Fix diagnostic suppression comment to include rationale - Add focused tests for parseOpenCodeMessage covering all easy-to-regress behaviors Addresses Macroscope high blocking and Cursor V2 findings --- apps/server/src/usage/usageOpenCodeDb.ts | 51 ++++++- .../server/src/usage/usageTranscripts.test.ts | 137 ++++++++++++++++++ apps/server/src/usage/usageTranscripts.ts | 28 +++- 3 files changed, 206 insertions(+), 10 deletions(-) diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index ee079040a99c..ab0357504b15 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -1,4 +1,4 @@ -// @effect-diagnostics nodeBuiltinImport:off +// @effect-diagnostics nodeBuiltinImport:off - `node:sqlite` has no Effect wrapper, so this reader stays on Node built-ins; isolated here so the rest of the usage code keeps using Effect's FileSystem/Path. /** * SQLite reader for OpenCode's `opencode.db`. * @@ -108,12 +108,51 @@ export async function readOpenCodeDbRecords( for (const table of tables) { try { - // `message` has `session_id`, `session_message` has `session_id` as well but different time columns; data is same shape. - const stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); - const rows = stmt.all() as Array<{ id: string; session_id: string; data: string }>; - for (const row of rows) { - const record = parseOpenCodeMessage(row.data, row.id, row.session_id); + // Use iterate() instead of all() to avoid loading all rows at once and to allow yielding. + let rows: Iterable>; + let isSessionMessage = table === "session_message"; + try { + if (isSessionMessage) { + const stmt = db.prepare( + `SELECT id, session_id, type, time_created, time_updated, data FROM "${table}"`, + ); + rows = stmt.iterate() as Iterable>; + } else { + const stmt = db.prepare( + `SELECT id, session_id, time_created, time_updated, data FROM "${table}"`, + ); + rows = stmt.iterate() as Iterable>; + } + } catch { + // Fallback for older schemas without time columns + const stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); + rows = stmt.iterate() as Iterable>; + isSessionMessage = false; + } + let count = 0; + for (const raw of rows) { + const row = raw as { + id: string; + session_id: string; + data: string; + type?: string; + time_created?: number; + time_updated?: number; + }; + const record = parseOpenCodeMessage( + row.data, + row.id, + row.session_id, + isSessionMessage ? row.type : undefined, + row.time_created, + row.time_updated, + ); if (record !== null) records.push(record); + // Yield to event loop every 100 rows to avoid blocking cold scans. + if (++count % 100 === 0) { + // eslint-disable-next-line no-await-in-loop + await new Promise((resolve) => setImmediate(resolve)); + } } } catch { // One table failing shouldn't hide the other; treat as empty for that table. diff --git a/apps/server/src/usage/usageTranscripts.test.ts b/apps/server/src/usage/usageTranscripts.test.ts index b09db613ed85..409ff8a4aea9 100644 --- a/apps/server/src/usage/usageTranscripts.test.ts +++ b/apps/server/src/usage/usageTranscripts.test.ts @@ -6,6 +6,7 @@ import { parseClaudeLine, parseCodexLine, parseGrokLine, + parseOpenCodeMessage, totalTokens, } from "./usageTranscripts.ts"; @@ -564,3 +565,139 @@ describe("parseGrokLine", () => { expect(records[0]?.timestampMs).toBe(1_786_372_566_000); }); }); + +describe("parseOpenCodeMessage", () => { + function openCodeData(overrides: Record = {}): string { + return JSON.stringify({ + role: "assistant", + tokens: { input: 100, output: 50, reasoning: 20, cache: { read: 30, write: 10 } }, + modelID: "mimo-v2.5-free", + providerID: "opencode", + time: { created: 1_786_372_566_000, completed: 1_786_372_567_000 }, + cost: 0.002, + ...overrides, + }); + } + + it("extracts token totals with input cache-exclusive", () => { + const record = parseOpenCodeMessage(openCodeData(), "msg_1", "ses_1"); + expect(record).not.toBeNull(); + expect(record?.provider).toBe("opencode"); + expect(record?.model).toBe("opencode/mimo-v2.5-free"); + // input is exclusive, not input - cache + expect(record?.totals).toEqual({ + uncachedInputTokens: 100, + cachedInputTokens: 30, + cacheCreationTokens: 10, + outputTokens: 50, + reasoningTokens: 20, + }); + expect(record?.dedupeKey).toBe("msg_1"); + }); + + it("returns null for user or non-assistant rows", () => { + expect(parseOpenCodeMessage(openCodeData({ role: "user" }), "msg_1", "ses_1")).toBeNull(); + expect( + parseOpenCodeMessage(JSON.stringify({ role: "assistant" }), "msg_1", "ses_1"), + ).toBeNull(); + expect(parseOpenCodeMessage("not json", "msg_1", "ses_1")).toBeNull(); + }); + + it("clamps reasoningTokens to outputTokens", () => { + const record = parseOpenCodeMessage( + openCodeData({ + tokens: { input: 10, output: 5, reasoning: 100, cache: { read: 0, write: 0 } }, + }), + "msg_1", + "ses_1", + ); + expect(record?.totals.reasoningTokens).toBe(5); + }); + + it("prefers time.completed over time.created, and falls back to column times", () => { + const withBoth = parseOpenCodeMessage(openCodeData(), "msg_1", "ses_1"); + expect(withBoth?.timestampMs).toBe(1_786_372_567_000); + + const withCreatedOnly = parseOpenCodeMessage( + openCodeData({ time: { created: 1_786_372_566_000 } }), + "msg_1", + "ses_1", + ); + expect(withCreatedOnly?.timestampMs).toBe(1_786_372_566_000); + + const withColumn = parseOpenCodeMessage( + JSON.stringify({ + role: "assistant", + tokens: { input: 10, output: 5, cache: { read: 0, write: 0 } }, + modelID: "m", + providerID: "opencode", + }), + "msg_1", + "ses_1", + undefined, + 1_786_372_566_000, + 1_786_372_567_000, + ); + expect(withColumn?.timestampMs).toBe(1_786_372_567_000); + + const noTime = parseOpenCodeMessage( + JSON.stringify({ + role: "assistant", + tokens: { input: 10, output: 5, cache: { read: 0, write: 0 } }, + modelID: "m", + providerID: "opencode", + }), + "msg_1", + "ses_1", + ); + expect(noTime).toBeNull(); + }); + + it("rejects zero-token rows", () => { + expect( + parseOpenCodeMessage( + openCodeData({ + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + }), + "msg_1", + "ses_1", + ), + ).toBeNull(); + }); + + it("falls back to null dedupeKey on empty message id", () => { + const record = parseOpenCodeMessage(openCodeData(), "", "ses_1"); + expect(record?.dedupeKey).toBeNull(); + }); + + it("handles V2 nested model and type column", () => { + const v2Data = JSON.stringify({ + tokens: { input: 200, output: 100, reasoning: 10, cache: { read: 50, write: 5 } }, + model: { id: "claude-sonnet-4", providerID: "anthropic" }, + time: { completed: 1_786_372_567_000 }, + cost: 0.01, + }); + const record = parseOpenCodeMessage( + v2Data, + "msg_v2", + "ses_v2", + "assistant", + undefined, + 1_786_372_567_000, + ); + expect(record).not.toBeNull(); + expect(record?.model).toBe("anthropic/claude-sonnet-4"); + expect(record?.provider).toBe("opencode"); + }); + + it("uses column type for V2 and ignores non-assistant", () => { + const data = JSON.stringify({ + tokens: { input: 10, output: 5, cache: { read: 0, write: 0 } }, + modelID: "m", + providerID: "opencode", + time: { completed: 1_786_372_567_000 }, + }); + expect(parseOpenCodeMessage(data, "msg_1", "ses_1", "user")).toBeNull(); + expect(parseOpenCodeMessage(data, "msg_1", "ses_1", "assistant")).not.toBeNull(); + }); +}); diff --git a/apps/server/src/usage/usageTranscripts.ts b/apps/server/src/usage/usageTranscripts.ts index b9220b987448..19f0f1733468 100644 --- a/apps/server/src/usage/usageTranscripts.ts +++ b/apps/server/src/usage/usageTranscripts.ts @@ -501,6 +501,9 @@ export function parseOpenCodeMessage( data: string, messageId: string, sessionId: string, + columnType?: string, + columnTimeCreated?: number, + columnTimeUpdated?: number, ): UsageRecord | null { let parsed: unknown; try { @@ -510,7 +513,9 @@ export function parseOpenCodeMessage( } if (typeof parsed !== "object" || parsed === null) return null; const record = parsed as Record; - if (record.role !== "assistant") return null; + // V2 stores role in the `type` column; V1 stores it in `data.role`. + const role = typeof columnType === "string" ? columnType : (record.role as string | undefined); + if (role !== "assistant") return null; const tokensRaw = record.tokens; if (typeof tokensRaw !== "object" || tokensRaw === null) return null; @@ -534,7 +539,7 @@ export function parseOpenCodeMessage( if (totalTokens(totals) === 0) return null; - // Prefer completed timestamp for bucketing, fall back to created. + // Prefer completed timestamp for bucketing, fall back to created, then column times (V2). let timestampMs: number | null = null; const timeRaw = record.time as Record | undefined; if (timeRaw) { @@ -544,10 +549,25 @@ export function parseOpenCodeMessage( timestampMs = timeRaw.created; } } + if (timestampMs === null) { + if (typeof columnTimeUpdated === "number" && Number.isFinite(columnTimeUpdated)) { + timestampMs = columnTimeUpdated; + } else if (typeof columnTimeCreated === "number" && Number.isFinite(columnTimeCreated)) { + timestampMs = columnTimeCreated; + } + } if (timestampMs === null) return null; - const modelId = typeof record.modelID === "string" ? record.modelID : ""; - const providerId = typeof record.providerID === "string" ? record.providerID : "opencode"; + let modelId = typeof record.modelID === "string" ? record.modelID : ""; + let providerId = typeof record.providerID === "string" ? record.providerID : "opencode"; + // V2 nests model under `model.id` / `model.providerID`. + if (modelId.length === 0 && typeof record.model === "object" && record.model !== null) { + const modelObj = record.model as Record; + if (typeof modelObj.id === "string") modelId = modelObj.id; + else if (typeof modelObj.modelID === "string") modelId = modelObj.modelID; + if (typeof modelObj.providerID === "string") providerId = modelObj.providerID; + else if (typeof modelObj.provider === "string") providerId = modelObj.provider; + } const model = modelId.length > 0 ? `${providerId}/${modelId}` : providerId; const cost = record.cost; From d0600ad8bcb4edcd5b301b0446311b31cf0bfdbb Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Mon, 31 Aug 2026 21:57:47 +0300 Subject: [PATCH 04/10] fix(usage): address third review round for opencode --- .../provider/Layers/OpenCodeAdapter.test.ts | 202 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 51 ++++- apps/server/src/usage/UsageService.ts | 8 +- apps/server/src/usage/usageOpenCodeDb.ts | 12 +- 4 files changed, 259 insertions(+), 14 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 9823a68708c2..8b2beb232404 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -4952,6 +4952,208 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("emits thread.token-usage.updated on message.updated with tokens", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-token-usage"); + const messageUpdatedEvent = promiseWithResolvers(); + runtimeMock.state.subscribedEvents = [messageUpdatedEvent.promise]; + + const tokenUsageFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "thread.token-usage.updated", + ), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "test tokens", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/mimo-test", + ), + }); + + messageUpdatedEvent.resolve({ + id: "evt-token-usage", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_token", + role: "assistant", + tokens: { input: 100, output: 50, reasoning: 10, cache: { read: 20, write: 5 } }, + cost: 0.002, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + const events = Array.from( + yield* Fiber.join(tokenUsageFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(events.length, 1); + const usage = (events[0] as { payload: { usage: Record } }).payload.usage; + NodeAssert.equal((usage as { inputTokens: number }).inputTokens, 100); + NodeAssert.equal((usage as { cachedInputTokens: number }).cachedInputTokens, 20); + NodeAssert.equal((usage as { outputTokens: number }).outputTokens, 50); + + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("does not emit token usage on zero-token payload", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-zero-token"); + const messageUpdatedEvent = promiseWithResolvers(); + const idleEvent = promiseWithResolvers(); + runtimeMock.state.subscribedEvents = [messageUpdatedEvent.promise, idleEvent.promise]; + + const turnCompletedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "zero tokens", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/mimo-test", + ), + }); + + messageUpdatedEvent.resolve({ + id: "evt-zero-token", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_zero", + role: "assistant", + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + cost: 0, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + yield* Effect.yieldNow; + + idleEvent.resolve({ + id: "evt-zero-idle", + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const events = Array.from( + yield* Fiber.join(turnCompletedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(events.length, 1); + const payload = (events[0] as { payload: Record }).payload; + NodeAssert.equal(payload.usage, undefined); + NodeAssert.equal(String((events[0] as { turnId: unknown }).turnId), String(turn.turnId)); + + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("idle turn completion carries usage and modelUsage", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-idle-usage"); + const messageUpdatedEvent = promiseWithResolvers(); + const idleEvent = promiseWithResolvers(); + runtimeMock.state.subscribedEvents = [messageUpdatedEvent.promise, idleEvent.promise]; + + const turnCompletedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "test idle usage", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/mimo-test", + ), + }); + + messageUpdatedEvent.resolve({ + id: "evt-idle-usage-msg", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_idle", + role: "assistant", + tokens: { input: 100, output: 20, reasoning: 5, cache: { read: 10, write: 2 } }, + cost: 0.001, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + // Yield to let token emission be processed + yield* Effect.yieldNow; + + idleEvent.resolve({ + id: "evt-idle-usage-idle", + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const events = Array.from( + yield* Fiber.join(turnCompletedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(events.length, 1); + const payload = (events[0] as { payload: Record }).payload; + NodeAssert.equal(payload.state, "completed"); + NodeAssert.ok(payload.usage); + NodeAssert.equal(payload.totalCostUsd as number, 0.001); + const modelUsage = payload.modelUsage as Record; + NodeAssert.ok(modelUsage["opencode/mimo-test"]); + + // Ensure turnId matches + NodeAssert.equal(String((events[0] as { turnId: unknown }).turnId), String(turn.turnId)); + + yield* adapter.stopSession(threadId); + }), + ); + it.effect("keeps the event pump alive when native event logging fails", () => Effect.gen(function* () { runtimeMock.state.subscribedEvents = [ diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index e33d428e008d..f35e6ca97e58 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -993,9 +993,10 @@ export function makeOpenCodeAdapter( tokens: OpenCodeTokens, turnId: TurnId | undefined, raw: unknown, + cost?: unknown, + model?: unknown, + providerId?: unknown, ) { - // Remember for turn.completed - context.lastTokens = tokens; if ( tokens.input + tokens.output + tokens.cache.read + tokens.cache.write === 0 && tokens.total === undefined @@ -1004,6 +1005,11 @@ export function makeOpenCodeAdapter( } const snapshot = openCodeTokensToSnapshot(tokens); if (snapshot.usedTokens <= 0) return; + // Remember for turn.completed only after validating non-zero tokens. + // This prevents zero-token events from overwriting previous valid usage + // and ensures cost/model are not leaked on empty turns. + context.lastTokens = tokens; + rememberOpenCodeCostAndModel(context, cost, model, providerId); yield* emit({ ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -1261,6 +1267,10 @@ export function makeOpenCodeAdapter( context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; context.reconcileIdleStatus = false; + // Clear pending usage so it does not leak to the next turn. + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; yield* updateProviderSession( context, { status: "error", lastError: detail }, @@ -1453,6 +1463,10 @@ export function makeOpenCodeAdapter( if (cancellation) { context.cancellation = undefined; } + // Clear any pending usage so it does not leak to the next turn. + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; if (context.activeTurnId === turnId) { context.activeTurnId = undefined; context.activeAgent = undefined; @@ -2091,13 +2105,15 @@ export function makeOpenCodeAdapter( const info = event.properties.info as unknown as Record; const tokens = readOpenCodeTokens(info.tokens); if (tokens) { - rememberOpenCodeCostAndModel( + yield* emitOpenCodeTokenUsage( context, + tokens, + turnId, + event, info.cost, (info.modelID as string | undefined) ?? (info.model as unknown), info.providerID as string | undefined, ); - yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); } } break; @@ -2166,13 +2182,15 @@ export function makeOpenCodeAdapter( const partRecord = part as unknown as Record; const tokens = readOpenCodeTokens(partRecord.tokens); if (tokens) { - rememberOpenCodeCostAndModel( + yield* emitOpenCodeTokenUsage( context, + tokens, + turnId, + event, partRecord.cost, (partRecord.modelID as string | undefined) ?? (partRecord.model as unknown), partRecord.providerID as string | undefined, ); - yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); } } @@ -2331,6 +2349,10 @@ export function makeOpenCodeAdapter( context.activeAgent = undefined; context.activeVariant = undefined; context.reconcileIdleStatus = false; + // Clear pending usage so it does not leak to the next turn. + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; yield* updateProviderSession( context, { @@ -2394,20 +2416,29 @@ export function makeOpenCodeAdapter( (props.providerID as string | undefined) ?? (props.info as Record | undefined)?.providerID ?? (props.part as Record | undefined)?.providerID; - rememberOpenCodeCostAndModel(context, cost, model, providerId); - yield* emitOpenCodeTokenUsage(context, tokens, turnId, event); + yield* emitOpenCodeTokenUsage( + context, + tokens, + turnId, + event, + cost, + model, + providerId, + ); } else { const data = props.data as Record | undefined; if (data) { const nestedTokens = readOpenCodeTokens(data.tokens); if (nestedTokens) { - rememberOpenCodeCostAndModel( + yield* emitOpenCodeTokenUsage( context, + nestedTokens, + turnId, + event, data.cost, (data.modelID as string | undefined) ?? (data.model as unknown), data.providerID as string | undefined, ); - yield* emitOpenCodeTokenUsage(context, nestedTokens, turnId, event); } } } diff --git a/apps/server/src/usage/UsageService.ts b/apps/server/src/usage/UsageService.ts index c454c42ddc24..868d837e4cc0 100644 --- a/apps/server/src/usage/UsageService.ts +++ b/apps/server/src/usage/UsageService.ts @@ -311,7 +311,13 @@ export const make = Effect.gen(function* () { return records; }); - /** Parses the OpenCode SQLite DB, reusing cached result when unchanged. */ + /** + * Parses the OpenCode SQLite DB, reusing cached result when unchanged. + * `readOpenCodeDbRecords` yields to the event loop every 100 rows via + * `setImmediate`, so a cold scan of ~4k messages does not block the server + * (observed 7 daily buckets from 4085 messages). The `size`/`mtimeMs` + * fingerprint includes the WAL file, so WAL-only writes invalidate the cache. + */ const readOpenCodeDbRecordsCached = ( dbPath: string, size: number, diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index ab0357504b15..2dae0b89003a 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -109,23 +109,27 @@ export async function readOpenCodeDbRecords( for (const table of tables) { try { // Use iterate() instead of all() to avoid loading all rows at once and to allow yielding. + // Keep the prepared statement alive for the duration of the iteration to avoid + // finalization mid-scan when yielding (node:sqlite finalizes statements that lose + // their JS reference during an async pause). let rows: Iterable>; let isSessionMessage = table === "session_message"; + let stmt: ReturnType["prepare"]> | null = null; try { if (isSessionMessage) { - const stmt = db.prepare( + stmt = db.prepare( `SELECT id, session_id, type, time_created, time_updated, data FROM "${table}"`, ); rows = stmt.iterate() as Iterable>; } else { - const stmt = db.prepare( + stmt = db.prepare( `SELECT id, session_id, time_created, time_updated, data FROM "${table}"`, ); rows = stmt.iterate() as Iterable>; } } catch { // Fallback for older schemas without time columns - const stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); + stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); rows = stmt.iterate() as Iterable>; isSessionMessage = false; } @@ -154,6 +158,8 @@ export async function readOpenCodeDbRecords( await new Promise((resolve) => setImmediate(resolve)); } } + // Retain reference to stmt until iteration completes (see comment above). + void stmt; } catch { // One table failing shouldn't hide the other; treat as empty for that table. continue; From bb91534f9d897b10cdbcb718dd2ffb5c26fd591d Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Tue, 1 Sep 2026 01:34:40 +0300 Subject: [PATCH 05/10] fix(server): accumulate opencode tokens and cost per turn --- .../provider/Layers/OpenCodeAdapter.test.ts | 106 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 64 +++++++++-- 2 files changed, 161 insertions(+), 9 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 8b2beb232404..ab915c333d38 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -5154,6 +5154,112 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("accumulates usage and cost across multiple tool steps", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-accumulated-usage"); + const firstMessage = promiseWithResolvers(); + const secondMessage = promiseWithResolvers(); + const idleEvent = promiseWithResolvers(); + runtimeMock.state.subscribedEvents = [ + firstMessage.promise, + secondMessage.promise, + idleEvent.promise, + ]; + + const turnCompletedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "test accumulation", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/mimo-test", + ), + }); + + firstMessage.resolve({ + id: "evt-accum-1", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_accum_1", + role: "assistant", + tokens: { input: 100, output: 20, reasoning: 5, cache: { read: 10, write: 2 } }, + cost: 0.001, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + yield* Effect.yieldNow; + + secondMessage.resolve({ + id: "evt-accum-2", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_accum_2", + role: "assistant", + tokens: { input: 50, output: 30, reasoning: 10, cache: { read: 5, write: 1 } }, + cost: 0.002, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + yield* Effect.yieldNow; + + idleEvent.resolve({ + id: "evt-accum-idle", + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const events = Array.from( + yield* Fiber.join(turnCompletedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(events.length, 1); + const payload = (events[0] as { payload: Record }).payload; + const usage = payload.usage as { + input: number; + output: number; + reasoning: number; + cache: { read: number; write: number }; + }; + // Input 100+50, output 20+30, cache read 10+5, write 2+1 + NodeAssert.equal(usage.input, 150); + NodeAssert.equal(usage.output, 50); + NodeAssert.equal(usage.cache.read, 15); + NodeAssert.equal(usage.cache.write, 3); + // Cost summed + NodeAssert.equal(payload.totalCostUsd as number, 0.003); + const modelUsage = payload.modelUsage as Record; + NodeAssert.ok(modelUsage["opencode/mimo-test"]); + NodeAssert.equal(modelUsage["opencode/mimo-test"].input, 150); + NodeAssert.equal(String((events[0] as { turnId: unknown }).turnId), String(turn.turnId)); + + yield* adapter.stopSession(threadId); + }), + ); + it.effect("keeps the event pump alive when native event logging fails", () => Effect.gen(function* () { runtimeMock.state.subscribedEvents = [ diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index f35e6ca97e58..ec5ec64630c3 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -341,6 +341,8 @@ interface OpenCodeSessionContext { lastTokens: OpenCodeTokens | null; lastCost: number | null; lastModel: string | null; + /** Accumulated per-model tokens for the active turn, for `modelUsage`. */ + modelUsage: Map; cancellation: OpenCodeCancellation | undefined; interruptedTurnId: TurnId | undefined; reconcileIdleStatus: boolean; @@ -696,6 +698,19 @@ function openCodeTokensToSnapshot(tokens: OpenCodeTokens): { }; } +function mergeOpenCodeTokens(a: OpenCodeTokens, b: OpenCodeTokens): OpenCodeTokens { + const totalA = a.total ?? a.input + a.output + a.cache.read + a.cache.write; + const totalB = b.total ?? b.input + b.output + b.cache.read + b.cache.write; + const hasTotal = a.total !== undefined || b.total !== undefined; + return { + input: a.input + b.input, + output: a.output + b.output, + reasoning: a.reasoning + b.reasoning, + cache: { read: a.cache.read + b.cache.read, write: a.cache.write + b.cache.write }, + ...(hasTotal ? { total: totalA + totalB } : {}), + }; +} + function updateProviderSession( context: OpenCodeSessionContext, patch: Partial, @@ -1005,11 +1020,15 @@ export function makeOpenCodeAdapter( } const snapshot = openCodeTokensToSnapshot(tokens); if (snapshot.usedTokens <= 0) return; - // Remember for turn.completed only after validating non-zero tokens. - // This prevents zero-token events from overwriting previous valid usage - // and ensures cost/model are not leaked on empty turns. - context.lastTokens = tokens; - rememberOpenCodeCostAndModel(context, cost, model, providerId); + // Accumulate for turn.completed — tool-calling turns emit multiple + // assistant records, so we sum tokens/cost and group by model for + // modelUsage instead of overwriting the last value. + if (context.lastTokens === null) { + context.lastTokens = tokens; + } else { + context.lastTokens = mergeOpenCodeTokens(context.lastTokens, tokens); + } + rememberOpenCodeCostAndModel(context, cost, model, providerId, tokens); yield* emit({ ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -1028,14 +1047,18 @@ export function makeOpenCodeAdapter( cost: unknown, model: unknown, providerId: unknown, + tokens?: OpenCodeTokens, ) => { - if (typeof cost === "number" && Number.isFinite(cost)) context.lastCost = cost; + if (typeof cost === "number" && Number.isFinite(cost)) { + context.lastCost = (context.lastCost ?? 0) + cost; + } + let modelKey: string | null = null; if (typeof model === "string" && model.trim().length > 0) { const provider = typeof providerId === "string" && providerId.trim().length > 0 ? providerId.trim() : "opencode"; - context.lastModel = `${provider}/${model.trim()}`; + modelKey = `${provider}/${model.trim()}`; } else if (typeof model === "object" && model !== null) { const record = model as Record; const id = @@ -1050,7 +1073,18 @@ export function makeOpenCodeAdapter( : typeof record.provider === "string" ? record.provider : "opencode"; - if (id) context.lastModel = `${prov}/${id}`; + if (id) modelKey = `${prov}/${id}`; + } + if (modelKey) { + context.lastModel = modelKey; + if (tokens) { + const existing = context.modelUsage.get(modelKey); + if (existing) { + context.modelUsage.set(modelKey, mergeOpenCodeTokens(existing, tokens)); + } else { + context.modelUsage.set(modelKey, tokens); + } + } } }; @@ -1105,9 +1139,17 @@ export function makeOpenCodeAdapter( const usage = context.lastTokens; const cost = context.lastCost; const model = context.lastModel; + const modelUsageEntries = [...context.modelUsage.entries()]; context.lastTokens = null; context.lastCost = null; context.lastModel = null; + context.modelUsage.clear(); + const modelUsage = + modelUsageEntries.length > 0 + ? Object.fromEntries(modelUsageEntries) + : model && usage + ? { [model]: usage } + : undefined; yield* emit({ ...(yield* buildEventBase({ threadId: context.session.threadId, @@ -1120,7 +1162,7 @@ export function makeOpenCodeAdapter( ...(usage ? { usage, - ...(model ? { modelUsage: { [model]: usage } } : {}), + ...(modelUsage ? { modelUsage } : {}), } : {}), ...(cost !== null ? { totalCostUsd: cost } : {}), @@ -1271,6 +1313,7 @@ export function makeOpenCodeAdapter( context.lastTokens = null; context.lastCost = null; context.lastModel = null; + context.modelUsage.clear(); yield* updateProviderSession( context, { status: "error", lastError: detail }, @@ -1467,6 +1510,7 @@ export function makeOpenCodeAdapter( context.lastTokens = null; context.lastCost = null; context.lastModel = null; + context.modelUsage.clear(); if (context.activeTurnId === turnId) { context.activeTurnId = undefined; context.activeAgent = undefined; @@ -2353,6 +2397,7 @@ export function makeOpenCodeAdapter( context.lastTokens = null; context.lastCost = null; context.lastModel = null; + context.modelUsage.clear(); yield* updateProviderSession( context, { @@ -2707,6 +2752,7 @@ export function makeOpenCodeAdapter( lastTokens: null, lastCost: null, lastModel: null, + modelUsage: new Map(), cancellation: undefined, interruptedTurnId: undefined, reconcileIdleStatus: false, From b34f72e2b242485933e10eb94ad39b1c08423540 Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Tue, 1 Sep 2026 01:41:57 +0300 Subject: [PATCH 06/10] fix(server): handle opencode fallback schema and clear stale tokens on prompt failure --- apps/server/src/provider/Layers/OpenCodeAdapter.ts | 8 ++++++++ apps/server/src/usage/usageOpenCodeDb.ts | 7 +++++-- 2 files changed, 13 insertions(+), 2 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index ec5ec64630c3..1dd5c571713b 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -3010,6 +3010,10 @@ export function makeOpenCodeAdapter( context.activeTurnId = undefined; context.activeAgent = undefined; context.activeVariant = undefined; + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; + context.modelUsage.clear(); yield* updateProviderSession( context, { @@ -3056,6 +3060,10 @@ export function makeOpenCodeAdapter( context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; context.reconcileIdleStatus = false; + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; + context.modelUsage.clear(); yield* updateProviderSession( context, { diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index 2dae0b89003a..513375439465 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -129,9 +129,12 @@ export async function readOpenCodeDbRecords( } } catch { // Fallback for older schemas without time columns - stmt = db.prepare(`SELECT id, session_id, data FROM "${table}"`); + stmt = db.prepare( + isSessionMessage + ? `SELECT id, session_id, type, data FROM "${table}"` + : `SELECT id, session_id, data FROM "${table}"`, + ); rows = stmt.iterate() as Iterable>; - isSessionMessage = false; } let count = 0; for (const raw of rows) { From 2f7bb3cb0a3fe4583df29e9a976c27b763daa373 Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Tue, 1 Sep 2026 01:53:17 +0300 Subject: [PATCH 07/10] fix(server): dedupe opencode token events and guard stale turn accumulation --- .../provider/Layers/OpenCodeAdapter.test.ts | 103 ++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 132 ++++++++++++------ 2 files changed, 194 insertions(+), 41 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index ab915c333d38..00f89d334fd6 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -5260,6 +5260,109 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("does not double count duplicate message.updated for same id", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-dedupe-usage"); + const firstMessage = promiseWithResolvers(); + const secondMessage = promiseWithResolvers(); + const idleEvent = promiseWithResolvers(); + runtimeMock.state.subscribedEvents = [ + firstMessage.promise, + secondMessage.promise, + idleEvent.promise, + ]; + + const turnCompletedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(1), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const turn = yield* adapter.sendTurn({ + threadId, + input: "test dedupe", + modelSelection: createModelSelection( + ProviderInstanceId.make("opencode"), + "opencode/mimo-test", + ), + }); + + firstMessage.resolve({ + id: "evt-dedupe-1", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_dedupe", + role: "assistant", + tokens: { input: 100, output: 20, reasoning: 5, cache: { read: 10, write: 2 } }, + cost: 0.001, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + yield* Effect.yieldNow; + + // Same message id, updated tokens (cumulative) — should replace, not sum + secondMessage.resolve({ + id: "evt-dedupe-2", + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg_dedupe", + role: "assistant", + tokens: { input: 150, output: 30, reasoning: 10, cache: { read: 15, write: 3 } }, + cost: 0.002, + modelID: "mimo-test", + providerID: "opencode", + }, + }, + }); + + yield* Effect.yieldNow; + + idleEvent.resolve({ + id: "evt-dedupe-idle", + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const events = Array.from( + yield* Fiber.join(turnCompletedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(events.length, 1); + const payload = (events[0] as { payload: Record }).payload; + const usage = payload.usage as { + input: number; + output: number; + cache: { read: number; write: number }; + }; + // Should be latest snapshot, not sum of both (100+150) + NodeAssert.equal(usage.input, 150); + NodeAssert.equal(usage.output, 30); + NodeAssert.equal(usage.cache.read, 15); + // Cost should be latest as well (deduped) — 0.002 not 0.003 + // For deduped message, cost is replaced, so total is 0.002 + NodeAssert.equal(payload.totalCostUsd as number, 0.002); + NodeAssert.equal(String((events[0] as { turnId: unknown }).turnId), String(turn.turnId)); + + yield* adapter.stopSession(threadId); + }), + ); + it.effect("keeps the event pump alive when native event logging fails", () => Effect.gen(function* () { runtimeMock.state.subscribedEvents = [ diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 1dd5c571713b..34c1cf6235f8 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -343,6 +343,13 @@ interface OpenCodeSessionContext { lastModel: string | null; /** Accumulated per-model tokens for the active turn, for `modelUsage`. */ modelUsage: Map; + /** Dedupe map for token-bearing events within the active turn, keyed by message/part id. */ + tokenDedupeMap: Map< + string, + { tokens: OpenCodeTokens; cost: number | null; modelKey: string | null } + >; + /** Counter for generic token events without a stable id. */ + dedupeCounter: number; cancellation: OpenCodeCancellation | undefined; interruptedTurnId: TurnId | undefined; reconcileIdleStatus: boolean; @@ -1011,6 +1018,7 @@ export function makeOpenCodeAdapter( cost?: unknown, model?: unknown, providerId?: unknown, + dedupeKey?: string, ) { if ( tokens.input + tokens.output + tokens.cache.read + tokens.cache.write === 0 && @@ -1020,38 +1028,26 @@ export function makeOpenCodeAdapter( } const snapshot = openCodeTokensToSnapshot(tokens); if (snapshot.usedTokens <= 0) return; - // Accumulate for turn.completed — tool-calling turns emit multiple - // assistant records, so we sum tokens/cost and group by model for - // modelUsage instead of overwriting the last value. - if (context.lastTokens === null) { - context.lastTokens = tokens; - } else { - context.lastTokens = mergeOpenCodeTokens(context.lastTokens, tokens); - } - rememberOpenCodeCostAndModel(context, cost, model, providerId, tokens); - yield* emit({ - ...(yield* buildEventBase({ - threadId: context.session.threadId, - turnId, - raw, - })), - type: "thread.token-usage.updated", - payload: { - usage: snapshot, - }, - }); - }); - - const rememberOpenCodeCostAndModel = ( - context: OpenCodeSessionContext, - cost: unknown, - model: unknown, - providerId: unknown, - tokens?: OpenCodeTokens, - ) => { - if (typeof cost === "number" && Number.isFinite(cost)) { - context.lastCost = (context.lastCost ?? 0) + cost; + // Guard: only accumulate for the active turn. Delayed events with + // undefined or stale turnId would otherwise repopulate after completion + // and leak into the next turn (see filtered issue at line 1026). + const isActiveTurn = turnId !== undefined && turnId === context.activeTurnId; + if (!isActiveTurn) { + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + raw, + })), + type: "thread.token-usage.updated", + payload: { usage: snapshot }, + }); + return; } + // Dedupe by message/part id to avoid double-counting when OpenCode + // re-emits `message.updated` with cumulative totals for the same id. + let costValue: number | null = + typeof cost === "number" && Number.isFinite(cost) ? cost : null; let modelKey: string | null = null; if (typeof model === "string" && model.trim().length > 0) { const provider = @@ -1075,18 +1071,43 @@ export function makeOpenCodeAdapter( : "opencode"; if (id) modelKey = `${prov}/${id}`; } - if (modelKey) { - context.lastModel = modelKey; - if (tokens) { - const existing = context.modelUsage.get(modelKey); - if (existing) { - context.modelUsage.set(modelKey, mergeOpenCodeTokens(existing, tokens)); - } else { - context.modelUsage.set(modelKey, tokens); - } + const key = dedupeKey ?? `generic:${context.dedupeCounter++}`; + context.tokenDedupeMap.set(key, { tokens, cost: costValue, modelKey }); + // Recompute aggregated totals from the dedupe map. + let aggTokens: OpenCodeTokens | null = null; + let aggCostSum = 0; + let hasCost = false; + const newModelUsage = new Map(); + for (const entry of context.tokenDedupeMap.values()) { + if (aggTokens === null) aggTokens = entry.tokens; + else aggTokens = mergeOpenCodeTokens(aggTokens, entry.tokens); + if (entry.cost !== null) { + aggCostSum += entry.cost; + hasCost = true; + } + if (entry.modelKey) { + const existing = newModelUsage.get(entry.modelKey); + if (existing) + newModelUsage.set(entry.modelKey, mergeOpenCodeTokens(existing, entry.tokens)); + else newModelUsage.set(entry.modelKey, entry.tokens); } } - }; + context.lastTokens = aggTokens; + context.lastCost = hasCost ? aggCostSum : null; + if (modelKey) context.lastModel = modelKey; + context.modelUsage = newModelUsage; + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + raw, + })), + type: "thread.token-usage.updated", + payload: { + usage: snapshot, + }, + }); + }); const cancelIdleReconciliation = Effect.fn("cancelIdleReconciliation")(function* ( context: OpenCodeSessionContext, @@ -1144,6 +1165,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; const modelUsage = modelUsageEntries.length > 0 ? Object.fromEntries(modelUsageEntries) @@ -1314,6 +1337,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; yield* updateProviderSession( context, { status: "error", lastError: detail }, @@ -1511,6 +1536,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; if (context.activeTurnId === turnId) { context.activeTurnId = undefined; context.activeAgent = undefined; @@ -2149,6 +2176,7 @@ export function makeOpenCodeAdapter( const info = event.properties.info as unknown as Record; const tokens = readOpenCodeTokens(info.tokens); if (tokens) { + const dedupeKey = typeof info.id === "string" ? `msg:${info.id}` : undefined; yield* emitOpenCodeTokenUsage( context, tokens, @@ -2157,6 +2185,7 @@ export function makeOpenCodeAdapter( info.cost, (info.modelID as string | undefined) ?? (info.model as unknown), info.providerID as string | undefined, + dedupeKey, ); } } @@ -2226,6 +2255,7 @@ export function makeOpenCodeAdapter( const partRecord = part as unknown as Record; const tokens = readOpenCodeTokens(partRecord.tokens); if (tokens) { + const dedupeKey = typeof part.id === "string" ? `part:${part.id}` : undefined; yield* emitOpenCodeTokenUsage( context, tokens, @@ -2234,6 +2264,7 @@ export function makeOpenCodeAdapter( partRecord.cost, (partRecord.modelID as string | undefined) ?? (partRecord.model as unknown), partRecord.providerID as string | undefined, + dedupeKey, ); } } @@ -2398,6 +2429,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; yield* updateProviderSession( context, { @@ -2753,6 +2786,8 @@ export function makeOpenCodeAdapter( lastCost: null, lastModel: null, modelUsage: new Map(), + tokenDedupeMap: new Map(), + dedupeCounter: 0, cancellation: undefined, interruptedTurnId: undefined, reconcileIdleStatus: false, @@ -2921,6 +2956,17 @@ export function makeOpenCodeAdapter( context.promptGeneration = promptGeneration; context.promptAdmission = promptAdmission; + // Fresh turn — clear any stale accumulated usage from a prior turn + // that may have been repopulated by a delayed event (see line 1026). + if (steeringTurnId === undefined) { + context.lastTokens = null; + context.lastCost = null; + context.lastModel = null; + context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; + } + context.activeTurnId = turnId; context.activeAgent = agent ?? (input.interactionMode === "plan" ? "plan" : undefined); context.activeVariant = variant; @@ -3014,6 +3060,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; yield* updateProviderSession( context, { @@ -3064,6 +3112,8 @@ export function makeOpenCodeAdapter( context.lastCost = null; context.lastModel = null; context.modelUsage.clear(); + context.tokenDedupeMap.clear(); + context.dedupeCounter = 0; yield* updateProviderSession( context, { From 10d08446f23fc589d3adabd1fa254995eae2eaff Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Tue, 1 Sep 2026 01:57:05 +0300 Subject: [PATCH 08/10] fix(server): remove ineffective eslint disable for opencode scan --- apps/server/src/usage/usageOpenCodeDb.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index 513375439465..a98a2c596751 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -157,7 +157,6 @@ export async function readOpenCodeDbRecords( if (record !== null) records.push(record); // Yield to event loop every 100 rows to avoid blocking cold scans. if (++count % 100 === 0) { - // eslint-disable-next-line no-await-in-loop await new Promise((resolve) => setImmediate(resolve)); } } From 897e1ff284115b554acc61a55befdceec6bacef8 Mon Sep 17 00:00:00 2001 From: Mina Sayed Date: Tue, 1 Sep 2026 02:03:38 +0300 Subject: [PATCH 09/10] fix(server): prevent stale token attribution and handle db scan failures --- .../src/provider/Layers/OpenCodeAdapter.ts | 25 +++++++++++++++++++ apps/server/src/usage/usageOpenCodeDb.ts | 3 +-- 2 files changed, 26 insertions(+), 2 deletions(-) diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 34c1cf6235f8..6e8feda42990 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -350,6 +350,8 @@ interface OpenCodeSessionContext { >; /** Counter for generic token events without a stable id. */ dedupeCounter: number; + /** Persistent map from dedupeKey to the turn that created it, for stale-event detection. */ + messageTurnMap: Map; cancellation: OpenCodeCancellation | undefined; interruptedTurnId: TurnId | undefined; reconcileIdleStatus: boolean; @@ -1071,6 +1073,28 @@ export function makeOpenCodeAdapter( : "opencode"; if (id) modelKey = `${prov}/${id}`; } + // Stale-event guard: if this dedupeKey was already seen for a different + // turn, this is a late event from a prior turn that would otherwise + // contaminate the new turn's totals. + if (dedupeKey) { + const existingTurnId = context.messageTurnMap.get(dedupeKey); + if (existingTurnId !== undefined) { + if (existingTurnId !== turnId) { + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId, + raw, + })), + type: "thread.token-usage.updated", + payload: { usage: snapshot }, + }); + return; + } + } else { + context.messageTurnMap.set(dedupeKey, turnId); + } + } const key = dedupeKey ?? `generic:${context.dedupeCounter++}`; context.tokenDedupeMap.set(key, { tokens, cost: costValue, modelKey }); // Recompute aggregated totals from the dedupe map. @@ -2788,6 +2812,7 @@ export function makeOpenCodeAdapter( modelUsage: new Map(), tokenDedupeMap: new Map(), dedupeCounter: 0, + messageTurnMap: new Map(), cancellation: undefined, interruptedTurnId: undefined, reconcileIdleStatus: false, diff --git a/apps/server/src/usage/usageOpenCodeDb.ts b/apps/server/src/usage/usageOpenCodeDb.ts index a98a2c596751..2a499958c8b1 100644 --- a/apps/server/src/usage/usageOpenCodeDb.ts +++ b/apps/server/src/usage/usageOpenCodeDb.ts @@ -163,8 +163,7 @@ export async function readOpenCodeDbRecords( // Retain reference to stmt until iteration completes (see comment above). void stmt; } catch { - // One table failing shouldn't hide the other; treat as empty for that table. - continue; + return null; } } return records; From 77512998485718eb1b6c336e20f196eb40a6a32f Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Wed, 2 Sep 2026 17:49:09 -0700 Subject: [PATCH 10/10] feat(web): preview document attachments in the file viewer (#9292) Co-authored-by: Julius Marminge Co-authored-by: Claude Fable 5 --- apps/server/src/assets/AssetAccess.test.ts | 61 ++++++ apps/server/src/assets/AssetAccess.ts | 32 ++- apps/server/src/http.test.ts | 13 ++ apps/server/src/http.ts | 26 ++- apps/web/src/components/ChatView.tsx | 33 ++- apps/web/src/components/RightPanelTabs.tsx | 6 +- .../components/chat/MessagesTimeline.test.tsx | 20 +- .../src/components/chat/MessagesTimeline.tsx | 59 ++++- .../components/files/FilePreviewPanel.test.ts | 42 +++- .../src/components/files/FilePreviewPanel.tsx | 207 +++++++++++++----- .../src/components/files/filePreviewMode.ts | 13 ++ apps/web/src/rightPanelStore.test.ts | 73 ++++++ apps/web/src/rightPanelStore.ts | 33 ++- apps/web/src/types.ts | 9 + docs/user/composer.md | 11 +- packages/contracts/src/assets.ts | 3 + 16 files changed, 540 insertions(+), 101 deletions(-) diff --git a/apps/server/src/assets/AssetAccess.test.ts b/apps/server/src/assets/AssetAccess.test.ts index 8d7c696cd625..539827a0a5ce 100644 --- a/apps/server/src/assets/AssetAccess.test.ts +++ b/apps/server/src/assets/AssetAccess.test.ts @@ -576,6 +576,67 @@ describe("AssetAccess", () => { }); }).pipe(Effect.provide(testLayer)), ); + + it.effect("serves document attachments inline when a viewer requests it", () => + Effect.gen(function* () { + const config = yield* ServerConfig.ServerConfig; + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const attachmentId = "thread-1-00000000-0000-4000-8000-000000000001-pdf"; + const attachmentPath = path.join(config.attachmentsDir, `${attachmentId}.pdf`); + yield* fileSystem.makeDirectory(config.attachmentsDir, { recursive: true }); + yield* fileSystem.writeFile(attachmentPath, new Uint8Array([1, 2, 3])); + + const result = yield* issueAssetUrl({ + resource: { + _tag: "attachment", + attachmentId, + fileName: "report.pdf", + mimeType: "application/pdf", + disposition: "inline", + }, + }); + const suffix = result.relativeUrl.slice(`${ASSET_ROUTE_PREFIX}/`.length); + const separatorIndex = suffix.indexOf("/"); + + expect( + yield* resolveAsset(suffix.slice(0, separatorIndex), suffix.slice(separatorIndex + 1)), + ).toEqual({ + kind: "file", + path: attachmentPath, + fileName: "report.pdf", + mimeType: "application/pdf", + }); + }).pipe(Effect.provide(testLayer)), + ); + + it.effect("keeps inline requests for other attachment types as downloads", () => + Effect.gen(function* () { + const config = yield* ServerConfig.ServerConfig; + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const attachmentId = "thread-1-00000000-0000-4000-8000-000000000002-zip"; + const attachmentPath = path.join(config.attachmentsDir, `${attachmentId}.zip`); + yield* fileSystem.makeDirectory(config.attachmentsDir, { recursive: true }); + yield* fileSystem.writeFile(attachmentPath, new Uint8Array([1, 2, 3])); + + const result = yield* issueAssetUrl({ + resource: { + _tag: "attachment", + attachmentId, + fileName: "archive.zip", + mimeType: "text/html", + disposition: "inline", + }, + }); + const suffix = result.relativeUrl.slice(`${ASSET_ROUTE_PREFIX}/`.length); + const separatorIndex = suffix.indexOf("/"); + + expect( + yield* resolveAsset(suffix.slice(0, separatorIndex), suffix.slice(separatorIndex + 1)), + ).toMatchObject({ kind: "file", path: attachmentPath, download: true }); + }).pipe(Effect.provide(testLayer)), + ); it.effect("issues project favicon capabilities with a signed fallback", () => Effect.gen(function* () { const fileSystem = yield* FileSystem.FileSystem; diff --git a/apps/server/src/assets/AssetAccess.ts b/apps/server/src/assets/AssetAccess.ts index ad7cb273cf07..36fa9130b683 100644 --- a/apps/server/src/assets/AssetAccess.ts +++ b/apps/server/src/assets/AssetAccess.ts @@ -51,6 +51,14 @@ const ASSET_TOKEN_TTL_MS = 60 * 60 * 1000; const PROJECT_FAVICON_TOKEN_BUCKET_MS = 30 * 60 * 1000; const PROJECT_FAVICON_VERSION_PREFIX = "v"; const INLINE_VIDEO_MIME_TYPE_PATTERN = /^video\/[\w!#$&^.+-]+$/i; +// Extensions a document viewer may request inline. The extension comes from +// the attachment id the server assigned, never from the client's mime type. +const INLINE_DOCUMENT_EXTENSIONS = new Set(["pdf", "html", "htm"]); +const INLINE_DOCUMENT_MIME_TYPES: Record = { + pdf: "application/pdf", + html: "text/html", + htm: "text/html", +}; const PREVIEW_ASSET_EXTENSIONS = new Set([ ...WORKSPACE_BROWSER_PREVIEW_EXTENSIONS, ...WORKSPACE_IMAGE_PREVIEW_EXTENSIONS, @@ -362,19 +370,31 @@ export const issueAssetUrl = Effect.fn("AssetAccess.issueAssetUrl")(function* (i } // Generic files carry their extension inside the attachment id (that // shape resolves the on-disk path); images do not. Videos and images - // render inline; other generic files download. - const isGenericFile = parseAttachmentFileExtension(input.resource.attachmentId) !== null; + // render inline. Other generic files download, unless a document viewer + // asked for inline and the stored extension is one a browser can show. + const extension = parseAttachmentFileExtension(input.resource.attachmentId); + const isGenericFile = extension !== null; const videoMimeType = input.resource.mimeType?.split(";", 1)[0]?.trim() ?? ""; const isVideo = INLINE_VIDEO_MIME_TYPE_PATTERN.test(videoMimeType); + const inlineDocumentMimeType = + input.resource.disposition === "inline" && + extension !== null && + INLINE_DOCUMENT_EXTENSIONS.has(extension) + ? INLINE_DOCUMENT_MIME_TYPES[extension] + : undefined; claims = { version: 1, kind: "attachment", attachmentId: input.resource.attachmentId, - ...(isGenericFile && !isVideo ? { download: true } : {}), - ...(input.resource.fileName !== undefined ? { fileName: input.resource.fileName } : {}), - ...(input.resource.mimeType !== undefined - ? { mimeType: isVideo ? videoMimeType : input.resource.mimeType } + ...(isGenericFile && !isVideo && inlineDocumentMimeType === undefined + ? { download: true } : {}), + ...(input.resource.fileName !== undefined ? { fileName: input.resource.fileName } : {}), + ...(inlineDocumentMimeType !== undefined + ? { mimeType: inlineDocumentMimeType } + : input.resource.mimeType !== undefined + ? { mimeType: isVideo ? videoMimeType : input.resource.mimeType } + : {}), expiresAt, }; fileName = input.resource.fileName ?? path.basename(attachmentPath); diff --git a/apps/server/src/http.test.ts b/apps/server/src/http.test.ts index 6b0940856659..0c253033ac52 100644 --- a/apps/server/src/http.test.ts +++ b/apps/server/src/http.test.ts @@ -324,6 +324,19 @@ describe("assetResponseHeaders", () => { "X-Content-Type-Options": "nosniff", }); }); + it("serves inline attachment documents with their declared mime type", () => { + expect( + assetResponseHeaders("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/attachments/upload.bin", { mimeType: "application/pdf" }), + ).toMatchObject({ + "Content-Type": "application/pdf", + }); + expect( + assetResponseHeaders("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/attachments/upload.bin", { mimeType: "text/html" }), + ).toMatchObject({ + "Content-Type": "text/html; charset=utf-8", + "Content-Security-Policy": "sandbox allow-scripts allow-forms allow-popups allow-modals", + }); + }); it("serves HTML assets as utf-8 inside a sandboxed origin", () => { for (const path of ["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/workspace/page.html", "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/workspace/PAGE.HTM", "/tmp/report.html"]) { expect(assetResponseHeaders(path)).toMatchObject({ diff --git a/apps/server/src/http.ts b/apps/server/src/http.ts index 4d5865335a39..3bcca884107c 100644 --- a/apps/server/src/http.ts +++ b/apps/server/src/http.ts @@ -64,6 +64,8 @@ const isSafeDownloadMimeType = (mimeType: string): boolean => !/(?:^text\/html$|\/xml(?:$|-)|\+xml$)/i.test(mimeType.trim().toLowerCase()); const isSafeInlineVideoMimeType = (mimeType: string): boolean => DOWNLOAD_MIME_TYPE_PATTERN.test(mimeType) && mimeType.toLowerCase().startsWith("video/"); +const isSafeInlineDocumentMimeType = (mimeType: string): boolean => + mimeType.toLowerCase() === "application/pdf" || mimeType.toLowerCase() === "text/html"; /** RFC 6266 disposition with an ASCII fallback name plus a UTF-8 `filename*`. */ export function downloadContentDisposition(fileName?: string): string { @@ -93,7 +95,7 @@ export function assetResponseHeaders( }, ): Record { const lowerPath = filePath.toLowerCase(); - const inlineVideoMimeType = options?.mimeType?.split(";", 1)[0]?.trim(); + const inlineMimeType = options?.mimeType?.split(";", 1)[0]?.trim(); return { "Cache-Control": "private, max-age=3600", "X-Content-Type-Options": "nosniff", @@ -106,14 +108,24 @@ export function assetResponseHeaders( ? options.mimeType : "application/octet-stream", } - : inlineVideoMimeType !== undefined && isSafeInlineVideoMimeType(inlineVideoMimeType) - ? { "Content-Type": inlineVideoMimeType } - : lowerPath.endsWith(".html") || lowerPath.endsWith(".htm") + : inlineMimeType !== undefined && isSafeInlineVideoMimeType(inlineMimeType) + ? { "Content-Type": inlineMimeType } + : inlineMimeType !== undefined && isSafeInlineDocumentMimeType(inlineMimeType) ? { - "Content-Type": "text/html; charset=utf-8", - "Content-Security-Policy": HTML_CONTENT_SECURITY_POLICY, + "Content-Type": + inlineMimeType.toLowerCase() === "text/html" + ? "text/html; charset=utf-8" + : "application/pdf", + ...(inlineMimeType.toLowerCase() === "text/html" + ? { "Content-Security-Policy": HTML_CONTENT_SECURITY_POLICY } + : {}), } - : {}), + : lowerPath.endsWith(".html") || lowerPath.endsWith(".htm") + ? { + "Content-Type": "text/html; charset=utf-8", + "Content-Security-Policy": HTML_CONTENT_SECURITY_POLICY, + } + : {}), ...(!options?.download && lowerPath.endsWith(".svg") ? { "Content-Security-Policy": SVG_CONTENT_SECURITY_POLICY } : {}), diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index cadb8b3028b3..a5335856be93 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -132,6 +132,7 @@ import { DEFAULT_THREAD_TERMINAL_ID, MAX_TERMINALS_PER_GROUP, type ChatMessage, + isBrowserPreviewAttachment, isImageAttachment, videoMimeType, type SessionPhase, @@ -2602,7 +2603,7 @@ function ChatViewContent(props: ChatViewProps) { }); }, []); const serverMessages = activeThread?.messages; - const openFileAttachment = useCallback( + const downloadFileAttachment = useCallback( async (attachment: ChatFileAttachment) => { const connection = readPreparedConnection(environmentId); if (!connection) { @@ -2654,6 +2655,16 @@ function ChatViewContent(props: ChatViewProps) { }, [createAttachmentAssetUrl, environmentId, routeThreadKey], ); + const openFileAttachment = useCallback( + (attachment: ChatFileAttachment) => { + if (isBrowserPreviewAttachment(attachment) && activeThreadRef) { + useRightPanelStore.getState().openAttachment(activeThreadRef, attachment); + return; + } + void downloadFileAttachment(attachment); + }, + [activeThreadRef, downloadFileAttachment], + ); const serverAttachmentIds = useMemo(() => { const attachmentIds = new Set(); for (const message of serverMessages ?? []) { @@ -7351,14 +7362,18 @@ function ChatViewContent(props: ChatViewProps) { /> ) : (renderedRightPanelSurface?.kind === "files" || renderedRightPanelSurface?.kind === "file") && - activeProject && - activeWorkspaceRoot ? ( + ((activeProject && activeWorkspaceRoot) || + (renderedRightPanelSurface.kind === "file" && renderedRightPanelSurface.attachment)) ? ( [] = []; - if (surface.kind === "file") { + if (surface.kind === "file" && surface.attachment === undefined) { items.push({ id: "copy-path", label: "Copy path" }); } const menuPreviewTabId = previewTabIdOf(surface, props.previewSessions); @@ -837,7 +837,9 @@ export function RightPanelTabs(props: RightPanelTabsProps) { const action = await api.contextMenu.show(items, { x: event.clientX, y: event.clientY }); switch (action) { case "copy-path": - if (surface.kind === "file") props.onCopyFilePath(surface.relativePath); + if (surface.kind === "file" && surface.attachment === undefined) { + props.onCopyFilePath(surface.relativePath); + } break; case "toggle-mute": { // menuOverlay repeats the disabled gate above: the desktop tab must diff --git a/apps/web/src/components/chat/MessagesTimeline.test.tsx b/apps/web/src/components/chat/MessagesTimeline.test.tsx index 1044a956bb66..3cd601ce09a3 100644 --- a/apps/web/src/components/chat/MessagesTimeline.test.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.test.tsx @@ -437,7 +437,7 @@ describe("MessagesTimeline", () => { expect(resolveTimelineMinimapInteractiveWidth(40, true)).toBe("22rem"); }); - it("renders generic attachments as download links instead of image previews", () => { + it("gives browser documents separate preview and download controls", () => { const entry = { ...buildUserTimelineEntry("Read the report."), message: { @@ -459,9 +459,9 @@ describe("MessagesTimeline", () => { , ); - expect(markup).toContain( - '', - ); + expect(markup).toContain('aria-label="Preview report.pdf"'); + expect(markup).toContain('aria-label="Download report.pdf"'); + expect(markup).not.toContain('download="report.pdf"'); expect(markup).not.toContain('alt="report.pdf"'); }); @@ -502,7 +502,7 @@ describe("MessagesTimeline", () => { expect(busyMarkup).not.toContain('disabled=""'); expect(busyMarkup).toContain(">Loading…"); }); - it("renders a file download button without creating its URL in advance", () => { + it("renders an ordinary file download button without creating its URL in advance", () => { const entry = { ...buildUserTimelineEntry("Read the report."), message: { @@ -511,8 +511,8 @@ describe("MessagesTimeline", () => { { type: "file" as const, id: "attachment-report-pdf", - name: "report.pdf", - mimeType: "application/pdf", + name: "archive.zip", + mimeType: "application/zip", sizeBytes: 42, }, ], @@ -524,9 +524,9 @@ describe("MessagesTimeline", () => { ); expect(markup).toContain( - ' + + ctx.onFileDownload(file)} + /> + } + > + + + Download {file.name} + + + ); + } + + const content = ( + <> + {fileIdentity} {file.downloadable === false ? null : ( )} ); - return file.previewUrl ? ( + return file.previewUrl && !opensInPreview ? ( ctx.onFileOpen(file)} className="flex min-w-0 cursor-pointer items-center gap-2 rounded-md py-1 text-left text-sm hover:underline focus-visible:outline-none focus-visible:ring-2 focus-visible:ring-inset focus-visible:ring-ring/70" > @@ -1257,7 +1300,11 @@ function UserTimelineRow({ row }: { row: Extract (
- + {attachment.name}
))} diff --git a/apps/web/src/components/files/FilePreviewPanel.test.ts b/apps/web/src/components/files/FilePreviewPanel.test.ts index 3b5295f180eb..5ef590847c4b 100644 --- a/apps/web/src/components/files/FilePreviewPanel.test.ts +++ b/apps/web/src/components/files/FilePreviewPanel.test.ts @@ -5,7 +5,11 @@ import { normalizeFileCommentRange, remapFileCommentAnnotations, } from "./fileCommentAnnotations"; -import { isMarkdownPreviewFile, setMarkdownTaskChecked } from "./filePreviewMode"; +import { + isMarkdownPreviewFile, + setMarkdownTaskChecked, + shouldShowFileExplorer, +} from "./filePreviewMode"; describe("file comment annotations", () => { it("normalizes and formats selected line ranges", () => { @@ -66,6 +70,42 @@ describe("isMarkdownPreviewFile", () => { }); }); +describe("shouldShowFileExplorer", () => { + it("hides the workspace tree for host files and attachments", () => { + expect( + shouldShowFileExplorer({ + relativePath: "/tmp/report.pdf", + explorerOpen: true, + attachmentOpen: false, + }), + ).toBe(false); + expect( + shouldShowFileExplorer({ + relativePath: "report.pdf", + explorerOpen: true, + attachmentOpen: true, + }), + ).toBe(false); + }); + + it("keeps the saved explorer preference for workspace files", () => { + expect( + shouldShowFileExplorer({ + relativePath: "docs/report.pdf", + explorerOpen: true, + attachmentOpen: false, + }), + ).toBe(true); + expect( + shouldShowFileExplorer({ + relativePath: "docs/report.pdf", + explorerOpen: false, + attachmentOpen: false, + }), + ).toBe(false); + }); +}); + describe("setMarkdownTaskChecked", () => { const markdown = "- [ ] First\n- [x] Second\n"; diff --git a/apps/web/src/components/files/FilePreviewPanel.tsx b/apps/web/src/components/files/FilePreviewPanel.tsx index 40c67bfce5e1..0ad18434d7e7 100644 --- a/apps/web/src/components/files/FilePreviewPanel.tsx +++ b/apps/web/src/components/files/FilePreviewPanel.tsx @@ -1,4 +1,5 @@ import type { + ChatFileAttachment, EditorId, EnvironmentId, ResolvedKeybindingsConfig, @@ -23,6 +24,7 @@ import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { isBrowserPreviewFile, openFileInPreview } from "~/browser/openFileInPreview"; import { useAssetUrlRefresh, useAssetUrlState } from "~/assets/assetUrls"; import { OpenInPicker } from "~/components/chat/OpenInPicker"; +import { PierreEntryIcon } from "~/components/chat/PierreEntryIcon"; import { MediaVideoPlayer } from "~/components/media/MediaVideoPlayer"; import { MediaActions, type MediaActionSource } from "~/components/media/MediaActions"; import { useRemoteOpenState } from "~/remoteOpen"; @@ -64,7 +66,11 @@ import { installFileEditorDismissal } from "./fileEditorDismissal"; import { resolveCenteredFileLineScrollTop } from "./fileLineReveal"; import { DiffCommentAnnotation } from "../diffs/DiffCommentAnnotation"; import { projectFileCacheKey, projectFileEditorCacheKey } from "./fileContentRevision"; -import { isMarkdownPreviewFile, setMarkdownTaskChecked } from "./filePreviewMode"; +import { + isMarkdownPreviewFile, + setMarkdownTaskChecked, + shouldShowFileExplorer, +} from "./filePreviewMode"; import { FileSaveCoordinator } from "./fileSaveCoordinator"; import { confirmProjectFileQueryData, @@ -78,6 +84,7 @@ interface FilePreviewPanelProps { cwd: string; projectName: string; relativePath: string | null; + attachment?: ChatFileAttachment; threadRef: ScopedThreadRef; composerDraftTarget: ScopedThreadRef | DraftId; keybindings: ResolvedKeybindingsConfig; @@ -200,6 +207,69 @@ function WorkspaceImagePreview(props: { const isPdfPreviewFile = (path: string): boolean => /\.pdf$/i.test(path.split(/[?#]/, 1)[0] ?? ""); +function BrowserDocumentFrame(props: { + readonly src: string; + readonly title: string; + readonly pdf: boolean; +}) { + const className = "min-h-0 flex-1 border-0 bg-white"; + // The built-in PDF viewer needs an unsandboxed frame; a PDF runs no scripts. + return props.pdf ? ( + // oxlint-disable-next-line react/iframe-missing-sandbox +