diff --git a/.superpowers/sdd/predeploy-fix-report.md b/.superpowers/sdd/predeploy-fix-report.md index ee9c497d..93f2ccb1 100644 --- a/.superpowers/sdd/predeploy-fix-report.md +++ b/.superpowers/sdd/predeploy-fix-report.md @@ -319,6 +319,152 @@ cd frontend && npm run build exit 0 ``` +## Integrated re-review closure (2026-07-15) + +This section supersedes the earlier cold same-session assertion that the replacement URL carries +`lastEventId=8`. That behavior was correct only while the backend process and its in-memory id +sequence survived. A restarted backend begins a fresh sequence, so a successful cold Resume now +explicitly discards the browser's cursor before replacing the EventSource. + +All four integrated re-review findings are closed: + +1. `AppShell` passes a dedicated cursor-reset epoch to `useSessionStream`. A cold same-session + Resume increments it only after `alreadyActive: false`; a high cursor such as `901` is omitted + from the replacement URL and fresh low-id events/gates are consumed. An already-active + same-session Resume still preserves its source, cursor, and store. +2. `useSessionStream` no longer mutates the cursor ref during render. Effect setup resets cursor + state on session/reset-epoch changes, callbacks are guarded by a captured active-source + identity, and cleanup clears only its own active identity. A queued event from the replaced + source cannot write the new store or poison its next reconnect URL. +3. Backend Resume is serialized per session and rechecks runtime state inside the lock. Manifest, + readiness, and reopen validation precede the transport commit. Idle/failed replacement creates + and binds the new runtime before `hub.clear`, which occurs synchronously immediately before the + first `Resuming session` publish. Reopen/create failure returns exactly + `Session could not be resumed. Check configuration and connectivity, then try again.`, keeps the + prior hub buffer/subscribers attached, and does not expose exception sentinels. Concurrent calls + perform one cold start and the waiter returns `alreadyActive: true`. +4. `SseHub.forget(id)` removes subscribers, buffered events, and the last id. Permanent session + DELETE invokes it after disk deletion; ordinary close and Resume continue to use `clear`, which + preserves the id sequence. + +### Re-review files + +Production: + +- `backend/src/pi/pi-process-manager.ts` +- `backend/src/routes/sessions.ts` +- `backend/src/sse/sse-hub.ts` +- `frontend/src/shell/AppShell.tsx` +- `frontend/src/stream/useSessionStream.ts` + +Tests/support: + +- `backend/test/pi-process-manager.test.ts` +- `backend/test/routes-sessions.test.ts` +- `backend/test/sse-hub.test.ts` +- `frontend/src/shell/AppShell.session-mgmt.test.tsx` +- `frontend/src/stream/useSessionStream.test.tsx` +- `frontend/src/test/fakeEventSource.ts` + +### Re-review TDD RED/GREEN evidence + +Frontend RED command: + +```text +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx \ + src/shell/AppShell.session-mgmt.test.tsx +``` + +RED output (exit 1): + +```text +Test Files 2 failed (2) +Tests 3 failed | 21 passed (24) + +reset epoch: expected the old source to close, received false +cold same-session: expected /sessions/s1/events, received ?lastEventId=901 +stale source: expected an empty transcript, received "stale session one" +``` + +Frontend GREEN command: + +```text +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx \ + src/shell/AppShell.session-mgmt.test.tsx +cd frontend && npx tsc -b +``` + +GREEN output (exit 0): + +```text +Test Files 2 passed (2) +Tests 24 passed (24) +TypeScript: no output, exit 0 +``` + +Backend RED command: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/routes-sessions.test.ts +``` + +RED output (exit 1): + +```text +Test Files 2 failed (2) +Tests 8 failed | 28 passed (36) + +three Resume ordering assertions observed clear before reopen/create +reopen and create sentinels escaped as raw HTTP 500 responses +the concurrent waiter cold-started again instead of returning alreadyActive: true +SseHub.forget was absent and DELETE did not invoke permanent cleanup +``` + +Backend GREEN command: + +```text +cd backend && npx vitest run test/sse-hub.test.ts test/routes-sessions.test.ts +cd backend && npx tsc --noEmit -p . +``` + +GREEN output (exit 0): + +```text +Test Files 2 passed (2) +Tests 36 passed (36) +TypeScript: no output, exit 0 +``` + +The failure tests publish a post-failure probe through the same hub and prove that a subscriber +attached before either reopen or create rejection still receives it. The concurrency test overlaps +two same-id requests behind a deferred reopen and proves one manifest/readiness/reopen/create/clear +sequence. + +### Initial re-review verification (before independent-review hardening) + +```text +cd backend && npx vitest run +Test Files 22 passed (22) +Tests 182 passed (182) + +cd frontend && npx vitest run +Test Files 43 passed (43) +Tests 273 passed (273) + +cd backend && npm run build +> tsc -p tsconfig.json +exit 0 + +cd frontend && npm run build +> tsc -b && vite build +✓ 4835 modules transformed. +✓ built in 8.46s +exit 0 +``` + +`git diff --check` produced no output (exit 0). The frontend build retains its pre-existing +large-chunk warning; no new build or type errors were introduced. + Final whitespace verification: ```text @@ -328,15 +474,15 @@ no output, exit 0 ## Self-review -- Resume sequencing: backend `clear` and runtime binding precede the HTTP success; frontend state - mutation and stream generation follow it. Failure catch only emits fixed UI copy. +- Resume sequencing: reopen and runtime binding precede backend `clear` and HTTP success; frontend + state mutation and cursor-reset epoch follow it. Failure catch only emits fixed UI copy. - Already active: same-session returns before reset/generation/manifest repaint; different session resets the single-session store and binds the new id only after success. - SSE exact-once: ids are transport identity, not content hashes; replay is strictly `id > cursor`; `clear` retains the counter; pending gate matching uses only descriptor id. -- Cursor behavior: hook tracks `MessageEvent.lastEventId`, carries it only to a same-id generation, - and resets it on session-id change. Native EventSource reconnect remains supported by the route - header. +- Cursor behavior: hook tracks `MessageEvent.lastEventId`, carries it only to an ordinary same-id + generation, and resets it on session-id/cold-runtime epoch change. Native EventSource reconnect + remains supported by the route header. - Gate defense: the Zustand set survives pending clear but resets with the session store. - Client boundary: raw `ensure.error` is unused in public responses; generic Pi system events are reconstructed rather than spread; frontend type mirrors the two-field event. @@ -348,12 +494,123 @@ no output, exit 0 - The 200-event SSE ring limit remains intentional. A brand-new page can reconstruct only retained backlog; an in-memory same-session reconnect is exact-once from its cursor. -- Per-session sequence counters remain in backend memory after `clear` by design so later cold - same-id Resume cannot reuse ids. This is one numeric map entry per session id for the process - lifetime. +- Per-session sequence counters remain in backend memory after `clear` by design so later in-process + cold same-id Resume cannot reuse ids. Permanent DELETE removes the counter via `forget`. - Frontend tests still print pre-existing MSW unhandled-request and React ref/`act` warnings even - though all 271 tests pass. The frontend production build still reports pre-existing large chunk + though all 276 tests pass. The frontend production build still reports pre-existing large chunk warnings. Neither warning class was introduced or expanded by this change. - No live Pi/DWH smoke was run; this wave changes only REST/SSE/frontend lifecycle boundaries and is covered by fake-Pi, live Fastify SSE, component, full-suite, typecheck, and production-build gates. + +## Independent-review hardening + +The required independent review was run repeatedly against the uncommitted diff. Its first pass +found four Important lifecycle edges beyond the integrated findings: queued old-runtime callbacks, +post-spawn construction cleanup, concurrent frontend Resume completions, and the passive-effect +commit window. Its second pass confirmed those fixes and identified one remaining Important +retention issue in the new runtime-identity map. The final pass reported no Critical, Important, or +Minor findings and assessed the diff ready to merge. + +The resulting hardening is: + +- Runtime bridge callbacks are gated by the bound runtime identity. Replacement, close, and DELETE + invalidate the old identity, so queued old events cannot publish or call `failSession`. An active + runtime removed by the manager can still publish its complete public failure sequence; after the + terminal unmanaged `agent_end`, its binding is released and later events are rejected. +- `PiProcessManager` kills the spawned child and removes any registered map entry if either + spawn-boundary stderr setup or later RPC/bridge/map initialization throws. +- Resume completion compares against synchronously maintained current active-session identity. + Concurrent `alreadyActive: false` then `alreadyActive: true` results preserve the cold source, + cursor, store, and replayed gate. +- Stream source replacement uses a layout effect. A deterministic later-layout-effect test delivers + a queued old event inside the former commit-to-passive-cleanup window and proves it is ignored. +- Cursor tests cover both a restarted backend's fresh low ids and an in-process hub's preserved high + ids followed by a cursor-bearing ordinary reconnect. + +### Hardening TDD RED/GREEN evidence + +Backend identity/construction RED command: + +```text +cd backend && npx vitest run test/pi-process-manager.test.ts test/routes-sessions.test.ts +``` + +```text +Test Files 2 failed (2) +Tests 3 failed | 69 passed (72) + +post-spawn reader initialization did not kill the child +replaced and deleted runtime callbacks still called failSession/published +``` + +Additional spawn-boundary and terminal-release RED checks: + +```text +cd backend && npx vitest run test/pi-process-manager.test.ts \ + -t "spawn boundary initialization" +Tests 1 failed | 38 skipped (39) + +cd backend && npx vitest run test/routes-sessions.test.ts -t "terminal sequence" +Tests 1 failed | 34 skipped (35) +``` + +Frontend concurrency/layout RED command: + +```text +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx \ + src/shell/AppShell.session-mgmt.test.tsx +``` + +```text +Test Files 2 failed (2) +Tests 2 failed | 25 passed (27) + +the later-layout-effect event wrote "commit-window stale text" +the false→true completion pair erased pending gate "cold-gate" +``` + +Final focused GREEN commands: + +```text +cd backend && npx vitest run test/pi-process-manager.test.ts \ + test/routes-sessions.test.ts test/sse-hub.test.ts +cd backend && npx tsc --noEmit -p . + +Test Files 3 passed (3) +Tests 79 passed (79) +TypeScript: no output, exit 0 + +cd frontend && npx vitest run src/stream/useSessionStream.test.tsx \ + src/shell/AppShell.session-mgmt.test.tsx +cd frontend && npx tsc -b + +Test Files 2 passed (2) +Tests 27 passed (27) +TypeScript: no output, exit 0 +``` + +### Final full verification after review hardening + +```text +cd backend && npx vitest run +Test Files 22 passed (22) +Tests 188 passed (188) + +cd frontend && npx vitest run +Test Files 43 passed (43) +Tests 276 passed (276) + +cd backend && npm run build +> tsc -p tsconfig.json +exit 0 + +cd frontend && npm run build +> tsc -b && vite build +✓ 4835 modules transformed. +✓ built in 8.47s +exit 0 +``` + +The final frontend run retains the repository's pre-existing MSW/ref/`act` warnings, and the build +retains the pre-existing large-chunk warning. No test, typecheck, or build failures remain. diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index 6933f079..ae9d3461 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -68,9 +68,14 @@ export class PiProcessManager { 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; + try { + // Log stderr for debugging (was silently drained) + child.stderr.on("data", (d: Buffer) => console.error(`[pi:${sessionId}] stderr:`, d.toString().trim())); + return child; + } catch (error) { + try { child.kill(); } catch { /* preserve the initialization error */ } + throw error; + } } count(): number { return this.runtimes.size; } @@ -91,31 +96,39 @@ export class PiProcessManager { 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" }); - } - }); + let rt: SessionRuntime | undefined; + try { + const rpc = new RpcClient(child); + const bridge = new SessionBridge(rpc); + const runtime: SessionRuntime = { rpc, bridge, child }; + rt = runtime; + bridge.beginTurn(); + this.runtimes.set(sessionId, runtime); + // 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) === runtime}`); + if (this.runtimes.get(sessionId) === runtime) { + runtime.bridge.markFailed(); + this.runtimes.delete(sessionId); + runtime.bridge.emitClientEvent({ + type: "info", + level: "error", + text: `Pi process exited unexpectedly (code ${code ?? "?"})`, + }); + runtime.bridge.emitClientEvent({ type: "system_event", event: "session_failed" }); + runtime.bridge.emitClientEvent({ type: "system_event", event: "agent_end" }); + } + }); - return rt; + return runtime; + } catch (error) { + if (rt && this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId); + try { child.kill(); } catch { /* preserve the initialization error */ } + throw error; + } } /** Configure model and thinking. Safe to run alongside deterministic retrieval. */ diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 3835c761..7dac44b4 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -10,6 +10,8 @@ const BOOTSTRAP_FAILURE_MESSAGE = "Session startup failed. Check configuration and connectivity, then Resume the session."; const READINESS_FAILURE_MESSAGE = "Session services are not ready. Check configuration and connectivity, then try again."; +const RESUME_FAILURE_MESSAGE = + "Session could not be resumed. Check configuration and connectivity, then try again."; function eventCursor(...values: unknown[]): number { let cursor = 0; @@ -25,16 +27,58 @@ export function sessionRoutes( app: FastifyInstance, d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager }, ) { + const resumeTails = new Map>(); + const boundRuntimes = new Map< + string, + ReturnType + >(); + + const withResumeLock = async (id: string, work: () => Promise): Promise => { + const previous = resumeTails.get(id) ?? Promise.resolve(); + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const tail = previous.then(() => gate); + resumeTails.set(id, tail); + await previous; + try { + return await work(); + } finally { + release(); + if (resumeTails.get(id) === tail) resumeTails.delete(id); + } + }; + const info = (id: string, text: string, level = "info") => d.hub.publish(id, "info", { type: "info", level, text }); - const bindRuntime = (id: string, rt: ReturnType) => - rt.bridge.onClientEvent((e) => { - if (e.type === "system_event" && e.event === "session_failed") { - void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined); - } - d.hub.publish(id, e.type, e); - }); + const bindRuntime = (id: string, rt: ReturnType) => { + const previous = boundRuntimes.get(id); + boundRuntimes.set(id, rt); + try { + rt.bridge.onClientEvent((e) => { + // Child termination is asynchronous. Ignore queued events from a runtime once a newer + // identity is bound or the session is explicitly closed/deleted. The active identity + // remains bound after an unexpected exit so its public failure events still reach SSE. + if (boundRuntimes.get(id) !== rt) return; + if (e.type === "system_event" && e.event === "session_failed") { + void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined); + } + d.hub.publish(id, e.type, e); + if ( + e.type === "system_event" + && e.event === "agent_end" + && d.mgr.get(id) !== rt + && boundRuntimes.get(id) === rt + ) { + boundRuntimes.delete(id); + } + }); + } catch (error) { + if (previous) boundRuntimes.set(id, previous); + else if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id); + throw error; + } + }; const bootstrap = ( id: string, @@ -114,42 +158,69 @@ export function sessionRoutes( }); app.post("/sessions/:id/resume", async (req, reply) => { const id = (req.params as any).id; - const existing = d.mgr.get(id); - if (existing) { - const state = existing.bridge.turnState(); - if (state === "running" || state === "waiting") { - return reply.code(200).send({ id, alreadyActive: true }); + return withResumeLock(id, async () => { + // This check belongs inside the per-session lock: a preceding cold Resume may have + // installed a running runtime while this request was waiting. + const existing = d.mgr.get(id); + if (existing) { + const state = existing.bridge.turnState(); + if (state === "running" || state === "waiting") { + return reply.code(200).send({ id, alreadyActive: true }); + } } - } - const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null; - if (manifest?.status === "finalized" || manifest?.archived) { - return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" }); - } - const settings = d.getSettings(); - const ensure = await d.readiness.ensure(settings.workspace ?? ""); - if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE }); - const saved = manifest as { provider?: string; model?: string; thinking?: string } | null; - const options = { - provider: saved?.provider, - model: saved?.model, - thinking: saved?.thinking ?? settings.thinking, - author: getUser(req).id, - mode: "resume" as const, - }; - if (existing) d.mgr.teardown(id); - d.hub.clear(id); - await d.tht.reopenSession(id, settings.workspace); - const rt = d.mgr.createFor(id, options); - bindRuntime(id, rt); - info(id, "Resuming session"); - bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options)); - return reply.code(200).send({ id, alreadyActive: false }); + const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null; + if (manifest?.status === "finalized" || manifest?.archived) { + return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" }); + } + const settings = d.getSettings(); + const ensure = await d.readiness.ensure(settings.workspace ?? ""); + if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE }); + const saved = manifest as { provider?: string; model?: string; thinking?: string } | null; + const options = { + provider: saved?.provider, + model: saved?.model, + thinking: saved?.thinking ?? settings.thinking, + author: getUser(req).id, + mode: "resume" as const, + }; + + // Reopening is validation, not the transport commit point. Keep the old hub intact if + // persistence cannot be reopened. + try { + await d.tht.reopenSession(id, settings.workspace); + } catch { + return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE }); + } + + let rt: ReturnType | undefined; + try { + if (existing) { + boundRuntimes.delete(id); + d.mgr.teardown(id); + } + rt = d.mgr.createFor(id, options); + bindRuntime(id, rt); + } catch { + // A created-but-unbound runtime is not usable. The old hub remains attached because + // clear() has not happened yet. + if (rt) d.mgr.teardown(id); + return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE }); + } + + // Commit the replacement only after reopen + runtime creation/binding succeeded, and + // immediately before the first event produced by the new Resume. + d.hub.clear(id); + info(id, "Resuming session"); + bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options)); + return reply.code(200).send({ id, alreadyActive: false }); + }); }); app.post("/sessions/:id/close", async (req) => { const id = (req.params as { id: string }).id; try { await d.tht.closeSession(id, d.getSettings().workspace); } finally { + boundRuntimes.delete(id); d.mgr.teardown(id); d.hub.clear(id); } @@ -202,8 +273,10 @@ export function sessionRoutes( }); app.delete("/sessions/:id", async (req, reply) => { const id = (req.params as any).id; + boundRuntimes.delete(id); d.mgr.teardown(id); // drop any live runtime before deleting on disk await d.tht.deleteSession(id, d.getSettings().workspace); + d.hub.forget(id); return reply.code(204).send(); }); app.get("/sessions/:id/documents", async (req) => d.tht.documents((req.params as any).id)); diff --git a/backend/src/sse/sse-hub.ts b/backend/src/sse/sse-hub.ts index 1e63abc7..367b1e03 100644 --- a/backend/src/sse/sse-hub.ts +++ b/backend/src/sse/sse-hub.ts @@ -91,4 +91,10 @@ export class SseHub { this.buffers.delete(sessionId); this.subs.delete(sessionId); } + + /** Permanently discard transport state for a deleted session id. */ + forget(sessionId: string): void { + this.clear(sessionId); + this.lastIds.delete(sessionId); + } } diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index b6d88466..41a99085 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -84,6 +84,36 @@ function recordingChild() { return ch; } +test("createFor kills a spawned child when post-spawn initialization throws", () => { + const child = recordingChild(); + child.kill = vi.fn(); + child.stdout = { + on: () => { throw new Error("READER_INIT_SENTINEL"); }, + }; + const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); + + expect(() => mgr.createFor("broken-init", {})).toThrow("READER_INIT_SENTINEL"); + + expect(child.kill).toHaveBeenCalledOnce(); + expect(mgr.get("broken-init")).toBeUndefined(); + expect(mgr.count()).toBe(0); +}); + +test("createFor kills a spawned child when spawn boundary initialization throws", () => { + const child = recordingChild(); + child.kill = vi.fn(); + child.stderr = { + on: () => { throw new Error("STDERR_INIT_SENTINEL"); }, + }; + const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); + + expect(() => mgr.createFor("broken-spawn-init", {})).toThrow("STDERR_INIT_SENTINEL"); + + expect(child.kill).toHaveBeenCalledOnce(); + expect(mgr.get("broken-spawn-init")).toBeUndefined(); + expect(mgr.count()).toBe(0); +}); + test("un exit INATTESO del child notifica il client (info error + agent_end)", async () => { const cfg = loadConfig({}); const child = recordingChild(); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 7f40980c..79eac85e 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -5,6 +5,7 @@ import os from "node:os"; import { chmodSync, unlinkSync, writeFileSync } from "node:fs"; import { buildApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; +import { SseHub } from "../src/sse/sse-hub.js"; const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs"); const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json"); @@ -186,7 +187,7 @@ test.each(["idle", "failed"])( const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); await new Promise((resolve) => setImmediate(resolve)); expect(response.json()).toEqual({ id: "s1", alreadyActive: false }); - expect(order).toEqual(["teardown:s1", "clear:s1", "reopen", "create", "start"]); + expect(order).toEqual(["reopen", "teardown:s1", "create", "clear:s1", "start"]); }, ); @@ -224,7 +225,7 @@ test("POST resume without a runtime clears stale SSE state before cold start", a await new Promise((resolve) => setImmediate(resolve)); expect(response.json()).toEqual({ id: "crashed", alreadyActive: false }); - expect(order).toEqual(["clear:crashed", "reopen", "create", "start"]); + expect(order).toEqual(["reopen", "create", "clear:crashed", "start"]); expect(createOptions).toMatchObject({ provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low", mode: "resume", }); @@ -275,6 +276,364 @@ test("POST resume keeps a failed runtime stream attached when readiness fails", expect(cleared).toBe(false); }); +test("POST resume sanitizes reopen failure and preserves the old hub attachment", async () => { + const delivered: string[] = []; + const actualHub = new SseHub(); + let clearCalls = 0; + let tornDown = false; + const hub = { + subscribe: actualHub.subscribe.bind(actualHub), + publish: actualHub.publish.bind(actualHub), + clear: (id: string) => { clearCalls += 1; actualHub.clear(id); }, + } as any; + hub.subscribe("s1", (_event: string, data: any) => delivered.push(data.text)); + hub.publish("s1", "info", { text: "before" }); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => ({ bridge: { turnState: () => "idle" } }), + teardown: () => { tornDown = true; }, + } as any, + hub, + thtRunner: { + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => { throw new Error("REOPEN_SENTINEL"); }, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + hub.publish("s1", "info", { text: "post-failure probe" }); + + expect(response.statusCode).toBe(503); + expect(response.json()).toEqual({ + error: "Session could not be resumed. Check configuration and connectivity, then try again.", + }); + expect(response.body).not.toContain("REOPEN_SENTINEL"); + expect(clearCalls).toBe(0); + expect(tornDown).toBe(false); + expect(delivered).toEqual(["before", "post-failure probe"]); +}); + +test("POST resume sanitizes runtime creation failure and preserves the old hub attachment", async () => { + const delivered: string[] = []; + const actualHub = new SseHub(); + let clearCalls = 0; + let createCalls = 0; + const hub = { + subscribe: actualHub.subscribe.bind(actualHub), + publish: actualHub.publish.bind(actualHub), + clear: (id: string) => { clearCalls += 1; actualHub.clear(id); }, + } as any; + hub.subscribe("s1", (_event: string, data: any) => delivered.push(data.text)); + hub.publish("s1", "info", { text: "before" }); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => ({ bridge: { turnState: () => "failed" } }), + teardown: () => {}, + createFor: () => { createCalls += 1; throw new Error("CREATE_SENTINEL"); }, + } as any, + hub, + thtRunner: { + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + hub.publish("s1", "info", { text: "post-failure probe" }); + + expect(response.statusCode).toBe(503); + expect(response.json()).toEqual({ + error: "Session could not be resumed. Check configuration and connectivity, then try again.", + }); + expect(response.body).not.toContain("CREATE_SENTINEL"); + expect(createCalls).toBe(1); + expect(clearCalls).toBe(0); + expect(delivered).toEqual(["before", "post-failure probe"]); +}); + +test("concurrent cold Resume requests serialize and create one runtime", async () => { + let runtime: any; + let showCalls = 0; + let readinessCalls = 0; + let reopenCalls = 0; + let createCalls = 0; + let clearCalls = 0; + let releaseReopen!: () => void; + let markReopenStarted!: () => void; + const reopenStarted = new Promise((resolve) => { markReopenStarted = resolve; }); + const reopenReleased = new Promise((resolve) => { releaseReopen = resolve; }); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => runtime, + createFor: () => { + createCalls += 1; + runtime = { + bridge: { + turnState: () => "running", + onClientEvent: () => {}, + emitClientEvent: () => {}, + }, + }; + return runtime; + }, + configure: async () => {}, + start: () => {}, + teardown: () => {}, + } as any, + hub: { + clear: () => { clearCalls += 1; }, + publish: () => {}, + } as any, + thtRunner: { + sessionShow: async () => { + showCalls += 1; + return { status: "open", archived: false }; + }, + reopenSession: async () => { + reopenCalls += 1; + markReopenStarted(); + await reopenReleased; + }, + } as any, + readiness: { + ensure: async () => { readinessCalls += 1; return { ok: true }; }, + } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const first = app.inject({ method: "POST", url: "/sessions/s1/resume" }); + await reopenStarted; + const second = app.inject({ method: "POST", url: "/sessions/s1/resume" }); + await new Promise((resolve) => setImmediate(resolve)); + releaseReopen(); + const [firstResponse, secondResponse] = await Promise.all([first, second]); + await new Promise((resolve) => setImmediate(resolve)); + + expect(firstResponse.json()).toEqual({ id: "s1", alreadyActive: false }); + expect(secondResponse.json()).toEqual({ id: "s1", alreadyActive: true }); + expect({ showCalls, readinessCalls, reopenCalls, createCalls, clearCalls }).toEqual({ + showCalls: 1, + readinessCalls: 1, + reopenCalls: 1, + createCalls: 1, + clearCalls: 1, + }); +}); + +test("a replaced runtime cannot publish or fail the newly resumed session", async () => { + const controlledBridge = (initialState: string) => { + let state = initialState; + let listener: ((event: any) => void) | undefined; + return { + turnState: () => state, + setState: (next: string) => { state = next; }, + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + pendingWidget: () => null, + }; + }; + const oldBridge = controlledBridge("running"); + const newBridge = controlledBridge("running"); + const runtimes = [{ bridge: oldBridge }, { bridge: newBridge }]; + let current: any; + let failed = 0; + const published: any[] = []; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtimes.shift(); return current; }, + configure: async () => {}, + start: () => {}, + teardown: () => { current = undefined; }, + } as any, + hub: { + clear: () => {}, + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + failSession: async () => { failed += 1; }, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); + await new Promise((resolve) => setImmediate(resolve)); + oldBridge.setState("idle"); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + published.length = 0; + // The manager removes an unexpectedly exited active runtime before its bridge publishes the + // public failure sequence. Binding identity, rather than mgr.get(), must still admit it. + current = undefined; + oldBridge.emitClientEvent({ type: "system_event", event: "session_failed" }); + newBridge.emitClientEvent({ type: "info", text: "new runtime" }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(response.json()).toEqual({ id: "s1", alreadyActive: false }); + expect(failed).toBe(0); + expect(published).toEqual([{ + event: "info", + data: { type: "info", text: "new runtime" }, + }]); +}); + +test("a deleted runtime cannot repopulate or fail the forgotten session", async () => { + let listener: ((event: any) => void) | undefined; + const bridge = { + turnState: () => "running", + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + }; + const runtime = { bridge }; + let current: any; + let failed = 0; + const published: any[] = []; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtime; return runtime; }, + configure: async () => {}, + start: () => {}, + teardown: () => { current = undefined; }, + } as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + forget: () => {}, + } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + deleteSession: async () => {}, + failSession: async () => { failed += 1; }, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); + await new Promise((resolve) => setImmediate(resolve)); + + const response = await app.inject({ method: "DELETE", url: "/sessions/s1" }); + published.length = 0; + bridge.emitClientEvent({ type: "system_event", event: "session_failed" }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(response.statusCode).toBe(204); + expect(failed).toBe(0); + expect(published).toEqual([]); +}); + +test("an unexpectedly exited runtime publishes its terminal sequence then releases its binding", async () => { + let listener: ((event: any) => void) | undefined; + const bridge = { + onClientEvent: (next: (event: any) => void) => { listener = next; }, + emitClientEvent: (event: any) => listener?.(event), + }; + const runtime = { bridge }; + let current: any; + let failed = 0; + const published: any[] = []; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + createFor: () => { current = runtime; return runtime; }, + configure: async () => {}, + start: () => {}, + teardown: () => { current = undefined; }, + } as any, + hub: { + publish: (_id: string, event: string, data: object) => { + published.push({ event, data }); + return published.length; + }, + } as any, + thtRunner: { + sessionNew: async () => ({ id: "s1" }), + searchPack: async () => {}, + failSession: async () => { failed += 1; }, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); + await new Promise((resolve) => setImmediate(resolve)); + published.length = 0; + + // PiProcessManager deletes an unexpectedly exited runtime before emitting this sequence. + current = undefined; + bridge.emitClientEvent({ type: "info", level: "error", text: "public failure" }); + bridge.emitClientEvent({ type: "system_event", event: "session_failed" }); + bridge.emitClientEvent({ type: "system_event", event: "agent_end" }); + bridge.emitClientEvent({ type: "text_delta", text: "too late" }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(failed).toBe(1); + expect(published).toEqual([ + { event: "info", data: { type: "info", level: "error", text: "public failure" } }, + { event: "system_event", data: { type: "system_event", event: "session_failed" } }, + { event: "system_event", data: { type: "system_event", event: "agent_end" } }, + ]); +}); + +test("POST resume tears down a created runtime when bridge binding fails", async () => { + const delivered: string[] = []; + const actualHub = new SseHub(); + let clearCalls = 0; + let teardownCalls = 0; + let current: any = { bridge: { turnState: () => "failed" } }; + const hub = { + subscribe: actualHub.subscribe.bind(actualHub), + publish: actualHub.publish.bind(actualHub), + clear: (id: string) => { clearCalls += 1; actualHub.clear(id); }, + } as any; + hub.subscribe("s1", (_event: string, data: any) => delivered.push(data.text)); + hub.publish("s1", "info", { text: "before" }); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr: { + get: () => current, + teardown: () => { teardownCalls += 1; current = undefined; }, + createFor: () => { + current = { + bridge: { + onClientEvent: () => { throw new Error("BIND_SENTINEL"); }, + }, + }; + return current; + }, + } as any, + hub, + thtRunner: { + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + } as any, + readiness: { ensure: async () => ({ ok: true }) } as any, + getSettings: () => ({ workspace: "local" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" }); + hub.publish("s1", "info", { text: "post-failure probe" }); + + expect(response.statusCode).toBe(503); + expect(response.body).not.toContain("BIND_SENTINEL"); + expect(teardownCalls).toBe(2); + expect(current).toBeUndefined(); + expect(clearCalls).toBe(0); + expect(delivered).toEqual(["before", "post-failure probe"]); +}); + test("POST /sessions/:id/response inoltra al bridge (no error)", async () => { const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { @@ -332,12 +691,31 @@ test("DELETE /sessions/:id tears down the runtime before deleting on disk", asyn const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { deleteSession: async (id: string) => { order.push(`del:${id}`); } } as any, mgr: { teardown: (id: string) => { order.push(`teardown:${id}`); } } as any, + hub: { forget: (id: string) => { order.push(`forget:${id}`); } } as any, getSettings: () => ({ workspace: "w" }) as any, spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any, }); const res = await app.inject({ method: "DELETE", url: "/sessions/s1" }); expect(res.statusCode).toBe(204); - expect(order).toEqual(["teardown:s1", "del:s1"]); + expect(order).toEqual(["teardown:s1", "del:s1", "forget:s1"]); +}); + +test("POST /sessions/:id/close clears transient hub state without forgetting its id", async () => { + const calls: string[] = []; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + thtRunner: { closeSession: async () => {} } as any, + mgr: { teardown: () => {} } as any, + hub: { + clear: (id: string) => { calls.push(`clear:${id}`); }, + forget: (id: string) => { calls.push(`forget:${id}`); }, + } as any, + getSettings: () => ({ workspace: "w" }) as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/s1/close" }); + + expect(response.statusCode).toBe(200); + expect(calls).toEqual(["clear:s1"]); }); test("GET /sessions/:id/documents returns the runner output", async () => { @@ -495,6 +873,7 @@ test("POST /sessions bootstrap failure emits only a fixed recovery message", asy "connect https://secret.invalid/bootstrap?token=DO_NOT_LEAK using /srv/private/model-key"; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { mgr: { + get: () => undefined, createFor: () => runtime, configure: async () => { throw new Error(rawFailure); }, start: () => {}, diff --git a/backend/test/sse-hub.test.ts b/backend/test/sse-hub.test.ts index 581b7a36..24627c9e 100644 --- a/backend/test/sse-hub.test.ts +++ b/backend/test/sse-hub.test.ts @@ -40,6 +40,21 @@ test("clear drops stale subscribers and buffers but preserves the per-session id expect(resumed).toEqual([{ event: "info", data: { text: "after" }, id: 2 }]); }); +test("forget drops subscribers, buffers, and the per-session id sequence", () => { + const hub = new SseHub(); + const stale: any[] = []; + hub.subscribe("s1", (event, data, id) => stale.push({ event, data, id })); + expect(hub.publish("s1", "info", { text: "before" })).toBe(1); + + hub.forget("s1"); + expect(hub.publish("s1", "info", { text: "after" })).toBe(1); + const fresh: any[] = []; + hub.subscribe("s1", (event, data, id) => fresh.push({ event, data, id })); + + expect(stale).toEqual([{ event: "info", data: { text: "before" }, id: 1 }]); + expect(fresh).toEqual([{ event: "info", data: { text: "after" }, id: 1 }]); +}); + test("a buffered pending gate is emitted exactly once and matched by descriptor id", () => { const hub = new SseHub(); const descriptor = { id: "gate-1", widget: "select", title: "Original" }; diff --git a/frontend/src/shell/AppShell.session-mgmt.test.tsx b/frontend/src/shell/AppShell.session-mgmt.test.tsx index ef7b4978..8894996f 100644 --- a/frontend/src/shell/AppShell.session-mgmt.test.tsx +++ b/frontend/src/shell/AppShell.session-mgmt.test.tsx @@ -134,6 +134,60 @@ test("an already-active same-session Resume preserves its EventSource and store" expect(useSessionStore.getState().activityLog).toEqual(before); }); +test("a false then true pair of concurrent Resume completions preserves the cold stream", async () => { + let resumeCalls = 0; + let releaseCold!: () => void; + let releaseAlreadyActive!: () => void; + let markColdStarted!: () => void; + let markAlreadyActiveStarted!: () => void; + const coldStarted = new Promise((resolve) => { markColdStarted = resolve; }); + const alreadyActiveStarted = new Promise((resolve) => { markAlreadyActiveStarted = resolve; }); + const coldReleased = new Promise((resolve) => { releaseCold = resolve; }); + const alreadyActiveReleased = new Promise((resolve) => { releaseAlreadyActive = resolve; }); + server.use( + http.post("http://localhost:8787/sessions/:id/resume", async () => { + resumeCalls += 1; + if (resumeCalls === 1) { + markColdStarted(); + await coldReleased; + return resumeResult("s1", false); + } + markAlreadyActiveStarted(); + await alreadyActiveReleased; + return resumeResult("s1", true); + }), + http.get("http://localhost:8787/sessions/:id", () => + HttpResponse.json({ id: "s1", status: "open", phase: 1 })), + ); + wrap(); + await userEvent.click(await screen.findByText("Attiva uno")); + const resumeButton = await screen.findByRole("button", { name: /resume/i }); + + await userEvent.click(resumeButton); + await coldStarted; + await userEvent.click(resumeButton); + await alreadyActiveStarted; + expect(resumeCalls).toBe(2); + + releaseCold(); + await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); + const stream = FakeEventSource.instances[0]; + act(() => stream.emitNamed( + "ui_request", + { type: "ui_request", ui_request: { id: "cold-gate", widget: "select" } }, + "1", + )); + const before = useSessionStore.getState().activityLog.map((entry) => ({ ...entry })); + + releaseAlreadyActive(); + await waitFor(() => expect(screen.queryByText("Domanda originale")).not.toBeInTheDocument()); + + expect(FakeEventSource.instances).toHaveLength(1); + expect(stream.closed).toBe(false); + expect(useSessionStore.getState().pendingWidget?.id).toBe("cold-gate"); + expect(useSessionStore.getState().activityLog).toEqual(before); +}); + test("cold same-session Resume keeps the old stream until success then receives post-clear events once", async () => { server.use( http.post("http://localhost:8787/sessions/:id/resume", () => @@ -146,7 +200,7 @@ test("cold same-session Resume keeps the old stream until success then receives await userEvent.click(await screen.findByRole("button", { name: /resume/i })); await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1)); const first = FakeEventSource.instances[0]; - act(() => first.emitNamed("info", { type: "info", text: "Before cold Resume" }, "7")); + act(() => first.emitNamed("info", { type: "info", text: "Before cold Resume" }, "900")); let releaseResume!: () => void; let markStarted!: () => void; @@ -168,7 +222,7 @@ test("cold same-session Resume keeps the old stream until success then receives expect(useSessionStore.getState().activityLog).toContainEqual({ kind: "status", phase: "F1", text: "Before cold Resume", level: "info", }); - act(() => first.emitNamed("text_delta", { type: "text_delta", text: "Still attached" }, "8")); + act(() => first.emitNamed("text_delta", { type: "text_delta", text: "Still attached" }, "901")); expect(useSessionStore.getState().transcript).toEqual([ { role: "assistant", text: "Still attached" }, ]); @@ -177,7 +231,7 @@ test("cold same-session Resume keeps the old stream until success then receives await waitFor(() => expect(FakeEventSource.instances).toHaveLength(2)); expect(first.closed).toBe(true); const replacement = FakeEventSource.instances[1]; - expect(replacement.url).toContain("/sessions/s1/events?lastEventId=8"); + expect(replacement.url).toBe("http://localhost:8787/sessions/s1/events"); expect(useSessionStore.getState().transcript).toEqual([]); expect(useSessionStore.getState().activityLog).toEqual([ { kind: "lifecycle", phase: null, text: "Resuming session" }, @@ -186,7 +240,7 @@ test("cold same-session Resume keeps the old stream until success then receives act(() => replacement.emitNamed( "text_delta", { type: "text_delta", text: "Post-resume delivery" }, - "9", + "1", )); expect(useSessionStore.getState().transcript).toEqual([ { role: "assistant", text: "Post-resume delivery" }, @@ -199,8 +253,8 @@ test("cold same-session Resume keeps the old stream until success then receives ui_request: { id: "post-resume-gate", widget: "select", title: "Review once" }, }; act(() => { - replacement.emitNamed("ui_request", gate, "10"); - replacement.emitNamed("ui_request", gate, "11"); + replacement.emitNamed("ui_request", gate, "2"); + replacement.emitNamed("ui_request", gate, "3"); }); expect(useSessionStore.getState().pendingWidget?.id).toBe("post-resume-gate"); expect(useSessionStore.getState().activityLog.filter((entry) => entry.kind === "gate")) diff --git a/frontend/src/shell/AppShell.tsx b/frontend/src/shell/AppShell.tsx index 6b2ecd72..ee454dfc 100644 --- a/frontend/src/shell/AppShell.tsx +++ b/frontend/src/shell/AppShell.tsx @@ -32,7 +32,8 @@ import { useEffect, useMemo, useRef, useState } from "react"; */ export function AppShell() { const [activeSessionId, setActiveSessionId] = useState(null); - const [streamGeneration, setStreamGeneration] = useState(0); + const activeSessionIdRef = useRef(null); + const [streamCursorResetEpoch, setStreamCursorResetEpoch] = useState(0); const [creatingSession, setCreatingSession] = useState(false); const [awaitingQuestion, setAwaitingQuestion] = useState(false); const { data: sessions = [] } = useQuery({ @@ -65,6 +66,12 @@ export function AppShell() { const selectedSessions = sessions.filter((session) => selectedSessionIds.has(session.id)); const allSessionsSelected = sessions.length > 0 && selectedSessions.length === sessions.length; + function selectActiveSession(id: string | null) { + // Keep async Resume completions synchronized before React commits the state update. + activeSessionIdRef.current = id; + setActiveSessionId(id); + } + // A background refresh can remove a session (for example from another browser). // Keep the local selection aligned with the authoritative list. useEffect(() => { @@ -100,9 +107,9 @@ export function AppShell() { }); } async function doResume(id: string) { - const reconnectSameSession = activeSessionId === id; try { const result = await resumeSession(id); + const reconnectSameSession = activeSessionIdRef.current === id; setPanelSession(null); setAwaitingQuestion(false); @@ -115,11 +122,11 @@ export function AppShell() { recordLifecycle("Resuming session"); setAgentActive(true); } - setActiveSessionId(id); - // The backend has now completed clear/rebind. Recreate a same-id source with its cursor; - // switching ids naturally creates a new source and resets the cursor in the stream hook. + selectActiveSession(id); + // A cold runtime starts a fresh SSE id sequence. Recreate a same-id source only after + // Resume succeeds, and explicitly discard the old runtime's cursor. if (!result.alreadyActive && reconnectSameSession) { - setStreamGeneration((value) => value + 1); + setStreamCursorResetEpoch((value) => value + 1); } // Paint the persisted re-entry phase while the replacement stream starts replaying. @@ -182,7 +189,7 @@ export function AppShell() { targets.filter((_, index) => results[index].status === "fulfilled").map((session) => session.id), ); if (deletedIds.has(panelSession?.id ?? "")) setPanelSession(null); - if (deletedIds.has(activeSessionId ?? "")) { resetSession(); setActiveSessionId(null); } + if (deletedIds.has(activeSessionId ?? "")) { resetSession(); selectActiveSession(null); } setSelectedSessionIds((current) => new Set([...current].filter((id) => !deletedIds.has(id)))); refresh(); if (deletedIds.size !== targets.length) { @@ -226,7 +233,7 @@ export function AppShell() { // session sits idle or a gate awaits the reviewer (pendingWidget). const running = working && !finalized; - useSessionStream(activeSessionId, streamGeneration); + useSessionStream(activeSessionId, 0, streamCursorResetEpoch); // A backend "session_exit" system event (e.g. the replay server emitting it // when the reviewer picks "Esci") asks us to leave the live session view and @@ -240,7 +247,7 @@ export function AppShell() { // Never let a streamed event terminate the managed Pi child. Only the // explicit “Stop & save” action is allowed to call /close. resetSession(); - setActiveSessionId(null); + selectActiveSession(null); setAwaitingQuestion(false); } // The final workflow turn ends with the session already finalized on disk: @@ -254,7 +261,7 @@ export function AppShell() { resetSession(); setAwaitingQuestion(true); setCreatingSession(false); - setActiveSessionId(null); + selectActiveSession(null); // Best effort only: session creation keeps the authoritative readiness gate. // Composer focus is deliberately independent of this network request. void prewarmRuntime().catch(() => undefined); @@ -269,7 +276,7 @@ export function AppShell() { function finishSessionCreation(id: string) { // React batches these updates, preserving the provisional view and timer // while useSessionStream opens the durable session's SSE channel. - setActiveSessionId(id); + selectActiveSession(id); setCreatingSession(false); setAwaitingQuestion(false); refresh(); @@ -287,7 +294,7 @@ export function AppShell() { await closeSession(activeSessionId); } finally { resetSession(); - setActiveSessionId(null); + selectActiveSession(null); setAwaitingQuestion(false); } } diff --git a/frontend/src/stream/useSessionStream.test.tsx b/frontend/src/stream/useSessionStream.test.tsx index f5c15d2b..58363e27 100644 --- a/frontend/src/stream/useSessionStream.test.tsx +++ b/frontend/src/stream/useSessionStream.test.tsx @@ -1,5 +1,5 @@ import { renderHook } from "@testing-library/react"; -import { act } from "react"; +import { act, useLayoutEffect } from "react"; import { FakeEventSource } from "../test/fakeEventSource"; import { useSessionStream } from "./useSessionStream"; import { useSessionStore } from "../store/sessionStore"; @@ -100,6 +100,58 @@ test("a same-session generation reconnect includes the last consumed SSE id", () .toEqual([{ kind: "assistant", phase: null, text: "hello world" }]); }); +test("a same-session reset epoch drops a high cursor and accepts fresh low-id events", () => { + const { rerender } = renderHook( + ({ resetEpoch }) => useSessionStream("s1", 0, resetEpoch), + { initialProps: { resetEpoch: 0 } }, + ); + const first = FakeEventSource.instances[0]; + act(() => first.emitNamed("info", { type: "info", text: "old runtime" }, "900")); + + rerender({ resetEpoch: 1 }); + + expect(first.closed).toBe(true); + expect(FakeEventSource.instances).toHaveLength(2); + expect(FakeEventSource.instances[1].url).toBe( + "http://localhost:8787/sessions/s1/events", + ); + + act(() => + FakeEventSource.instances[1].emitNamed( + "ui_request", + { type: "ui_request", ui_request: { id: "fresh-gate", widget: "select" } }, + "1", + ), + ); + expect(useSessionStore.getState().pendingWidget?.id).toBe("fresh-gate"); +}); + +test("a reset epoch also accepts an in-process preserved high id and resumes from it", () => { + const { rerender } = renderHook( + ({ generation, resetEpoch }) => useSessionStream("s1", generation, resetEpoch), + { initialProps: { generation: 0, resetEpoch: 0 } }, + ); + const first = FakeEventSource.instances[0]; + act(() => first.emitNamed("info", { type: "info", text: "old runtime" }, "900")); + + rerender({ generation: 0, resetEpoch: 1 }); + const replacement = FakeEventSource.instances[1]; + expect(replacement.url).toBe("http://localhost:8787/sessions/s1/events"); + act(() => + replacement.emitNamed( + "ui_request", + { type: "ui_request", ui_request: { id: "preserved-high-gate", widget: "select" } }, + "902", + ), + ); + + expect(useSessionStore.getState().pendingWidget?.id).toBe("preserved-high-gate"); + rerender({ generation: 1, resetEpoch: 1 }); + expect(FakeEventSource.instances[2].url).toBe( + "http://localhost:8787/sessions/s1/events?lastEventId=902", + ); +}); + test("changing the session resets the manual reconnect cursor", () => { const { rerender } = renderHook( ({ sessionId }) => useSessionStream(sessionId), @@ -113,3 +165,55 @@ test("changing the session resets the manual reconnect cursor", () => { expect(first.closed).toBe(true); expect(FakeEventSource.instances[1].url).toBe("http://localhost:8787/sessions/s2/events"); }); + +test("a queued event from a replaced source cannot mutate the new session or its cursor", () => { + const { rerender } = renderHook( + ({ sessionId, generation }) => useSessionStream(sessionId, generation, 0), + { initialProps: { sessionId: "s1" as string | null, generation: 0 } }, + ); + const first = FakeEventSource.instances[0]; + act(() => first.emitNamed("info", { type: "info", text: "session one" }, "12")); + + rerender({ sessionId: "s2", generation: 0 }); + expect(first.closed).toBe(true); + expect(FakeEventSource.instances[1].url).toBe("http://localhost:8787/sessions/s2/events"); + useSessionStore.getState().resetSession(); + + act(() => + first.emitQueuedNamed( + "text_delta", + { type: "text_delta", text: "stale session one" }, + "999", + ), + ); + expect(useSessionStore.getState().transcript).toEqual([]); + + rerender({ sessionId: "s2", generation: 1 }); + expect(FakeEventSource.instances[2].url).toBe("http://localhost:8787/sessions/s2/events"); +}); + +test("the old source is invalid before later layout effects can deliver a queued event", () => { + let first: FakeEventSource; + const { rerender } = renderHook( + ({ sessionId, emitDuringLayout }) => { + useSessionStream(sessionId); + useLayoutEffect(() => { + if (emitDuringLayout) { + first.emitQueuedNamed( + "text_delta", + { type: "text_delta", text: "commit-window stale text" }, + "777", + ); + } + }, [emitDuringLayout, sessionId]); + }, + { initialProps: { sessionId: "s1" as string | null, emitDuringLayout: false } }, + ); + first = FakeEventSource.instances[0]; + useSessionStore.getState().resetSession(); + + rerender({ sessionId: "s2", emitDuringLayout: true }); + + expect(first.closed).toBe(true); + expect(useSessionStore.getState().transcript).toEqual([]); +}); diff --git a/frontend/src/stream/useSessionStream.ts b/frontend/src/stream/useSessionStream.ts index 76d3086a..ca6f7d63 100644 --- a/frontend/src/stream/useSessionStream.ts +++ b/frontend/src/stream/useSessionStream.ts @@ -1,29 +1,51 @@ -import { useEffect, useRef, useState } from "react"; +import { useLayoutEffect, useRef, useState } from "react"; import { BASE } from "../api/client"; import { joinBackendPath } from "../api/runtime-config"; import { useSessionStore } from "../store/sessionStore"; import type { StreamEvent } from "../api/types"; -export function useSessionStream(sessionId: string | null, generation = 0) { +export function useSessionStream( + sessionId: string | null, + generation = 0, + cursorResetEpoch = 0, +) { const [connected, setConnected] = useState(false); const applyEvent = useSessionStore((s) => s.applyEvent); - const cursor = useRef({ sessionId: null as string | null, lastEventId: "" }); + const cursor = useRef({ + sessionId: null as string | null, + cursorResetEpoch, + lastEventId: "", + }); + const activeSource = useRef(null); - if (cursor.current.sessionId !== sessionId) { - cursor.current = { sessionId, lastEventId: "" }; - } - - useEffect(() => { - if (!sessionId) return; + useLayoutEffect(() => { + if ( + cursor.current.sessionId !== sessionId + || cursor.current.cursorResetEpoch !== cursorResetEpoch + ) { + cursor.current = { sessionId, cursorResetEpoch, lastEventId: "" }; + } + if (!sessionId) { + activeSource.current = null; + setConnected(false); + return; + } const query = cursor.current.lastEventId ? `?lastEventId=${encodeURIComponent(cursor.current.lastEventId)}` : ""; const es = new EventSource(joinBackendPath(BASE, `/sessions/${sessionId}/events${query}`)); - es.onopen = () => setConnected(true); - es.onerror = () => setConnected(false); + const identity = { source: es, sessionId, cursorResetEpoch }; + activeSource.current = identity; + es.onopen = () => { + if (activeSource.current === identity) setConnected(true); + }; + es.onerror = () => { + if (activeSource.current === identity) setConnected(false); + }; const handle = (ev: MessageEvent) => { + if (activeSource.current !== identity) return; if (ev.lastEventId) cursor.current.lastEventId = ev.lastEventId; try { applyEvent(JSON.parse(ev.data) as StreamEvent); @@ -47,10 +69,14 @@ export function useSessionStream(sessionId: string | null, generation = 0) { return () => { for (const name of namedEvents) es.removeEventListener(name, handle); + es.onmessage = null; es.close(); - setConnected(false); + if (activeSource.current === identity) { + activeSource.current = null; + setConnected(false); + } }; - }, [sessionId, generation, applyEvent]); + }, [sessionId, generation, cursorResetEpoch, applyEvent]); return { connected }; } diff --git a/frontend/src/test/fakeEventSource.ts b/frontend/src/test/fakeEventSource.ts index b239e777..bbce75dd 100644 --- a/frontend/src/test/fakeEventSource.ts +++ b/frontend/src/test/fakeEventSource.ts @@ -5,6 +5,10 @@ export class FakeEventSource { onerror: (() => void) | null = null; closed = false; private listeners = new Map void>>(); + private historicalListeners = new Map< + string, + Set<(e: { data: string; lastEventId: string }) => void> + >(); constructor(public url: string) { FakeEventSource.instances.push(this); } @@ -17,9 +21,16 @@ export class FakeEventSource { const data = JSON.stringify(obj); for (const h of this.listeners.get(type) ?? []) h({ data, lastEventId }); } + /** Simulate an event already queued before removeEventListener/close completed. */ + emitQueuedNamed(type: string, obj: unknown, lastEventId = "") { + const data = JSON.stringify(obj); + for (const h of this.historicalListeners.get(type) ?? []) h({ data, lastEventId }); + } addEventListener(type: string, handler: (e: { data: string; lastEventId: string }) => void) { if (!this.listeners.has(type)) this.listeners.set(type, new Set()); + if (!this.historicalListeners.has(type)) this.historicalListeners.set(type, new Set()); this.listeners.get(type)!.add(handler); + this.historicalListeners.get(type)!.add(handler); } removeEventListener(type: string, handler: (e: { data: string; lastEventId: string }) => void) { this.listeners.get(type)?.delete(handler);