import { spawn as nodeSpawn, type ChildProcessWithoutNullStreams } from "node:child_process"; 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"; import { buildPiChildEnv, canonicalPiProvider } from "./provider-credentials.js"; import { secretValue } from "../config/secret-bundle.js"; export interface SessionRuntime { rpc: RpcClient; bridge: SessionBridge; child: ChildProcessWithoutNullStreams; } /** Injectable child-process boundary; callbacks may ignore arguments in simpler tests. */ type SpawnFn = ( command: string, args: string[], options: { cwd: string; env: NodeJS.ProcessEnv }, ) => ChildProcessWithoutNullStreams; export class PiProcessManager { private runtimes = new Map(); private spawnFn: ( sessionId: string, author: string, provider: string | undefined, ) => ChildProcessWithoutNullStreams; constructor(private cfg: AppConfig, opts?: { spawnFn?: SpawnFn }) { if (opts?.spawnFn) { this.spawnFn = (sessionId, author, provider) => this.spawnPi(opts.spawnFn!, sessionId, author, provider); } else { this.spawnFn = (sessionId, author, provider) => this.spawnPi(nodeSpawn, sessionId, author, provider); } } private spawnPi( spawnFn: SpawnFn, sessionId: string, author: string, provider: string | undefined, ): ChildProcessWithoutNullStreams { const env = buildPiChildEnv({ provider, credentialValue: secretValue(this.cfg, "THT_MODEL_API_KEY"), credentialFile: this.cfg.modelApiKeyFile, additions: { THT_SESSION: sessionId, THT_AUTHOR: author }, }); delete env.THT_DATA_ROOT; if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot; // pi 0.73 removed `--approve`: rpc mode is headless and its argv is intentionally minimal. const child = spawnFn(this.cfg.piBin, ["--mode", "rpc"], { cwd: this.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 provider = canonicalPiProvider(o.provider ?? this.cfg.defaults.provider); const child = this.spawnFn(sessionId, author, provider); 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. // Expected teardowns (teardown()/respawn) delete the runtime from the map BEFORE the // exit event fires, so reaching this branch with `rt` still mapped means the child // died on its own: tell the client, or the UI spins forever waiting for a turn end. child.on("exit", (code) => { if (this.runtimes.get(sessionId) === rt) { this.runtimes.delete(sessionId); rt.bridge.emitClientEvent({ type: "info", level: "error", text: `Pi process exited unexpectedly (code ${code ?? "?"})`, }); rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" }); } }); 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); } } }