Files
orca/src/main/ai-vault/session-scanner-opencode-sqlite-worker-client.ts
T
Jinwoo Hong 65631e449a refactor(ai-vault): split the session scanner into a transcript reader and consumers (#19666)
* refactor(ai-vault): split the session scanner into a transcript reader and consumers

The scanner's only output was the Session History summary; a second reader
of the same transcripts (a search index) had nowhere to plug in without
hooking the parse itself. Extract a reader that owns each file read, keeps
the resumable cursor and publishes every decoded message to registered
consumers. The session list stays a fold inside the parser and is the only
consumer here. Parsers take an optional message sink instead of a scope.

Also: read Cursor chats/<md5>/<uuid>/meta.json for cwd, title and
timestamps (Cursor transcripts carry only role and message); share the
lazily spawned worker-thread host between the OpenCode SQLite reader and
the port-scan probe; probe OpenCode's schema before querying; keep the
newest-N discovery set with a bounded insert instead of sort+slice.

Session list output is byte-identical to main across all 18 providers
cold and append-resumed; the one Cursor session gains cwd/timestamps
from meta.json.

* fix(ai-vault): serialize per-path parses and report unpublished reads

Overlapping parses of one transcript share the cached resume point's
message channel, so the second beginRead dropped the first read's
consumers and the first finishRead handed them the wrong outcome. Two
callers really do overlap: a forced refresh restarts a scan while the
aborted scan's parse is still in flight, and the title reader parses
outside any scan. Restore the per-path lane around the whole
lookup-read-store sequence.

OpenCode's SQLite sessions are decoded on a worker thread the channel
cannot reach, so their reads published no messages while reporting a
complete span. Finish those reads as incomplete instead, so a consumer
never records a cursor for a stream it did not receive.

* fix(ai-vault): degrade a refused cursor chats read instead of dropping sessions

A refused WSL read of Cursor's chats tree rethrew, and the per-file catch
in discovery then recorded an issue and skipped the transcript. Before
the meta.json join Cursor had no content dependency, so a stalled distro
could not hide a Cursor session at all. Degrade to no metadata for the
scan and report the chats root once. The parse cache stays honest without
the throw: discovery stats no meta.json on a refused scan, so the entry's
recorded size omits it and the next healthy scan re-reads the transcript.

The per-scan index scope covered discovery only, so every Cursor finalize
re-read the chats root to validate the module cache. Move the scope to
scanAiVaultSessions, which spans discovery and parse.

Also drop the unused signal parameters the sink threading added to the
Devin and Hermes content parsers, by giving each file parser a private
record parser instead.

* fix(ai-vault): do not cache a cursor parse whose meta.json read was refused

Discovery stats meta.json into the candidate's cache key, so when only
the meta.json read is refused the un-enriched session was stored under a
key that looks unchanged and reuseCachedSession never re-ran the enrich
hook. The session stayed without cwd until Cursor rewrote the file. The
enrich hook now reports 'refused', the resumable state exposes
isCacheable, and the parse cache drops the entry instead of storing it,
so the next healthy scan re-parses. The index-read branch is unaffected:
it never stats meta.json, so its key is honest already.

* fix(ai-vault): separate the transcript's size from its cache key

sizeBytes folds a content dependency's size in, so it is a cache key
rather than a file length. The reader compared a transcript byte offset
against it and reported it as a whole-file read offset, which for Cline
handed consumers an offset past the end of the file it read. Carry the
dependency's own size on FileWithMtime and subtract it in the reader.

A refused sibling stat rethrew, so discovery recorded an issue and
skipped the transcript, the same drop removed for the readdir and read
paths. Degrade to no dependency, note the tree once, and mark the key
untrustworthy.

An untrustworthy key no longer costs the resume cursor: the entry is
stored under an mtime no stat can produce, so unchanged is false while
the resume point survives and the next scan resumes instead of re-reading
the whole transcript.

* test(ai-vault): pin the untrustworthy-key mechanism, not just its effect

Both refusal tests asserted that a later healthy scan re-enriches, which
a plain store would also satisfy once the resume cursor was preserved.
Assert the cache entry directly: its mtime is the unmatchable sentinel
and its resume point survives. The sentinel is exported so the tests name
the contract instead of repeating -1.

* refactor(ai-vault): track a session's sidecar file apart from its transcript

Folding Cursor's meta.json stat into the transcript's mtime/size made one
key mean two things, and every round of review found another consequence:
a byte offset could not be compared against it, a refused sibling read
took the transcript down with it, and an un-enriched parse cached under
it looked current forever. Main already had the answer for a file the
transcript key cannot see: Codex titles are refreshed at reuse time over
the cached session, not folded into the key.

Discovery now records the sibling as its own observation, unknown when it
could not be read. A cache hit needs both the transcript key and the
sidecar to match. When only the sidecar moved, Cursor re-merges it over
the stored un-enriched fold result and never re-reads the transcript;
Cline, which reads its sibling as part of the parse, re-parses.

Merging over the fold result rather than the accumulator makes enrichment
pure, so a meta.json rewritten with a new cwd replaces the old one
instead of losing to it. That was unreachable while the merge used ??= on
a session it had already enriched.

Cline and the remote scanner move to the same field, so the fold is gone
from both discovery paths.

* fix(ai-vault): tell an absent sidecar from an unreadable one

Three places collapsed the two. sidecarUnchanged returned true for any
observed 'none' without reading the entry, so a sidecar that was deleted,
or one that was unreadable last scan, both read as cache hits. Native
discovery mapped every non-WSL stat failure to 'none', so an EACCES on
meta.json left a session enriched from a file nobody can see, with no
scan issue. Remote discovery could not tell a missing sibling from a
failed stat, because statRemoteSessionFile returns null for both.

'none' is now a claim: absent-now is a hit only when it was absent before
or the agent never had a sidecar, and only ENOENT/ENOTDIR reads as
absent. statRemoteSessionFile grows an opt-in rethrow so its caller can
distinguish the two failures it already reports.

Also rewrites three comments in the cursor chat-meta reader that still
described the deleted fold.
2026-09-09 03:36:41 -04:00

282 lines
10 KiB
TypeScript

import { LazyWorkerThreadHost, type WorkerThreadFactory } from '../lazy-worker-thread-host'
import type { AiVaultScanIssue, AiVaultSession } from '../../shared/ai-vault-types'
import type {
OpenCodeSqliteListRequest,
OpenCodeSqliteListValue,
OpenCodeSqliteParseRequest,
OpenCodeSqliteWorkerRequest,
OpenCodeSqliteWorkerResponse
} from './session-scanner-opencode-sqlite-worker-protocol'
import type { SessionFileCandidate } from './session-scanner-types'
import { errorMessage } from './session-scanner-values'
// Why (#8864): a lazily-spawned, unref'd worker runs OpenCode SQLite reads off
// the main-process event loop. This module owns the request half (FIFO
// one-at-a-time dispatch, per-call timeouts, respawn-on-fault); the thread's
// own lifetime belongs to LazyWorkerThreadHost, shared with the port-scan probe
// client. The default spawn + shared singleton live in
// session-scanner-opencode-sqlite-worker-spawn.ts.
export const LIST_TIMEOUT_MS = 30_000
export const PARSE_TIMEOUT_MS = 15_000
export const IDLE_TEARDOWN_MS = 30_000
// After this many consecutive worker deaths, fail the remaining queued calls to
// scan issues instead of respawning so a DB that reliably kills the worker can't
// spin a crash loop. Reset on any successful response, after draining, and when a
// fresh scan burst starts from idle (so the cap is per-scan, not process-wide).
export const MAX_CONSECUTIVE_DEATHS = 3
// Omit<union, 'id'> collapses to the shared keys, so omit each member and let
// the client stamp the correlation id.
type OpenCodeSqliteRequestBody =
| Omit<OpenCodeSqliteListRequest, 'id'>
| Omit<OpenCodeSqliteParseRequest, 'id'>
type PendingCall = {
request: OpenCodeSqliteWorkerRequest
timeoutMs: number
resolve: (value: unknown) => void
reject: (error: Error) => void
timer: NodeJS.Timeout | null
}
// Distinguishes "no worker available at all" from a timeout or crash so callers
// can surface a precise issue while keeping synchronous SQLite off the main thread.
class OpenCodeSqliteWorkerUnavailableError extends Error {}
/**
* Main-thread bridge that runs OpenCode SQLite reads on a persistent worker
* thread. Dispatches one request at a time (FIFO), times each request out from
* dispatch, respawns after faults (capped by `MAX_CONSECUTIVE_DEATHS`), tears
* the worker down after `IDLE_TEARDOWN_MS` of inactivity, and fails closed when
* no worker can be spawned rather than moving SQLite work onto the main thread.
*/
export class OpenCodeSqliteWorkerClient {
private active: PendingCall | null = null
private queue: PendingCall[] = []
private consecutiveDeaths = 0
private nextId = 1
private readonly host: LazyWorkerThreadHost<OpenCodeSqliteWorkerResponse>
constructor(options: { workerFactory: WorkerThreadFactory; log?: (message: string) => void }) {
const log = options.log ?? ((message: string) => console.warn(message))
this.host = new LazyWorkerThreadHost<OpenCodeSqliteWorkerResponse>({
factory: options.workerFactory,
idleTeardownMs: IDLE_TEARDOWN_MS,
onMessage: (response) => this.onMessage(response),
onError: (error) => this.onWorkerFault(error),
onExit: (code) => this.onWorkerExit(code),
isIdle: () => !this.active && this.queue.length === 0,
// Why (#8864): never fall back to synchronous SQLite reads here; a missing
// bundle or resource-exhausted spawn must omit OpenCode history rather than
// reintroduce the main-process hang this worker boundary prevents.
onUnavailable: (err) =>
log(`OpenCode SQLite worker unavailable; skipping its history. ${errorMessage(err)}`)
})
}
/**
* List session candidates from the given OpenCode databases on the worker.
* @param args.dbPaths - Absolute paths to opencode.db files to scan.
* @param args.limit - Maximum number of sessions to return per database.
* @param args.issues - Collected scan issues (worker issues are merged in).
* @returns Synthetic candidates sorted by effective recency; empty (with a
* scan issue) when the worker is unavailable, times out, or crashes.
*/
async list(args: {
dbPaths: readonly string[]
limit: number
issues: AiVaultScanIssue[]
}): Promise<SessionFileCandidate[]> {
if (args.dbPaths.length === 0) {
return []
}
try {
const value = (await this.dispatch(
{ kind: 'list', dbPaths: args.dbPaths, limit: args.limit },
LIST_TIMEOUT_MS
)) as OpenCodeSqliteListValue
args.issues.push(...value.issues)
return value.candidates
} catch (err) {
if (err instanceof OpenCodeSqliteWorkerUnavailableError) {
// Kinded: a whole source failed, not a transcript.
args.issues.push({
agent: 'opencode',
kind: 'scope',
path: args.dbPaths[0] ?? 'opencode.db',
message:
'OpenCode history was skipped because its background scanner could not start; the app remains responsive.'
})
return []
}
// Timeout/crash: this storage dir's SQLite DBs contribute no sessions this
// scan, surfaced as one scan issue rather than an unbounded stall.
args.issues.push({
agent: 'opencode',
kind: 'scope',
path: args.dbPaths[0] ?? 'opencode.db',
message: `OpenCode history scan did not complete: ${errorMessage(err)}`
})
return []
}
}
/**
* Parse a single OpenCode session on the worker.
* @param args.dbPath - Absolute path to the opencode.db file.
* @param args.sessionId - Primary key in the `session` table.
* @param args.platform - Platform used for resume-command generation.
* @returns The parsed session, or `null` when it does not exist; rejects on
* worker timeout/crash so the scanner records a per-session scan issue.
*/
async parse(args: {
dbPath: string
sessionId: string
platform: NodeJS.Platform
}): Promise<AiVaultSession | null> {
try {
const value = await this.dispatch(
{ kind: 'parse', dbPath: args.dbPath, sessionId: args.sessionId, platform: args.platform },
PARSE_TIMEOUT_MS
)
return value as AiVaultSession | null
} catch (err) {
if (err instanceof OpenCodeSqliteWorkerUnavailableError) {
throw new Error('OpenCode SQLite background scanner could not start.')
}
// Reject only this session; the scanner turns the throw into a scan issue.
throw err instanceof Error ? err : new Error(String(err))
}
}
private dispatch(request: OpenCodeSqliteRequestBody, timeoutMs: number): Promise<unknown> {
return new Promise((resolve, reject) => {
const id = this.nextId++
// A fresh burst from full idle starts a new scan: clear any death count
// carried from a prior scan so the respawn cap can't drain this scan early.
if (!this.active && this.queue.length === 0) {
this.consecutiveDeaths = 0
}
this.queue.push({
request: { ...request, id } as OpenCodeSqliteWorkerRequest,
timeoutMs,
resolve,
reject,
timer: null
})
this.pump()
})
}
private pump(): void {
if (this.active || this.queue.length === 0) {
return
}
const worker = this.host.ensure()
if (!worker) {
this.failQueuedAsUnavailable()
return
}
const call = this.queue.shift()
if (!call) {
return
}
this.active = call
this.host.clearIdleTimer()
// Timeout clock starts at dispatch (not enqueue): a batch may enqueue up to
// 8 parses at once, and a queue-inclusive timeout would fire falsely.
call.timer = setTimeout(() => this.onTimeout(call), call.timeoutMs)
call.timer.unref?.()
worker.postMessage(call.request)
}
private onMessage(response: OpenCodeSqliteWorkerResponse): void {
const call = this.active
if (!call || call.request.id !== response.id) {
return
}
this.consecutiveDeaths = 0
if (response.ok) {
this.settle(call, () => call.resolve(response.value))
} else {
this.settle(call, () => call.reject(new Error(response.error)))
}
this.afterSettle()
}
private onTimeout(call: PendingCall): void {
if (this.active !== call) {
return
}
this.onWorkerFault(new Error(`OpenCode SQLite worker timed out after ${call.timeoutMs}ms`))
}
private onWorkerExit(code: number): void {
// A clean self-exit is not a death, but the stale handle must be dropped
// or the next dispatch would post into the dead worker and stall to timeout.
if (code === 0 && !this.active && this.queue.length === 0) {
this.host.destroy()
return
}
this.onWorkerFault(new Error(`OpenCode SQLite worker exited with code ${code}`))
}
private onWorkerFault(error: Error): void {
const failed = this.active
this.host.destroy()
this.consecutiveDeaths++
if (failed) {
this.settle(failed, () => failed.reject(error))
}
if (this.consecutiveDeaths >= MAX_CONSECUTIVE_DEATHS) {
this.drainQueueAfterCrashLoop(error)
return
}
if (this.queue.length > 0) {
this.pump()
}
}
private drainQueueAfterCrashLoop(error: Error): void {
const pending = this.queue
this.queue = []
this.consecutiveDeaths = 0
const drainError = new Error(
`OpenCode SQLite worker crashed repeatedly; skipping remaining sessions (${error.message})`
)
for (const call of pending) {
this.settle(call, () => call.reject(drainError))
}
}
private failQueuedAsUnavailable(): void {
const pending = this.queue
this.queue = []
for (const call of pending) {
this.settle(call, () =>
call.reject(new OpenCodeSqliteWorkerUnavailableError('worker spawn failed'))
)
}
}
private settle(call: PendingCall, run: () => void): void {
if (call.timer) {
clearTimeout(call.timer)
call.timer = null
}
if (this.active === call) {
this.active = null
}
run()
}
private afterSettle(): void {
if (this.queue.length > 0) {
this.pump()
} else {
this.host.scheduleIdleTeardown()
}
}
}