import { spawn as nodeSpawn, type ChildProcessWithoutNullStreams } from "node:child_process"; import { join } from "node:path"; import type { AppConfig } from "../config.js"; import { RpcClient } from "../rpc/rpc-client.js"; import { SessionBridge } from "../bridge/session-bridge.js"; import type { ThtRunner } from "../tht/tht-runner.js"; export interface SessionRuntime { rpc: RpcClient; bridge: SessionBridge; child: ChildProcessWithoutNullStreams; } /** Injected test double signature: produce a child process, no args needed. */ type SpawnFn = () => ChildProcessWithoutNullStreams; export class PiProcessManager { private runtimes = new Map(); private spawnFn: (sessionId: string, author: string) => ChildProcessWithoutNullStreams; constructor(private cfg: AppConfig, opts?: { spawnFn?: SpawnFn }) { if (opts?.spawnFn) { this.spawnFn = () => opts.spawnFn!(); } else { this.spawnFn = (sessionId: string, author: string) => { const harnessVenvBin = join(cfg.harnessDir, ".venv", "bin"); const env: NodeJS.ProcessEnv = { ...process.env, THT_SESSION: sessionId, THT_AUTHOR: author, PATH: `${harnessVenvBin}:${process.env.PATH ?? ""}`, }; const child = nodeSpawn(cfg.piBin, ["--mode", "rpc", "--approve"], { cwd: cfg.harnessDir, env, }); // Drain stderr so the child's stderr buffer never blocks the process. child.stderr.resume(); return child; }; } } count(): number { return this.runtimes.size; } get(id: string): SessionRuntime | undefined { return this.runtimes.get(id); } async spawnFor( sessionId: string, o: { provider?: string; model?: string; thinking?: string; author?: string; mode?: "new" | "resume" }, ): Promise { // Idempotent per session id: tear down any existing runtime for this id // first (before the cap check) so a resume/respawn neither leaks the old // child nor falsely hits the process cap. const existing = this.runtimes.get(sessionId); if (existing) { existing.child.kill(); this.runtimes.delete(sessionId); } if (this.runtimes.size >= this.cfg.maxPiProcesses) { throw new Error("max Pi processes reached"); } const author = o.author ?? "dev@local"; const child = this.spawnFn(sessionId, author); const rpc = new RpcClient(child); const bridge = new SessionBridge(rpc); const rt: SessionRuntime = { rpc, bridge, child }; this.runtimes.set(sessionId, rt); // Identity-checked: a stale child's exit must not evict a newer runtime. child.on("exit", () => { if (this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId); }); const provider = o.provider ?? this.cfg.defaults.provider; const model = o.model ?? this.cfg.defaults.model; const thinking = o.thinking ?? this.cfg.defaults.thinking; if (provider && model) { await rpc.request({ type: "set_model", provider, modelId: model } as object & { type: string }); } if (thinking) { await rpc.request({ type: "set_thinking_level", level: thinking } as object & { type: string }); } const message = o.mode === "resume" ? `/riprendi-sessione ${sessionId}` : `/nuova-domanda "kickoff"`; rpc.send({ type: "prompt", message }); return rt; } async resume(sessionId: string, tht: ThtRunner): Promise { const manifest = await tht.sessionShow(sessionId) as { provider?: string; model?: string; thinking?: string } | null; return this.spawnFor(sessionId, { provider: manifest?.provider, model: manifest?.model, thinking: manifest?.thinking, mode: "resume", }); } teardown(id: string): void { const rt = this.runtimes.get(id); if (rt) { rt.child.kill(); this.runtimes.delete(id); } } }