From c63b2bd1268aefaf9565f8328e346d62db5d066d Mon Sep 17 00:00:00 2001 From: mptyl Date: Sat, 27 Jun 2026 21:53:12 +0200 Subject: [PATCH] fix(backend): idempotent spawnFor + identity-checked exit + unified SSE re-emit shape - spawnFor now tears down any existing runtime for the same session id before the cap check, so resume/respawn neither leaks the old child nor falsely hits maxPiProcesses - exit handler is identity-checked (captures rt) so a stale child's late exit cannot evict a newer runtime - SSE pending re-emit now sends the full ClientEvent shape { type: "ui_request", ui_request } to match hub.publish live events - tests: same-id respawn replaces runtime (count 1); old child exit does not evict new runtime; sse-hub re-emit asserts unified shape Co-Authored-By: Claude Sonnet 4.6 --- backend/src/pi/pi-process-manager.ts | 11 +++++++++- backend/src/sse/sse-hub.ts | 4 +++- backend/test/pi-process-manager.test.ts | 27 +++++++++++++++++++++++++ backend/test/sse-hub.test.ts | 2 +- 4 files changed, 41 insertions(+), 3 deletions(-) diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index e06c98ad..cb0daf15 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -49,6 +49,14 @@ export class PiProcessManager { sessionId: string, o: { provider?: string; model?: string; thinking?: string; author?: string }, ): 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"); } @@ -58,7 +66,8 @@ export class PiProcessManager { const bridge = new SessionBridge(rpc); const rt: SessionRuntime = { rpc, bridge, child }; this.runtimes.set(sessionId, rt); - child.on("exit", () => this.runtimes.delete(sessionId)); + // Identity-checked: a stale child's exit must not evict a newer runtime. + child.on("exit", () => { if (this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId); }); const provider = o.provider ?? this.cfg.defaults.provider; const model = o.model ?? this.cfg.defaults.model; diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts index 0ccd9813..23d750d2 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -4,7 +4,9 @@ export class SseHub { subscribe(sessionId: string, send: Send, pending?: object | null): () => void { if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set()); this.subs.get(sessionId)!.add(send); - if (pending) send("ui_request", { ui_request: pending }); + // Match the live-event shape (hub.publish sends the full ClientEvent): + // both carry { type, ui_request } so re-emit and live widgets are identical. + if (pending) send("ui_request", { type: "ui_request", ui_request: pending }); return () => this.subs.get(sessionId)?.delete(send); } publish(sessionId: string, event: string, data: object): void { diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index 1b0edee7..fb53bbfb 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -34,6 +34,33 @@ test("l'exit del child rimuove il runtime dalla mappa (exit handler)", async () expect(mgr.get("exit-test")).toBeUndefined(); }); +test("spawnFor sullo STESSO id uccide il vecchio child e sostituisce il runtime (count resta 1)", async () => { + const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); + const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); + const first = await mgr.spawnFor("dup-id", {}); + expect(mgr.count()).toBe(1); + const firstExited = new Promise((res) => first.child.on("exit", () => res())); + const second = await mgr.spawnFor("dup-id", {}); + await firstExited; // the old child was killed by the idempotent respawn + expect(mgr.count()).toBe(1); + expect(mgr.get("dup-id")).toBe(second); + expect(second).not.toBe(first); + mgr.teardown("dup-id"); +}); + +test("l'exit del VECCHIO child non elimina il nuovo runtime (exit identity-checked)", async () => { + const cfg = loadConfig({ THT_HARNESS_DIR: "../harness" }); + const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); + const first = await mgr.spawnFor("respawn-id", {}); + const second = await mgr.spawnFor("respawn-id", {}); + // The old child's exit handler fires after the respawn; it must NOT evict `second`. + await new Promise((res) => setImmediate(res)); + expect(mgr.get("respawn-id")).toBe(second); + expect(mgr.count()).toBe(1); + void first; + mgr.teardown("respawn-id"); +}); + test("oltre maxPiProcesses solleva errore", async () => { const cfg = { ...loadConfig({}), maxPiProcesses: 1 }; const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any }); diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index 5c1501d8..7f4f6fe0 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -4,7 +4,7 @@ import { SseHub } from "../src/sse/sse-hub.js"; test("re-emette il widget pendente alla sottoscrizione", () => { const hub = new SseHub(); const sent: any[] = []; hub.subscribe("s1", (ev, data) => sent.push({ ev, data }), { id: "u1", widget: "select" }); - expect(sent[0]).toEqual({ ev: "ui_request", data: { ui_request: { id: "u1", widget: "select" } } }); + expect(sent[0]).toEqual({ ev: "ui_request", data: { type: "ui_request", ui_request: { id: "u1", widget: "select" } } }); }); test("publish raggiunge i subscriber e unsubscribe li stacca", () => {