feat(backend): PiProcessManager (one Pi per session, cap, set_model/thinking)
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -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<string, SessionRuntime>();
|
||||
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<SessionRuntime> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user