import { access, stat } from "node:fs/promises"; import { constants } from "node:fs"; import { randomUUID } from "node:crypto"; import { getTalkPaths, type TalkPaths } from "../config/paths.js"; import { Provisioner, type ProvisionerProgress, } from "../provisioning/provisioner.js"; import { type RuntimeMessage } from "./protocol.js"; import { RuntimeProcess } from "./runtime-process.js"; export interface ProvisionProgress { asset: "runtime" | "model"; receivedBytes: number; totalBytes?: number; } export interface RuntimeDiagnostics { runtimePath: string; modelPath: string; runtimeInstalled: boolean; modelInstalled: boolean; processAlive: boolean; stderr: string; } export interface RecordingOptions { sessionId: string; language: string; onEvent: (message: RuntimeMessage) => void; } export interface TalkRuntime { ensureProvisioned( onProgress?: (progress: ProvisionProgress) => void, signal?: AbortSignal, ): Promise; ensureReady(signal?: AbortSignal): Promise; startRecording(options: RecordingOptions): Promise; stopRecording(sessionId: string): Promise; cancelRecording(sessionId: string): Promise; getDiagnostics(): Promise; shutdown(): Promise; } export class RuntimeUnavailableError extends Error { readonly code = "ASSET_NOT_PROVISIONED" as const; constructor(message: string) { super(message); this.name = "RuntimeUnavailableError"; } } export class RuntimeManager implements TalkRuntime { private process: RuntimeProcess | undefined; private startPromise: Promise | undefined; private restartPromise: Promise | undefined; private restartError = ""; private modelProvisioned = false; private readonly listeners = new Map void>>(); private readonly provisioner: Provisioner; constructor(private readonly paths: TalkPaths = getTalkPaths()) { this.provisioner = new Provisioner(paths); } async ensureProvisioned( onProgress?: (progress: ProvisionProgress) => void, signal?: AbortSignal, ): Promise { if (signal?.aborted) throw new DOMException("The operation was aborted.", "AbortError"); if (await this.hasExplicitLocalAssets()) return; if (this.modelProvisioned) return; if (this.hasRuntimeOnlyOverride()) { await this.provisioner.ensureModel(onProgress, signal); this.modelProvisioned = true; return; } try { await this.provisioner.ensure( onProgress ? (progress: ProvisionerProgress) => onProgress(progress) : undefined, signal, ); } catch (error) { if (error instanceof RuntimeUnavailableError) throw error; throw new RuntimeUnavailableError( error instanceof Error ? error.message : String(error), ); } const [runtimeInstalled, modelInstalled] = await Promise.all([ this.isExecutable(this.paths.runtimePath), this.isRegularFile(this.paths.modelPath), ]); if (!runtimeInstalled || !modelInstalled) { throw new RuntimeUnavailableError( "Talk-to-Pi provisioning completed without installing all required assets.", ); } this.modelProvisioned = true; } async ensureReady(signal?: AbortSignal): Promise { await this.ensureProvisioned(undefined, signal); await this.restartPromise; if (signal?.aborted) throw new DOMException("The operation was aborted.", "AbortError"); if (this.process?.isAlive) return; if (!this.startPromise) { this.startPromise = (async () => { const process = this.createProcess(); await process.start(); this.process = process; this.restartError = ""; })().finally(() => { this.startPromise = undefined; }); } await this.startPromise; } async startRecording(options: RecordingOptions): Promise { await this.ensureReady(); const process = this.requireProcess(); this.cleanupSession(options.sessionId); const cleanup = () => this.cleanupSession(options.sessionId); const offMessage = process.onMessage((message) => { if (message.sessionId !== options.sessionId && message.type !== "error") return; options.onEvent(message); if (message.type === "error") { cleanup(); this.restartInBackground(process); return; } if (["recording_finalized", "recording_cancelled"].includes(message.type)) cleanup(); }); const offFailure = process.onFailure((error) => { options.onEvent({ v: 1, type: "error", sessionId: options.sessionId, code: "RUNTIME_CRASHED", message: error.message, recoverable: true, }); cleanup(); this.restartInBackground(process); }); this.listeners.set(options.sessionId, [offMessage, offFailure]); try { await process.sendCommand({ type: "start", sessionId: options.sessionId, language: options.language, }); } catch (error) { cleanup(); throw error; } } async stopRecording(sessionId: string): Promise { await this.requireProcess().sendCommand({ type: "stop", sessionId }); } async cancelRecording(sessionId: string): Promise { try { await this.requireProcess().sendCommand({ type: "cancel", sessionId }); } finally { this.cleanupSession(sessionId); } } async getDiagnostics(): Promise { const runtimeInstalled = await this.isExecutable(this.paths.runtimePath); const modelInstalled = await this.isRegularFile(this.paths.modelPath); return { runtimePath: this.paths.runtimePath, modelPath: this.paths.modelPath, runtimeInstalled, modelInstalled, processAlive: this.process?.isAlive ?? false, stderr: this.process?.stderr || this.restartError, }; } async shutdown(): Promise { for (const sessionId of this.listeners.keys()) this.cleanupSession(sessionId); await this.restartPromise; await this.startPromise; const process = this.process; this.process = undefined; await process?.shutdown(); } private createProcess(): RuntimeProcess { const process = new RuntimeProcess({ runtimePath: this.paths.runtimePath, modelPath: this.paths.modelPath, protocolVersion: 1, }); process.onFailure(() => this.restartInBackground(process)); return process; } private restartInBackground(failedProcess: RuntimeProcess): void { if (this.process !== failedProcess || this.restartPromise) return; this.process = undefined; this.restartError = ""; this.restartPromise = (async () => { await failedProcess.shutdown(); const replacement = this.createProcess(); await replacement.start(); this.process = replacement; })() .catch((error: unknown) => { this.restartError = error instanceof Error ? error.message : String(error); }) .finally(() => { this.restartPromise = undefined; }); } private cleanupSession(sessionId: string): void { for (const off of this.listeners.get(sessionId) ?? []) off(); this.listeners.delete(sessionId); } private requireProcess(): RuntimeProcess { if (!this.process) throw new Error("Talk-to-Pi runtime is not ready."); return this.process; } private async hasExplicitLocalAssets(): Promise { const runtimeOverride = process.env.TALK_TO_PI_RUNTIME_PATH; const modelOverride = process.env.TALK_TO_PI_MODEL_PATH; if (!runtimeOverride && !modelOverride) return false; if (runtimeOverride && !modelOverride) { if (!(await this.isExecutable(runtimeOverride))) { throw new RuntimeUnavailableError( "The configured TALK_TO_PI_RUNTIME_PATH is not executable.", ); } return false; } if (!runtimeOverride || !modelOverride) { throw new RuntimeUnavailableError( "TALK_TO_PI_RUNTIME_PATH and TALK_TO_PI_MODEL_PATH must be set together unless the model is provisioned from Hugging Face.", ); } const [runtimeInstalled, modelInstalled] = await Promise.all([ this.isExecutable(runtimeOverride), this.isRegularFile(modelOverride), ]); if (!runtimeInstalled || !modelInstalled) { throw new RuntimeUnavailableError( "The configured local Talk-to-Pi runtime or model does not exist.", ); } return true; } private hasRuntimeOnlyOverride(): boolean { return Boolean( process.env.TALK_TO_PI_RUNTIME_PATH && !process.env.TALK_TO_PI_MODEL_PATH, ); } private async isExecutable(path: string): Promise { try { await access(path, constants.X_OK); return (await stat(path)).isFile(); } catch { return false; } } private async isRegularFile(path: string): Promise { try { return (await stat(path)).isFile(); } catch { return false; } } } export function createSessionId(): string { return randomUUID(); }