Files
ThothII/backend/src/pi/pi-process-manager.ts
T

170 lines
6.7 KiB
TypeScript

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;
}
export interface RuntimeOptions {
provider?: string;
model?: string;
thinking?: string;
author?: string;
question?: string;
mode?: "new" | "resume";
}
/** 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<string, SessionRuntime>();
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 },
});
// The Thoth gate executes the deterministic `tht` CLI as a Pi tool. Give only
// this managed session process the adapter values already loaded by the core
// entrypoint; the generic provider helper continues to scrub them by default.
for (const name of [
"THT_DWH_API_KEY", "THT_VEC_API_KEY", "THT_VEC_WRITE_API_KEY", "THT_SSL_CA",
] as const) {
if (process.env[name] !== undefined) env[name] = process.env[name];
}
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,
});
// Log stderr for debugging (was silently drained)
child.stderr.on("data", (d: Buffer) => console.error(`[pi:${sessionId}] stderr:`, d.toString().trim()));
return child;
}
count(): number { return this.runtimes.size; }
get(id: string): SessionRuntime | undefined { return this.runtimes.get(id); }
/** Spawn and register a runtime synchronously, without starting a model turn. */
createFor(sessionId: string, o: RuntimeOptions = {}): SessionRuntime {
// A duplicate start must never tear down a live session: that used to send
// SIGTERM to the in-flight Pi process and lose its pending gate.
const existing = this.runtimes.get(sessionId);
if (existing) {
throw new Error(`session runtime already active: ${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 };
bridge.beginTurn();
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) => {
console.error(`[pi:${sessionId}] exited code=${code ?? "?"} mapped=${this.runtimes.get(sessionId) === rt}`);
if (this.runtimes.get(sessionId) === rt) {
rt.bridge.markFailed();
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: "session_failed" });
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
}
});
return rt;
}
/** Configure model and thinking. Safe to run alongside deterministic retrieval. */
async configure(rt: SessionRuntime, o: RuntimeOptions = {}): Promise<void> {
const provider = canonicalPiProvider(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 rt.rpc.request({ type: "set_model", provider, modelId: model } as object & { type: string });
}
if (thinking) {
await rt.rpc.request({ type: "set_thinking_level", level: thinking } as object & { type: string });
}
}
/** Start the first turn only after callers have attached the runtime bridge. */
start(sessionId: string, rt: SessionRuntime, o: RuntimeOptions = {}): void {
if (this.runtimes.get(sessionId) !== rt) throw new Error("session runtime is no longer active");
const message = o.mode === "resume"
? `/riprendi-sessione ${sessionId}`
: `/nuova-domanda ${JSON.stringify(o.question ?? "")}`;
rt.bridge.beginTurn();
rt.rpc.send({ type: "prompt", message });
}
async spawnFor(sessionId: string, o: RuntimeOptions = {}): Promise<SessionRuntime> {
const rt = this.createFor(sessionId, o);
await this.configure(rt, o);
this.start(sessionId, rt, o);
return rt;
}
async resume(sessionId: string, tht: ThtRunner): Promise<SessionRuntime> {
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);
}
}
}