repositories / will
will
owned by admin
src/service/worker-client.ts
Rawimport { existsSync } from 'node:fs'
import { join } from 'node:path'
import { Worker } from 'node:worker_threads'
/**
* Typed RPC over the service worker thread. Method names mirror Ops exactly;
* args and results cross the worker boundary as structured clones.
*/
interface Pending {
resolve: (value: unknown) => void
reject: (err: Error) => void
}
export interface WorkerClient {
call(op: string, ...args: unknown[]): Promise<unknown>
onNotify(listener: (kind: 'audit' | 'chat', seqs: number[]) => void): void
close(timeoutMs: number): Promise<void>
readonly ready: Promise<void>
}
export function startWorkerClient(workerUrl: string, dataDir: string): WorkerClient {
if (!existsSync(workerUrl)) throw new Error(`worker script not found: ${workerUrl}`)
const worker = new Worker(workerUrl, {
workerData: { dataDir },
resourceLimits: { maxOldGenerationSizeMb: 512 },
})
const pending = new Map<number, Pending>()
const notifyListeners: ((kind: 'audit' | 'chat', seqs: number[]) => void)[] = []
let nextId = 1
let readyResolve: () => void
const ready = new Promise<void>((res) => {
readyResolve = res
})
worker.on(
'message',
(msg: {
type?: string
id?: number
ok?: boolean
result?: unknown
error?: string
kind?: 'audit' | 'chat'
seqs?: number[]
}) => {
if (msg.type === 'ready') {
readyResolve()
return
}
if (msg.type === 'notify' && msg.kind) {
for (const l of notifyListeners) l(msg.kind, msg.seqs ?? [])
return
}
if (msg.id !== undefined) {
const p = pending.get(msg.id)
if (!p) return
pending.delete(msg.id)
if (msg.ok) p.resolve(msg.result)
else p.reject(new Error(msg.error ?? 'worker error'))
}
},
)
return {
ready,
onNotify(listener) {
notifyListeners.push(listener)
},
call(op, ...args) {
const id = nextId++
return new Promise((resolve, reject) => {
pending.set(id, { resolve, reject })
worker.postMessage({ id, op, args })
})
},
async close(timeoutMs) {
const exited = new Promise<void>((res) => worker.once('exit', () => res()))
worker.postMessage({ id: 0, op: 'close', args: [] })
const force = setTimeout(() => void worker.terminate(), timeoutMs)
force.unref()
await Promise.race([exited, new Promise((res) => setTimeout(res, timeoutMs).unref())])
clearTimeout(force)
await worker.terminate().catch(() => undefined)
},
}
}
export function workerScriptCandidates(dir: string): string[] {
return [join(dir, 'service-worker.js'), join(dir, 'worker.ts')]
}