Luigit
repositories / will

will

owned by admin

src/service/worker-client.ts

Raw
import { 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')]
}