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 <noreply@anthropic.com>
This commit is contained in:
@@ -49,6 +49,14 @@ export class PiProcessManager {
|
|||||||
sessionId: string,
|
sessionId: string,
|
||||||
o: { provider?: string; model?: string; thinking?: string; author?: string },
|
o: { provider?: string; model?: string; thinking?: string; author?: string },
|
||||||
): Promise<SessionRuntime> {
|
): 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) {
|
if (this.runtimes.size >= this.cfg.maxPiProcesses) {
|
||||||
throw new Error("max Pi processes reached");
|
throw new Error("max Pi processes reached");
|
||||||
}
|
}
|
||||||
@@ -58,7 +66,8 @@ export class PiProcessManager {
|
|||||||
const bridge = new SessionBridge(rpc);
|
const bridge = new SessionBridge(rpc);
|
||||||
const rt: SessionRuntime = { rpc, bridge, child };
|
const rt: SessionRuntime = { rpc, bridge, child };
|
||||||
this.runtimes.set(sessionId, rt);
|
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 provider = o.provider ?? this.cfg.defaults.provider;
|
||||||
const model = o.model ?? this.cfg.defaults.model;
|
const model = o.model ?? this.cfg.defaults.model;
|
||||||
|
|||||||
@@ -4,7 +4,9 @@ export class SseHub {
|
|||||||
subscribe(sessionId: string, send: Send, pending?: object | null): () => void {
|
subscribe(sessionId: string, send: Send, pending?: object | null): () => void {
|
||||||
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
if (!this.subs.has(sessionId)) this.subs.set(sessionId, new Set());
|
||||||
this.subs.get(sessionId)!.add(send);
|
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);
|
return () => this.subs.get(sessionId)?.delete(send);
|
||||||
}
|
}
|
||||||
publish(sessionId: string, event: string, data: object): void {
|
publish(sessionId: string, event: string, data: object): void {
|
||||||
|
|||||||
@@ -34,6 +34,33 @@ test("l'exit del child rimuove il runtime dalla mappa (exit handler)", async ()
|
|||||||
expect(mgr.get("exit-test")).toBeUndefined();
|
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<void>((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 () => {
|
test("oltre maxPiProcesses solleva errore", async () => {
|
||||||
const cfg = { ...loadConfig({}), maxPiProcesses: 1 };
|
const cfg = { ...loadConfig({}), maxPiProcesses: 1 };
|
||||||
const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any });
|
const mgr = new PiProcessManager(cfg, { spawnFn: () => spawn("node", [FAKE, SCRIPT]) as any });
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import { SseHub } from "../src/sse/sse-hub.js";
|
|||||||
test("re-emette il widget pendente alla sottoscrizione", () => {
|
test("re-emette il widget pendente alla sottoscrizione", () => {
|
||||||
const hub = new SseHub(); const sent: any[] = [];
|
const hub = new SseHub(); const sent: any[] = [];
|
||||||
hub.subscribe("s1", (ev, data) => sent.push({ ev, data }), { id: "u1", widget: "select" });
|
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", () => {
|
test("publish raggiunge i subscriber e unsubscribe li stacca", () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user