diff --git a/backend/src/bridge/session-bridge.ts b/backend/src/bridge/session-bridge.ts index 830cade4..e7834b73 100644 --- a/backend/src/bridge/session-bridge.ts +++ b/backend/src/bridge/session-bridge.ts @@ -7,7 +7,10 @@ export type ClientEvent = | { type: "info"; [k: string]: any } | { type: "system_event"; [k: string]: any }; +export type TurnState = "idle" | "running" | "waiting" | "failed"; + export class SessionBridge { + private state: TurnState = "idle"; private pending: any = null; // Pi (rpc-mode createDialogPromise) assegna a ogni ctx.ui.input un id RPC PROPRIO // (crypto.randomUUID) e correla extension_ui_response su quell'id — NON sull'id interno @@ -23,7 +26,19 @@ export class SessionBridge { try { descriptor = JSON.parse(m.title as string); } catch { return; } this.pending = descriptor; this.pendingPiId = m.id as string | null; + this.state = "waiting"; this.fan({ type: "ui_request", ui_request: descriptor }); + } else if ( + m.type === "message_end" && + m.message?.role === "assistant" && + m.message?.stopReason === "error" + ) { + this.state = "failed"; + this.fan({ + type: "info", + level: "error", + text: "Model request failed. Check provider connectivity, then Resume the session.", + }); } else if (m.type === "extension_ui_request" && m.method === "notify") { this.fan({ type: "info", level: m.notifyType ?? "info", text: m.message ?? "" }); } else if (m.type === "message_update" && m.assistantMessageEvent?.type === "text_delta") { @@ -37,8 +52,10 @@ export class SessionBridge { } else if (m.type === "system_event") { this.fan(m as ClientEvent); } else if (m.type === "agent_end") { + if (this.state !== "failed" && !this.pending) this.state = "idle"; this.fan({ type: "system_event", event: "agent_end" }); } else if (m.type === "agent_start") { + this.state = "running"; this.fan({ type: "system_event", event: "agent_start" }); } else if (m.type === "turn_end") { this.fan({ type: "system_event", event: "turn_end" }); @@ -48,6 +65,10 @@ export class SessionBridge { private fan(e: ClientEvent) { for (const cb of this.cbs) cb(e); } + turnState(): TurnState { return this.state; } + + beginTurn(): void { this.state = "running"; } + onClientEvent(cb: (e: ClientEvent) => void): void { this.cbs.add(cb); } /** Eventi generati dal backend stesso (es. exit inatteso del child Pi), non da Pi. */ @@ -57,11 +78,15 @@ export class SessionBridge { // Correla sull'id RPC di Pi; `value` porta l'uiResponse (con l'id del descriptor) cosi' // il check interno del gate (resp.id === descriptor.id) regge. const piId = this.pendingPiId ?? uiResponse.id; + this.state = "running"; this.rpc.send({ type: "extension_ui_response", id: piId, value: JSON.stringify(uiResponse) }); if (this.pending && uiResponse.id === this.pending.id) { this.pending = null; this.pendingPiId = null; } } - steer(text: string): void { this.rpc.send({ type: "steer", message: text }); } + steer(text: string): void { + this.state = "running"; + this.rpc.send({ type: "steer", message: text }); + } pendingWidget(): object | null { return this.pending; } } diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index bfecedb1..7f6af95b 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -94,6 +94,7 @@ export class PiProcessManager { 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 @@ -136,6 +137,7 @@ export class PiProcessManager { const message = o.mode === "resume" ? `/riprendi-sessione ${sessionId}` : `/nuova-domanda ${JSON.stringify(o.question ?? "")}`; + rt.bridge.beginTurn(); rt.rpc.send({ type: "prompt", message }); } diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index 5ffb179a..f10d6ce8 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -155,6 +155,34 @@ test("createFor does not prompt until start is called", async () => { mgr.teardown("sid-deferred"); }); +test("a created runtime is active during configure/bootstrap", () => { + const child = recordingChild(); + const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); + const runtime = mgr.createFor("starting-id", {}); + expect(runtime.bridge.turnState()).toBe("running"); + mgr.teardown("starting-id"); +}); + +test("start reactivates the runtime before sending the prompt", () => { + const child = recordingChild(); + const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); + const runtime = mgr.createFor("restart-id", {}); + child.stdout.emit("data", `${JSON.stringify({ type: "agent_end", messages: [] })}\n`); + expect(runtime.bridge.turnState()).toBe("idle"); + let stateAtWrite = runtime.bridge.turnState(); + child.stdin.write = (data: unknown) => { + stateAtWrite = runtime.bridge.turnState(); + child._writes.push(String(data)); + return true; + }; + + mgr.start("restart-id", runtime, { question: "q" }); + + expect(stateAtWrite).toBe("running"); + expect(runtime.bridge.turnState()).toBe("running"); + mgr.teardown("restart-id"); +}); + test("production spawn uses explicit Pi path and passes portable data root without rewriting PATH", async () => { vi.stubEnv("PATH", "/usr/local/bin:/usr/bin"); vi.stubEnv("PI_PROVIDER_API_KEY", "provider-secret"); diff --git a/backend/test/session-bridge.test.ts b/backend/test/session-bridge.test.ts index 73d15d94..7069a1e6 100644 --- a/backend/test/session-bridge.test.ts +++ b/backend/test/session-bridge.test.ts @@ -42,6 +42,78 @@ test("real Pi thinking_delta becomes a dedicated activity_delta to the FE", () = expect(seen).toEqual([{ type: "activity_delta", text: "Valuto le ambiguità" }]); }); +test("assistant provider errors are sanitized and leave the turn failed", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + bridge.beginTurn(); + + fire({ + type: "message_end", + message: { + role: "assistant", + stopReason: "error", + errorMessage: "Connection failed for https://secret.invalid/?api_key=DO_NOT_LEAK", + }, + }); + fire({ type: "agent_end", messages: [] }); + + expect(bridge.turnState()).toBe("failed"); + expect(seen).toContainEqual({ + type: "info", + level: "error", + text: "Model request failed. Check provider connectivity, then Resume the session.", + }); + expect(JSON.stringify(seen)).not.toContain("DO_NOT_LEAK"); +}); + +test("reviewer wait and response transition waiting back to running", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + bridge.beginTurn(); + fire({ + type: "extension_ui_request", + id: "pi-1", + method: "input", + title: JSON.stringify({ id: "gate-1", widget: "select" }), + }); + expect(bridge.turnState()).toBe("waiting"); + bridge.respond({ id: "gate-1", choices: ["approve"] }); + expect(bridge.turnState()).toBe("running"); + fire({ type: "agent_end", messages: [] }); + expect(bridge.turnState()).toBe("idle"); +}); + +test("agent_start transitions an idle turn to running", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + + expect(bridge.turnState()).toBe("idle"); + fire({ type: "agent_start" }); + + expect(bridge.turnState()).toBe("running"); + expect(seen).toContainEqual({ type: "system_event", event: "agent_start" }); +}); + +test("agent_end leaves a pending reviewer wait intact", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + bridge.beginTurn(); + fire({ + type: "extension_ui_request", + id: "pi-1", + method: "input", + title: JSON.stringify({ id: "gate-1", widget: "select" }), + }); + + fire({ type: "agent_end", messages: [] }); + + expect(bridge.turnState()).toBe("waiting"); +}); + test("extension_ui_request nativo (method:input, title=json) diventa ui_request col descriptor ed è il pendente", () => { const { rpc, fire } = fakeRpc(); const b = new SessionBridge(rpc); @@ -79,8 +151,20 @@ test("agent_end di Pi diventa un system_event agent_end per il FE", () => { expect(seen).toEqual([{ type: "system_event", event: "agent_end" }]); }); -test("steer invia un comando steer", () => { - const { rpc, sent } = fakeRpc(); - new SessionBridge(rpc).steer("considera solo il 2024"); +test("steer invia un comando steer e riattiva il turno", () => { + const { rpc, sent, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + bridge.beginTurn(); + fire({ + type: "extension_ui_request", + id: "pi-1", + method: "input", + title: JSON.stringify({ id: "gate-1", widget: "select" }), + }); + expect(bridge.turnState()).toBe("waiting"); + + bridge.steer("considera solo il 2024"); + expect(sent.at(-1)).toEqual({ type: "steer", message: "considera solo il 2024" }); + expect(bridge.turnState()).toBe("running"); });