repositories / will
will
owned by admin
src/service/worker.ts
Rawimport { mkdir } from 'node:fs/promises'
import { join } from 'node:path'
import { parentPort, workerData } from 'node:worker_threads'
import { migrate, openDatabase } from '../shared/db.ts'
import { systemClock } from '../shared/time.ts'
import { createOps, type Ops, type OpsDirs } from './ops.ts'
/**
* The single SQLite owner (spec: service-only writes, worker thread so the
* main event loop never blocks on db or cold archive scans).
*/
interface WorkerRequest {
id: number
op: string
args: unknown[]
}
interface WorkerReply {
id: number
ok: boolean
result?: unknown
error?: string
}
export interface WorkerNotify {
type: 'notify'
kind: 'audit' | 'chat'
seqs: number[]
}
export async function initWorker(
dataDir: string,
dirs?: Partial<OpsDirs>,
clock = systemClock,
): Promise<Ops> {
const resolved: OpsDirs = {
dataDir,
archiveDir: dirs?.archiveDir ?? join(dataDir, 'archive'),
mediaDir: dirs?.mediaDir ?? join(dataDir, 'media'),
}
await mkdir(resolved.archiveDir, { recursive: true })
await mkdir(resolved.mediaDir, { recursive: true })
const db = openDatabase(join(dataDir, 'will.db'))
migrate(db)
return createOps(db, resolved, clock)
}
// Direct module import in tests returns ops; as a real worker we serve the RPC protocol.
if (parentPort && workerData?.dataDir) {
const port = parentPort
void (async () => {
const ops = await initWorker(
workerData.dataDir as string,
workerData.dirs as Partial<OpsDirs> | undefined,
)
port?.postMessage({ type: 'ready' })
port?.on('message', (msg: WorkerRequest) => {
void (async () => {
const reply: WorkerReply = { id: msg.id, ok: true }
try {
const method = (ops as unknown as Record<string, (...a: unknown[]) => unknown>)[msg.op]
if (typeof method !== 'function') throw new Error(`unknown op ${msg.op}`)
const result = await method.call(ops, ...msg.args)
reply.result = result
if (msg.op === 'ingestAudit') {
const r = result as { insertedSeqs: number[] }
const notify: WorkerNotify = { type: 'notify', kind: 'audit', seqs: r.insertedSeqs }
port?.postMessage(notify)
}
} catch (err) {
reply.ok = false
reply.error = err instanceof Error ? err.message : String(err)
}
port?.postMessage(reply)
})()
})
})()
}