mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 00:02:39 +00:00
* refactor(sqlite): rename the OpenCode SQLite worker entry to foreign-sqlite-reader (STA-9122) The worker thread that reads OpenCode's database off the main thread is about to read other apps' databases too, so its entry is renamed to what it is: src/main/foreign-sqlite-readers/foreign-sqlite-reader-entry.ts, built as out/main/foreign-sqlite-reader-entry.js. Why now: #24572 fixed the Cursor focus freeze with a second, dedicated worker. Rather than grow one worker per foreign app, the next commit moves Cursor onto this entry and deletes that worker. This commit is the rename only; #24572's cursor-desktop-profile-worker-entry lines stay until then. It moves out of ai-vault/ into a new foreign-sqlite-readers/ module because it will no longer be session-scanner code; the module will own the readers, their dispatch, protocol and main-process client. The entry still routes only OpenCode kinds in this commit. The OpenCode dispatch, protocol and process entry stay in ai-vault/ and stay OpenCode-only, because the SSH/WSL relay reader bundles them (build-relay.mjs). Every reference is updated: electron.vite.config.ts input key, knip entry, the plain-node entry guard and its test, the asarUnpack list (the scanner service still spawns this entry under ELECTRON_RUN_AS_NODE), and the electron-builder test that reads the filename. The filename and the beside-or-one-up (Rollup chunks) lookup now live in foreign-sqlite-reader-entry-path.ts, which the OpenCode spawn reuses, plus an Electron-main resolver that uses the packaged app.asar path. * fix(cursor): move the desktop-login read from its dedicated worker onto the foreign SQLite reader (STA-9122) #24572 fixed the Cursor focus freeze (#24360) with a dedicated worker (rate-limits/cursor-desktop-profile-worker*.ts). Orca already runs OpenCode's database reads on a worker, and more foreign-app SQLite reads are coming, so keeping one worker per app means one entry, build input, asarUnpack line, knip entry and guard line each. This keeps one pattern instead: Cursor's state.vscdb read runs on the shared foreign SQLite reader entry, and the dedicated worker, its entry and its config lines are deleted. What moves: - The read itself is a pure cursorProfile reader in foreign-sqlite-readers/readers/ (was rate-limits/cursor-desktop-state-db.ts), run only on the worker. A separate dispatch owns the new kinds and refuses an unknown kind. The entry routes OpenCode kinds to the untouched OpenCode dispatch, so the relay's OpenCode reader stays byte-identical. - ForeignSqliteReaderClient gives each reader its own WorkerThreadRequestQueue lane (own lazily started, idle-torn-down thread; one shared factory) with in-flight dedupe per database path. Any failure resolves to the reader's existing failure value and never falls back to the main thread. Kept from #24572, so every reader gets them: - Await worker retirement before respawning. Worker.terminate() cannot interrupt a native SQLite call (e.g. a WAL-index rebuild), so the old thread lives on until that call returns; respawning at once stacked a new thread on the same work for every timed-out read (#24572 measured three live workers). This belongs in the shared host, which fire-and-forgot terminate(): LazyWorkerThreadHost now takes awaitRetirement and refuses to spawn until the terminated worker settles, and the queue fails calls closed meanwhile. Opt-in, because pure-JS clients (session scanner abort, port scan) respawn right after an abort. A rejected terminate() also ends retirement, so it cannot latch the reader off (raised in #24572's review). - 10 s Cursor timeout, 2 consecutive deaths, a queue cap of 8. - #24572's worker tests, rewritten against the shared client: responsive caller plus coalesced probes, unavailable worker without path leaks, stalled-worker recovery, no respawn before retirement, dispose settles. Tests: reader, dispatch, client (timeout, 10 s default, crash, malformed, unavailable without a main-thread read, dedupe, queue cap, own thread per reader), queue retirement (stalled and rejected terminate), import boundary, and an event-loop test reading a ~50 MB WAL with no -shm on a real worker. * test(sqlite): walk the reader import boundary with the shared source-tree scan (STA-9122) * fix(sqlite): key reader dedupe on a caller-supplied key, not the path alone (STA-9122)
281 lines
9.3 KiB
TypeScript
281 lines
9.3 KiB
TypeScript
import { LazyWorkerThreadHost, type WorkerThreadFactory } from './lazy-worker-thread-host'
|
|
|
|
/**
|
|
* FIFO one-at-a-time request half shared by every main-process worker-thread
|
|
* client: per-call timeout armed at dispatch, respawn-on-fault capped so a
|
|
* payload that reliably kills the worker cannot spin a crash loop, idle
|
|
* teardown, and — the rule that matters — failing queued calls closed instead
|
|
* of moving their work back onto the main thread. `LazyWorkerThreadHost` owns
|
|
* the thread's lifetime; this owns which call a message belongs to.
|
|
*/
|
|
|
|
export type WorkerThreadRequestQueueOptions<TRequest> = {
|
|
factory: WorkerThreadFactory
|
|
idleTeardownMs: number
|
|
/** Consecutive deaths after which the remaining queue is failed rather than respawned. */
|
|
maxConsecutiveDeaths: number
|
|
/**
|
|
* Omit for an unbounded queue; set it where pile-up is itself the bug.
|
|
* `describeFull` gets the rejected request so the message can name the work
|
|
* that was dropped, which is the only detail a log has to identify it.
|
|
*/
|
|
queueCap?: { maxQueuedCalls: number; describeFull: (request: TRequest) => string }
|
|
/**
|
|
* Marks a message as liveness for the active call rather than its result.
|
|
* Omit it and `describeTimeout`'s deadline is a wall clock on the whole call;
|
|
* supply it and the deadline becomes a no-progress window, re-armed by every
|
|
* progress message the active call sends. Work that is slow but still moving
|
|
* must not be killed for being slow.
|
|
*/
|
|
isProgress?: (message: { id: number }) => boolean
|
|
/** The client's own error subclass, so callers can tell "no worker" from a fault. */
|
|
createUnavailableError: (message: string) => Error
|
|
describeTimeout: (timeoutMs: number) => string
|
|
describeExit: (code: number) => string
|
|
describeCrashLoop: (lastError: string) => string
|
|
/** First spawn failure only; a repeating one must not repeat the log. */
|
|
onUnavailable: (error: unknown) => void
|
|
/**
|
|
* Fail calls closed while a terminated worker has yet to exit, instead of
|
|
* spawning beside it. For workers whose native calls delay termination.
|
|
*/
|
|
awaitRetirement?: boolean
|
|
}
|
|
|
|
type PendingCall<TRequest, TResponse> = {
|
|
request: TRequest
|
|
timeoutMs: number
|
|
resolve: (value: TResponse) => void
|
|
reject: (error: Error) => void
|
|
timer: NodeJS.Timeout | null
|
|
signal?: AbortSignal
|
|
cleanupAbort: () => void
|
|
}
|
|
|
|
export class WorkerThreadRequestQueue<
|
|
TRequest extends { id: number },
|
|
TResponse extends { id: number }
|
|
> {
|
|
private active: PendingCall<TRequest, TResponse> | null = null
|
|
private queue: PendingCall<TRequest, TResponse>[] = []
|
|
private consecutiveDeaths = 0
|
|
private nextId = 1
|
|
private disposed = false
|
|
private readonly host: LazyWorkerThreadHost<TResponse>
|
|
|
|
constructor(private readonly options: WorkerThreadRequestQueueOptions<TRequest>) {
|
|
this.host = new LazyWorkerThreadHost<TResponse>({
|
|
factory: options.factory,
|
|
idleTeardownMs: options.idleTeardownMs,
|
|
onMessage: (response) => this.onMessage(response),
|
|
onError: (error) => this.onWorkerFault(error),
|
|
onExit: (code) => this.onWorkerExit(code),
|
|
isIdle: () => !this.active && this.queue.length === 0,
|
|
onUnavailable: options.onUnavailable,
|
|
awaitRetirement: options.awaitRetirement
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Queue one request and resolve with the worker's matching response.
|
|
* @param buildRequest - Builds the request body around the correlation id this queue stamps.
|
|
* @param timeoutMs - Deadline measured from dispatch, not from enqueue.
|
|
* @returns The worker's response; rejects on timeout, crash, or an unspawnable worker.
|
|
*/
|
|
dispatch(
|
|
buildRequest: (id: number) => TRequest,
|
|
timeoutMs: number,
|
|
signal?: AbortSignal
|
|
): Promise<TResponse> {
|
|
return new Promise((resolve, reject) => {
|
|
if (this.disposed) {
|
|
reject(new Error('Worker request queue disposed'))
|
|
return
|
|
}
|
|
if (signal?.aborted) {
|
|
reject(signal.reason ?? new Error('Worker request aborted'))
|
|
return
|
|
}
|
|
// Built before the cap check so a rejection can name the dropped work;
|
|
// the id it burns is only a correlation token, so a gap costs nothing.
|
|
const request = buildRequest(this.nextId++)
|
|
const cap = this.options.queueCap
|
|
if (cap && this.queue.length >= cap.maxQueuedCalls) {
|
|
reject(new Error(cap.describeFull(request)))
|
|
return
|
|
}
|
|
// A fresh burst from full idle starts new work: clear any death count
|
|
// carried from a prior burst so the respawn cap can't drain it early.
|
|
if (!this.active && this.queue.length === 0) {
|
|
this.consecutiveDeaths = 0
|
|
}
|
|
const call: PendingCall<TRequest, TResponse> = {
|
|
request,
|
|
timeoutMs,
|
|
resolve,
|
|
reject,
|
|
timer: null,
|
|
signal,
|
|
cleanupAbort: () => signal?.removeEventListener('abort', abort)
|
|
}
|
|
const abort = (): void => {
|
|
if (this.active === call) {
|
|
this.host.destroy()
|
|
} else {
|
|
this.queue = this.queue.filter((queued) => queued !== call)
|
|
}
|
|
this.settle(call, () => reject(signal?.reason ?? new Error('Worker request aborted')))
|
|
this.afterSettle()
|
|
}
|
|
signal?.addEventListener('abort', abort, { once: true })
|
|
this.queue.push(call)
|
|
this.pump()
|
|
})
|
|
}
|
|
|
|
dispose(): void {
|
|
this.disposed = true
|
|
this.host.destroy()
|
|
const pending = this.active ? [this.active, ...this.queue] : this.queue
|
|
this.queue = []
|
|
for (const call of pending) {
|
|
this.settle(call, () => call.reject(new Error('Worker request queue disposed')))
|
|
}
|
|
}
|
|
|
|
private pump(): void {
|
|
if (this.active || this.queue.length === 0) {
|
|
return
|
|
}
|
|
// A shared signal is already aborted before its remaining listeners run.
|
|
while (this.queue[0]?.signal?.aborted) {
|
|
const cancelled = this.queue.shift()
|
|
if (cancelled) {
|
|
this.settle(cancelled, () =>
|
|
cancelled.reject(cancelled.signal?.reason ?? new Error('Worker request aborted'))
|
|
)
|
|
}
|
|
}
|
|
if (this.queue.length === 0) {
|
|
this.host.scheduleIdleTeardown()
|
|
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()
|
|
this.armDeadline(call)
|
|
try {
|
|
worker.postMessage(call.request)
|
|
} catch (error) {
|
|
this.onWorkerFault(error instanceof Error ? error : new Error(String(error)))
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Clock starts at dispatch, not enqueue: a queue-inclusive deadline would fire
|
|
* falsely on the calls waiting behind a long one. Re-armed on every progress.
|
|
*/
|
|
private armDeadline(call: PendingCall<TRequest, TResponse>): void {
|
|
if (call.timer) {
|
|
clearTimeout(call.timer)
|
|
}
|
|
call.timer = setTimeout(() => this.onTimeout(call), call.timeoutMs)
|
|
call.timer.unref?.()
|
|
}
|
|
|
|
private onMessage(response: TResponse): void {
|
|
const call = this.active
|
|
if (!call || call.request.id !== response.id) {
|
|
return
|
|
}
|
|
// Liveness, not a result: keep waiting, but restart the no-progress window.
|
|
if (this.options.isProgress?.(response)) {
|
|
this.armDeadline(call)
|
|
return
|
|
}
|
|
this.consecutiveDeaths = 0
|
|
this.settle(call, () => call.resolve(response))
|
|
this.afterSettle()
|
|
}
|
|
|
|
private onTimeout(call: PendingCall<TRequest, TResponse>): void {
|
|
if (this.active !== call) {
|
|
return
|
|
}
|
|
this.onWorkerFault(new Error(this.options.describeTimeout(call.timeoutMs)))
|
|
}
|
|
|
|
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(this.options.describeExit(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 >= this.options.maxConsecutiveDeaths) {
|
|
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(this.options.describeCrashLoop(error.message))
|
|
for (const call of pending) {
|
|
this.settle(call, () => call.reject(drainError))
|
|
}
|
|
}
|
|
|
|
private failQueuedAsUnavailable(): void {
|
|
const pending = this.queue
|
|
this.queue = []
|
|
const reason = this.host.isRetiring ? 'previous worker still exiting' : 'worker spawn failed'
|
|
for (const call of pending) {
|
|
this.settle(call, () => call.reject(this.options.createUnavailableError(reason)))
|
|
}
|
|
}
|
|
|
|
private settle(call: PendingCall<TRequest, TResponse>, run: () => void): void {
|
|
call.cleanupAbort()
|
|
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()
|
|
}
|
|
}
|
|
}
|