Files
orca/src/main/opencode-usage/scanner.ts
T
Jinwoo Hong 631b51f508 perf(usage): run the Claude/Codex/OpenCode usage scans on a worker thread (#21114)
* perf(codex-usage): resume rollout scans at the last parsed byte

Codex rollout files are append-only and grow all day, but any append
changed both mtime and size, so `canReuse` discarded the cached entry and
the scanner re-read the whole file from byte 0 on the Electron main
process. On one real corpus that was 6.59 GB re-read per cycle across
26.63 GB / 21,110 files.

Each parsed file now persists a resume point: the offset just past the
last newline-terminated line, the parse context at that offset (session
id, cwd, model, running totals), a sha256 of the 4 KiB before it, and the
file's dev:ino. A grown file resumes there and merges the appended
rollup into the cached one; anything unproven falls back to a full
reparse — truncation, an in-place rewrite, rotation, a counted tail with
no trailing newline, a legacy copied-session suffix offset, or a file
that must reclaim deferred fork claims. Resume never depends on mtime
equality, so a coarse-mtime filesystem cannot hide an append.

Fixture: a 75,737-byte rollout with a 758-byte append re-read 76,495
bytes before and 8,950 after (the append plus two bounded 4 KiB boundary
windows).

Also bounds the automation-attribution force predicate for both Codex and
Claude: it keyed on `lastScanError`, so a persistently failing scan forced
a fresh full rescan on every single lookup. It now keys on the most recent
scan attempt, which is one forced scan per run regardless of outcome.

* perf(usage): run the Claude/Codex/OpenCode usage scans on a worker thread

The three first-party usage scans walk whole rollout and transcript corpora
and read OpenCode's SQLite synchronously, all on the Electron main process.
They rarely produce a long stall — the JSONL reader streams, so it yields to
the loop between chunks — but they pin the main-process event loop at ~95%
utilization for the scan's whole duration, which is what every IPC message,
timer and window event then queues behind.

Move that work to one lazily-spawned, unref'd worker thread shared by all
three providers, following the OpenCode SQLite scanner precedent (#8864).
Measured on a synthetic 4,000-rollout corpus (25.8 MB cache): a cold scan
drops from 2,147 ms of main-thread time to 31 ms, and a steady-state
incremental scan from 165 ms to 64 ms.

The worker is stateless and the cache crosses the boundary both ways. That
costs ~64 ms of structured clone at this corpus size, against 2,147 ms saved
on the cold path, and it keeps the persisted cache the single source of
truth — a worker-owned copy would need an invalidation protocol and a second
resident copy of the same multi-MB array.

Failure is closed, never a silent empty result: a worker that cannot spawn,
times out, or crash-loops rejects, and the store records the scan error and
keeps the previous projection.

Two clients already carried the same FIFO/timeout/crash-cap machinery, so
extract it once as WorkerThreadRequestQueue (with the packaged entry-path
resolver as worker-thread-entry-path) and move all three onto it, rather
than adding a third copy. Their existing tests pass unchanged.

The oracle is event-loop utilization on the calling thread, not a stopwatch:
usage-scan-worker-event-loop.test.ts runs the same scan both ways and asserts
the worker leg leaves the caller idle while the main-thread leg does not, so
CI load moves both legs together (#18788).

* test(usage): compare the two scan arms instead of two fixed thresholds

The event-loop oracle claimed to be self-calibrating — its header said "the
ratio is self-calibrating, so CI load moves both legs together (#18788)
instead of tipping a fixed millisecond threshold." It computed no ratio. Two
separate `it()` blocks each asserted an absolute threshold against its own
arm, run separately, so load moved them independently. The comment described
a test nobody wrote, and the flake it promised was impossible is the one that
landed: `activeRatio > 0.8` on the calling-thread arm measured 0.764 on an
ubuntu runner.

Fixing the comment is not enough, because the fraction is the wrong quantity.
CPU contention drags the calling-thread arm's active/wall fraction *down*
toward the worker's, since the loop parks waiting on a contended libuv pool.
A 4-vCPU Linux container measured that arm at 0.175-0.756 across twenty runs,
idle and loaded — never once above 0.8. Active *milliseconds* move the other
way: contention stretches the caller's JS time far more than it stretches the
worker arm's fixed post-and-deserialize cost, so the gap widens under load.

Merge the two arms into one case over one corpus and assert the worker arm
costs the caller under a fifth of the inline arm's active milliseconds. Same
twenty Linux runs: 10.9x-83.6x, passing throughout. Keep the presence
preconditions on both arms — an arm that silently scanned nothing satisfies
the comparison trivially — and extend them to the calling-thread arm, which
previously checked only file and session counts.

* fix(ports): name the dropped command when the probe queue is full

The shared-queue extraction turned `Port scan command queue is full; dropped
${command}.` into a constant string, because `describeFull` was given no way
to see the request. Pile-up is per-probe, so the name is the only thing in
that log that identifies which of lsof/ps/netstat was shed.

Pass the rejected request to `describeFull` and restore the name. The request
is built before the cap check so it exists to be named; the id it burns is a
correlation token, so a gap costs nothing.

The existing overflow test asserted only the error class, which is why the
regression escaped a 29-test suite. It now dispatches the overflow under a
different command than the accepted ones and asserts the message text, so a
message that names the wrong request fails too.

Also add a direct WorkerThreadRequestQueue test. Three subsystems share the
queue and each client test only sees the parts its own protocol exercises,
with `queueCap` reachable from port-scan alone. Covers one-at-a-time FIFO
dispatch, the deadline starting at dispatch rather than enqueue, the
consecutive-death cap, and both points where that count clears.

And record the child-process hazard at the usage worker entry. `terminate()`
reaps nothing the thread spawned, and OpenCode discovery reaches a fork
today: `wslGated*` forks the WSL transcript sidecar for a `\\wsl$\...` path,
which a Windows `OPENCODE_DB` or `XDG_DATA_HOME` can be. One scan through
that entry with a UNC `OPENCODE_DB` forked a sidecar that outlived
`terminate()`.

* test(ai-vault): assert the OpenCode worker messages exactly, not by fragment

Checked every message string in the two clients the shared-queue extraction
rewrote against origin/main. Only the port-scan queue-full one regressed
(fixed in the previous commit); the OpenCode SQLite client's four messages
render identically, the remaining source diffs being renames — `error.message`
to `lastError`, `call.timeoutMs` and `CALL_DEADLINE_MS` to `timeoutMs`.
`session-scanner-worker-client.ts` was not touched by the extraction.

But its suite could not have caught it either. `/timed out/`, `/exited with
code/` and a bare `rejects.toThrow()` all still match a message that has lost
its interpolated value, which is the same blind spot that let the port-scan
regression through. Assert the rendered text instead: the timeout names its
deadline, the exit names its code, and the crash-loop drain still carries the
text of the fault that killed the run.

* fix(usage): correct the worker entry's child-process note

The previous note said `worker.terminate()` leaves a forked sidecar orphaned.
It does not, and the reproduction that appeared to show it used a stub sidecar
missing the `process.on('disconnect', () => process.exit(0))` the real entry
has. With a faithful one: the sidecar lives exactly as long as the thread and
is gone within 2s of `terminate()`, because tearing the thread down closes the
IPC channel it owned. Two worker lifecycles forked two sidecars and leaked
neither, and the pre-worker main-thread path reaps its sidecar the same way,
on host exit.

What is true and worth recording: a fork is reachable from this bundle at all,
which is easy to miss; it survives only as long as the channel does; and the
sidecar is now re-forked per worker lifecycle instead of pooled for the app's
life. State those, and warn that a future child which does not exit on channel
close would not get the same free cleanup.

* fix(usage): kill a wedged scan worker on no progress, not on wall clock

`USAGE_SCAN_TIMEOUT_MS` was a 10-minute deadline on the whole scan. A cold
scan of a real history is legitimately minutes — 637 s measured on a 30 GB
corpus with 300 worktrees before the per-cwd memo, ~51 s after — so a
larger corpus or a slower disk crosses it. Crossing it killed the worker,
recorded a scan error and left the cache unadvanced, so the next refresh
started cold and died at the same point, forever.

The deadline is now a no-progress window. The worker posts a file counter
as it walks the corpus (`UsageScanWorkerProgress`, rate-limited to one
message a second), and `WorkerThreadRequestQueue` re-arms the active
call's timer on each one via the new optional `isProgress`. Clients that
do not pass it keep the plain wall-clock deadline. `MAX_CONSECUTIVE_DEATHS`
and idle teardown are unchanged.

* refactor(usage): report scan progress as a file count, not one call per file

Claude's scanner walks batches, so a per-file callback made it loop just
to bump a counter.
2026-09-16 23:03:37 -04:00

227 lines
8.1 KiB
TypeScript

import { yieldToEventLoop } from '../../shared/event-loop-yield'
import Database from '../sqlite/sync-database'
import { createUsageEventAggregation } from '../usage/usage-event-aggregation'
import {
createUsageWorktreeResolver,
type UsageWorktreeResolver
} from '../usage/usage-worktree-resolver'
import {
compareOpenCodeClaimPriority,
getProcessedDatabaseInfo,
listOpenCodeDatabases
} from './opencode-database-discovery'
import { parseOpenCodeUsageRow } from './opencode-usage-row-parsing'
import { selectUsageRows } from './opencode-usage-row-queries'
import {
attributeOpenCodeUsageEvent,
type OpenCodeUsageWorktreeRef
} from './opencode-usage-worktree-attribution'
import type {
OpenCodeUsageAttributedEvent,
OpenCodeUsageDailyAggregate,
OpenCodeUsagePersistedDatabase,
OpenCodeUsageSession
} from './types'
const YIELD_EVERY_DATABASES = 2
function addCost(left: number | null, right: number | null): number | null {
if (left === null && right === null) {
return null
}
return (left ?? 0) + (right ?? 0)
}
type OpenCodeUsageMetric = { estimatedCostUsd: number | null }
const openCodeUsageAggregation = createUsageEventAggregation<
OpenCodeUsageAttributedEvent,
OpenCodeUsageMetric
>({
metric: {
empty: () => ({ estimatedCostUsd: null }),
fromEvent: (event) => ({ estimatedCostUsd: event.estimatedCostUsd }),
fold: (target, source) => {
target.estimatedCostUsd = addCost(target.estimatedCostUsd, source.estimatedCostUsd)
}
},
cloneSessionForMerge: (session) => structuredClone(session)
})
const { finalizeSessions, mergeSessions, mergeDailyAggregates, sortDailyAggregates } =
openCodeUsageAggregation
export async function parseOpenCodeUsageDatabase(
dbPath: string,
resolveWorktree: UsageWorktreeResolver,
options: { claimSession?: (sessionId: string) => boolean } = {}
): Promise<OpenCodeUsagePersistedDatabase> {
const processedDatabase = await getProcessedDatabaseInfo(dbPath)
const db = new Database(dbPath, { readonly: true, fileMustExist: true })
try {
db.pragma('query_only = ON')
const events: OpenCodeUsageAttributedEvent[] = []
const claimedBySessionId = new Map<string, boolean>()
let hasDeferredClaims = false
for (const row of selectUsageRows(db)) {
const parsed = parseOpenCodeUsageRow(row)
if (!parsed) {
continue
}
// Why: a stale sibling copy of opencode.db carries the same sessions, so
// each session must be counted from exactly one database (#8006).
let owned = claimedBySessionId.get(parsed.sessionId)
if (owned === undefined) {
owned = options.claimSession ? options.claimSession(parsed.sessionId) : true
claimedBySessionId.set(parsed.sessionId, owned)
}
if (!owned) {
hasDeferredClaims = true
continue
}
const attributed = await attributeOpenCodeUsageEvent(parsed, resolveWorktree)
if (attributed) {
events.push(attributed)
}
}
return {
...processedDatabase,
...openCodeUsageAggregation.aggregate(events),
ownedSessionIds: [...claimedBySessionId.entries()]
.filter(([, owned]) => owned)
.map(([sessionId]) => sessionId),
hasDeferredClaims
}
} finally {
db.close()
}
}
export async function scanOpenCodeUsageDatabases(
worktrees: OpenCodeUsageWorktreeRef[],
previousProcessedDatabases: OpenCodeUsagePersistedDatabase[],
onFilesScanned?: (count: number) => void
): Promise<{
processedDatabases: OpenCodeUsagePersistedDatabase[]
sessions: OpenCodeUsageSession[]
dailyAggregates: OpenCodeUsageDailyAggregate[]
}> {
const dbPaths = await listOpenCodeDatabases()
const previousByPath = new Map(
previousProcessedDatabases.map((database) => [database.path, database])
)
// Why: one resolver for the whole scan so every database shares the per-cwd memo.
const resolveWorktree = await createUsageWorktreeResolver(worktrees)
const currentPaths = new Set(dbPaths)
// Why: when a database that owned sessions is deleted, remaining siblings
// still contain those sessions but their caches record them as unowned.
// Only databases that previously deferred claims can reclaim.
const lostOwnerPath = previousProcessedDatabases.some(
(database) =>
!currentPaths.has(database.path) &&
Array.isArray(database.ownedSessionIds) &&
database.ownedSessionIds.length > 0
)
const reusedByPath = new Map<string, OpenCodeUsagePersistedDatabase>()
const pathsToParse: string[] = []
for (const dbPath of dbPaths) {
const databaseInfo = await getProcessedDatabaseInfo(dbPath)
const previous = previousByPath.get(dbPath)
// When an owner disappears, only deferred-claim databases need reparse.
const mustReclaimDeferred = lostOwnerPath && previous?.hasDeferredClaims !== false
const canReuse =
!mustReclaimDeferred &&
previous &&
previous.mtimeMs === databaseInfo.mtimeMs &&
previous.size === databaseInfo.size &&
Array.isArray(previous.ownedSessionIds) &&
typeof previous.hasDeferredClaims === 'boolean'
if (canReuse) {
reusedByPath.set(dbPath, previous)
} else {
pathsToParse.push(dbPath)
}
}
// Why: a sticky backup claim from a scan where opencode.db was missing would
// otherwise freeze a still-growing session at the backup snapshot when the
// live db reappears. Reparse a lower-priority sibling only when it still owns
// sessions a higher-priority path could reclaim; a sibling that owns nothing
// (the common case once the live db has claimed every shared session) has no
// claim to give back, so reparsing it every time the live db changes is pure
// work.
const demotedReusePaths: string[] = []
for (const [dbPath, reused] of reusedByPath) {
if ((reused.ownedSessionIds?.length ?? 0) === 0) {
continue
}
const higherPriorityParsing = pathsToParse.some(
(candidate) => compareOpenCodeClaimPriority(candidate, dbPath) < 0
)
if (higherPriorityParsing) {
demotedReusePaths.push(dbPath)
}
}
for (const dbPath of demotedReusePaths) {
reusedByPath.delete(dbPath)
pathsToParse.push(dbPath)
}
// Why: `opencode-*.db` siblings are typically stale copies of `opencode.db`
// (backups), so mergeSessions would double every duplicated session (#8006).
// Each session is counted from exactly one database. The canonical live db
// claims first so a stale backup cannot freeze a still-growing session at
// its snapshot totals; cached databases keep the claims they persisted.
const sessionOwnerById = new Map<string, string>()
for (const dbPath of [...reusedByPath.keys()].sort(compareOpenCodeClaimPriority)) {
const previous = reusedByPath.get(dbPath)
for (const sessionId of previous?.ownedSessionIds ?? []) {
if (!sessionOwnerById.has(sessionId)) {
sessionOwnerById.set(sessionId, dbPath)
}
}
}
const parsedByPath = new Map<string, OpenCodeUsagePersistedDatabase>()
const orderedPathsToParse = [...pathsToParse].sort(compareOpenCodeClaimPriority)
for (const [index, dbPath] of orderedPathsToParse.entries()) {
const processed = await parseOpenCodeUsageDatabase(dbPath, resolveWorktree, {
claimSession: (sessionId) => {
const owner = sessionOwnerById.get(sessionId)
if (owner !== undefined && owner !== dbPath) {
return false
}
sessionOwnerById.set(sessionId, dbPath)
return true
}
})
parsedByPath.set(dbPath, processed)
onFilesScanned?.(1)
if ((index + 1) % YIELD_EVERY_DATABASES === 0) {
await yieldToEventLoop()
}
}
const processedDatabases: OpenCodeUsagePersistedDatabase[] = []
const sessionsById = new Map<string, OpenCodeUsageSession>()
const dailyByKey = new Map<string, OpenCodeUsageDailyAggregate>()
for (const dbPath of dbPaths) {
const processed = reusedByPath.get(dbPath) ?? parsedByPath.get(dbPath)
if (!processed) {
continue
}
processedDatabases.push(processed)
mergeSessions(sessionsById, processed.sessions)
mergeDailyAggregates(dailyByKey, processed.dailyAggregates)
}
return {
processedDatabases,
sessions: finalizeSessions(sessionsById),
dailyAggregates: sortDailyAggregates(dailyByKey)
}
}