Files
ThothII/backend/src/pi/pi-process-manager.ts
T
marcopanandClaude Fable 5 2410f01b34 fix(bridge): forward Pi agent_end so the spinner stops at workflow completion
The FE derived 'working' purely as activeSession && !pendingWidget, so the
final workflow turn — the only one that ends without a follow-up gate —
left the spinner on forever (observed live: 21592s after F8 approve).

- SessionBridge maps Pi's agent_end -> SSE system_event {event: agent_end}
- PiProcessManager notifies the client (info error + synthetic agent_end)
  when the child dies unexpectedly; expected teardowns stay silent
- sessionStore tracks agentActive (on: user entry/text_delta/ui_request,
  off: agent_end); AppShell working now requires it; resume sets it
  optimistically

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-07 10:14:43 +02:00

124 lines
4.6 KiB
TypeScript

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<string, SessionRuntime>();
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 ?? ""}`,
};
// pi 0.73 (the @mariozechner rebrand) removed the `--approve` flag: rpc mode is
// headless and runs tools without an approval gate, so passing it makes pi exit
// with "Unknown option: --approve". Args are intentionally just `--mode rpc`.
const child = nodeSpawn(cfg.piBin, ["--mode", "rpc"], {
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<SessionRuntime> {
// 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.
// 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 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<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);
}
}
}