import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; import { randomUUID } from "node:crypto"; import { Buffer } from "node:buffer"; import { existsSync } from "node:fs"; import { basename, delimiter, dirname, join, resolve } from "node:path"; import { encodeProtocolMessage, JsonlDecoder, parseProtocolLine, ProtocolError, type ProtocolEnvelope, type RuntimeCommandInput, type RuntimeMessage, } from "./protocol.js"; export interface RuntimeProcessOptions { runtimePath: string; modelPath: string; protocolVersion: 1; threads?: number; startupTimeoutMs?: number; expectedEngine?: string; minimumEngineAbi?: number; maximumEngineAbi?: number; env?: NodeJS.ProcessEnv; } export type RuntimeMessageHandler = (message: RuntimeMessage) => void; export type RuntimeFailureHandler = (error: Error) => void; function runtimeLibraryDirectories(runtimePath: string): string[] { const runtimeDirectory = resolve(dirname(runtimePath)); const configuration = basename(runtimeDirectory); const candidates = [ join(runtimeDirectory, "lib"), // Visual Studio multi-config builds place DLLs in bin\\ while // the executable remains in . join(runtimeDirectory, "..", "bin", configuration), runtimeDirectory, ]; return [ ...new Set( candidates .map((directory) => resolve(directory)) .filter((directory) => existsSync(directory)), ), ]; } interface PendingResponse { resolve: (message: RuntimeMessage) => void; reject: (error: Error) => void; timer: NodeJS.Timeout; } export class RuntimeProcess { private child: ChildProcessWithoutNullStreams | undefined; private decoder = new JsonlDecoder(); private readonly pending = new Map(); private readonly messageHandlers = new Set(); private readonly failureHandlers = new Set(); private sequence = 0; private stderrBuffer = ""; private recentMessages: RuntimeMessage[] = []; private starting: Promise | undefined; constructor(private readonly options: RuntimeProcessOptions) {} get isAlive(): boolean { return ( this.child !== undefined && this.child.exitCode === null && !this.child.killed ); } get stderr(): string { return this.stderrBuffer; } onMessage(handler: RuntimeMessageHandler): () => void { this.messageHandlers.add(handler); return () => this.messageHandlers.delete(handler); } onFailure(handler: RuntimeFailureHandler): () => void { this.failureHandlers.add(handler); return () => this.failureHandlers.delete(handler); } async start(): Promise { if (this.isAlive) return; if (this.starting) return this.starting; this.starting = this.startInternal() .catch(async (error: unknown) => { await this.shutdown(); throw error; }) .finally(() => { this.starting = undefined; }); return this.starting; } private async startInternal(): Promise { this.decoder = new JsonlDecoder(); this.sequence = 0; this.stderrBuffer = ""; this.recentMessages = []; const env = { ...process.env, ...this.options.env }; const libraryVariable = process.platform === "win32" ? "PATH" : process.platform === "darwin" ? "DYLD_LIBRARY_PATH" : "LD_LIBRARY_PATH"; const libraryDirectories = runtimeLibraryDirectories( this.options.runtimePath, ); const inheritedLibraryPath = env[libraryVariable] ?? (process.platform === "win32" ? env.Path : undefined); const libraryPath = inheritedLibraryPath ? `${libraryDirectories.join(delimiter)}${delimiter}${inheritedLibraryPath}` : libraryDirectories.join(delimiter); if (process.platform === "win32") { // Windows environment keys are case-insensitive, but Node can expose // the inherited key as `Path`. Keep one canonical value for spawn(). delete env.PATH; env.Path = libraryPath; } else { env[libraryVariable] = libraryPath; } this.child = spawn( this.options.runtimePath, [ "--stdio", "--model", this.options.modelPath, "--protocol-version", String(this.options.protocolVersion), "--threads", String(this.options.threads ?? 0), ], { stdio: ["pipe", "pipe", "pipe"], shell: false, windowsHide: true, env, }, ); const child = this.child; child.stdout.setEncoding("utf8"); child.stderr.setEncoding("utf8"); child.stdout.on("data", (chunk: string) => this.handleStdout(chunk)); child.stderr.on("data", (chunk: string) => this.handleStderr(chunk)); child.on("error", (error) => this.fail(error)); child.on("exit", (code, signal) => { this.child = undefined; const detail = signal ? `signal ${signal}` : `exit code ${code ?? "unknown"}`; this.fail(new Error(`Talk-to-Pi runtime exited with ${detail}.`)); }); const hello = await this.waitForMessage( (message) => message.type === "hello", this.options.startupTimeoutMs ?? 10_000, "Runtime did not send hello.", ); this.validateHello(hello); await this.waitForMessage( (message) => message.type === "ready", this.options.startupTimeoutMs ?? 30_000, "Runtime did not become ready.", ); } async sendCommand( command: RuntimeCommandInput, timeoutMs = 5_000, ): Promise { if (!this.child || !this.isAlive) throw new Error("Talk-to-Pi runtime is not running."); const id = command.id ?? randomUUID(); const message = { ...command, v: 1, id } as RuntimeCommandInput & { v: 1; id: string; }; const encoded = encodeProtocolMessage(message as ProtocolEnvelope); return new Promise((resolve, reject) => { const timer = setTimeout(() => { this.pending.delete(id); reject(new Error(`Runtime command timed out: ${message.type}`)); }, timeoutMs); this.pending.set(id, { resolve, reject, timer }); this.child?.stdin.write(encoded, (error) => { if (!error) return; clearTimeout(timer); this.pending.delete(id); reject(error); }); }); } async shutdown(): Promise { const child = this.child; if (!child) return; try { await this.sendCommand({ type: "shutdown" }, 2_000); } catch { // The escalation below is the fallback for a crashed or unresponsive runtime. } await this.waitForExit(child, 2_000); if (this.child === child && child.exitCode === null) child.kill("SIGTERM"); await this.waitForExit(child, 1_000); if (this.child === child && child.exitCode === null) child.kill("SIGKILL"); this.child = undefined; this.rejectPending(new Error("Runtime shut down.")); } private handleStdout(chunk: string): void { try { for (const message of this.decoder.feed(chunk)) this.handleMessage(message); } catch (error) { this.fail(error instanceof Error ? error : new Error(String(error))); } } private handleMessage(message: RuntimeMessage): void { this.recentMessages.push(message); if (this.recentMessages.length > 32) this.recentMessages.shift(); if (typeof message.seq === "number") { if (message.seq <= this.sequence) this.fail( new ProtocolError( "SEQUENCE_REGRESSION", "Runtime sequence regressed.", ), ); this.sequence = Math.max(this.sequence, message.seq); } if (typeof message.id === "string") { const pending = this.pending.get(message.id); if (pending) { this.pending.delete(message.id); clearTimeout(pending.timer); if (message.type === "error") { pending.reject( new Error( typeof message.message === "string" ? message.message : "Runtime rejected command.", ), ); } else { pending.resolve(message); } } } for (const handler of this.messageHandlers) handler(message); } private handleStderr(chunk: string): void { this.stderrBuffer = `${this.stderrBuffer}${chunk}`.slice(-16 * 1024); } private async waitForMessage( predicate: (message: RuntimeMessage) => boolean, timeoutMs: number, timeoutMessage: string, ): Promise { const existing = this.recentMessages.find(predicate); if (existing) return existing; return new Promise((resolve, reject) => { const timer = setTimeout(() => { cleanup(); reject(new Error(timeoutMessage)); }, timeoutMs); const off = this.onMessage((message) => { if (!predicate(message)) return; cleanup(); resolve(message); }); const onFailure = this.onFailure((error) => { cleanup(); reject(error); }); const cleanup = () => { clearTimeout(timer); off(); onFailure(); }; }); } private validateHello(message: RuntimeMessage): void { const engine = this.options.expectedEngine ?? "nemo-speech.cpp"; if (message.engine !== engine) throw new ProtocolError( "ENGINE_MISMATCH", `Expected runtime engine ${engine}, received ${String(message.engine)}.`, ); if ( !Array.isArray(message.protocolVersions) || !message.protocolVersions.includes(this.options.protocolVersion) ) throw new ProtocolError( "PROTOCOL_MISMATCH", "Runtime hello does not advertise the requested protocol version.", ); const abi = message.nemoAbi; const minimum = this.options.minimumEngineAbi ?? 1; const maximum = this.options.maximumEngineAbi ?? 1; if ( typeof abi !== "number" || !Number.isSafeInteger(abi) || abi < minimum || abi > maximum ) throw new ProtocolError( "ENGINE_ABI_MISMATCH", `Unsupported NeMo-Speech.cpp ABI: ${String(abi)} (supported ${minimum}-${maximum}).`, ); } private async waitForExit( child: ChildProcessWithoutNullStreams, timeoutMs: number, ): Promise { if (child.exitCode !== null) return; await new Promise((resolve) => { const timer = setTimeout(() => { child.off("exit", onExit); resolve(); }, timeoutMs); const onExit = () => { clearTimeout(timer); resolve(); }; child.once("exit", onExit); }); } private fail(error: Error): void { this.rejectPending(error); for (const handler of this.failureHandlers) handler(error); } private rejectPending(error: Error): void { for (const [id, pending] of this.pending) { clearTimeout(pending.timer); pending.reject(error); this.pending.delete(id); } } }