From de7200f7dd39b787b0628b7a6558debf26dc5bbf Mon Sep 17 00:00:00 2001 From: mptyl Date: Sat, 27 Jun 2026 21:05:03 +0200 Subject: [PATCH] feat(backend): PiProcessManager (one Pi per session, cap, set_model/thinking) Co-Authored-By: Claude Sonnet 4.6 --- backend/src/pi/pi-process-manager.ts | 83 +++++++++++++++++++++++++ backend/test/pi-process-manager.test.ts | 28 +++++++++ 2 files changed, 111 insertions(+) create mode 100644 backend/src/pi/pi-process-manager.ts create mode 100644 backend/test/pi-process-manager.test.ts diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts new file mode 100644 index 00000000..a3fc079f --- /dev/null +++ b/backend/src/pi/pi-process-manager.ts @@ -0,0 +1,83 @@ +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"; + +export interface SessionRuntime { + rpc: RpcClient; + bridge: SessionBridge; + child: ChildProcessWithoutNullStreams; +} + +type SpawnFn = () => ChildProcessWithoutNullStreams; + +export class PiProcessManager { + private runtimes = new Map(); + private spawnFn: (sessionId: string, author: string) => ChildProcessWithoutNullStreams; + + constructor(private cfg: AppConfig, opts?: { spawnFn?: () => ChildProcessWithoutNullStreams }) { + 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 }, + ): Promise { + 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); + child.on("exit", () => 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 }); + } + + rpc.send({ type: "prompt", message: `/nuova-domanda "kickoff"` }); + return rt; + } + + teardown(id: string): void { + const rt = this.runtimes.get(id); + if (rt) { + rt.child.kill(); + this.runtimes.delete(id); + } + } +} diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts new file mode 100644 index 00000000..3aa3907b --- /dev/null +++ b/backend/test/pi-process-manager.test.ts @@ -0,0 +1,28 @@ +import { test, expect } from "vitest"; +import { spawn } from "node:child_process"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { PiProcessManager } from "../src/pi/pi-process-manager.js"; +import { loadConfig } from "../src/config.js"; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +const FAKE = path.resolve(__dirname, "../../harness/tests/fake_pi/fake_pi_rpc.mjs"); +const SCRIPT = path.resolve(__dirname, "../../harness/tests/fake_pi/scripts/f1_disambiguation.json"); + +test("spawnFor avvia un runtime e il bridge emette il widget F1", async () => { + const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); + const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); + const rt = await mgr.spawnFor("2026-06-27-100000-x", {}); + const widget = await new Promise((res) => rt.bridge.onClientEvent((e) => e.type === "ui_request" && res(e))); + expect(widget.ui_request.widget).toBe("select"); + mgr.teardown("2026-06-27-100000-x"); + expect(mgr.count()).toBe(0); +}); + +test("oltre maxPiProcesses solleva errore", async () => { + const cfg = { ...loadConfig({}), maxPiProcesses: 1 }; + const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); + await mgr.spawnFor("a", {}); + await expect(mgr.spawnFor("b", {})).rejects.toThrow(/max/i); + mgr.teardown("a"); +});